When Spark stops compiling your query


Spark SQL compiles your query into Java. Most of the time it compiles all of it: each stage of the plan becomes one generated method, and the operators inside it pass values to each other in local variables instead of rows. Sometimes it compiles less than all of it, and it rarely says so. An operator drops out of its stage and nothing is logged. A stage stays in and runs slower than it would outside. A limit is crossed and the only line about it is at a level the shell hides. A method fails to compile and Spark quietly interprets it instead. And once in a while the query just fails.

This post is about those cases, in the order you would meet them: what each looks like, what it costs, how to see it on your own cluster, and what to do. It is written for people who run Spark, and it follows The 8000-byte cliff in Spark SQL, which covered the best known of them. Every snippet runs on a stock Spark 4.2.0 distribution.

If you have five minutes

  1. An operator without a *(n) in front of it in EXPLAIN runs outside whole-stage codegen. That is usually by design and cheap; section 1 says which ones are worth a look.
  2. A projection of fifty or more cheap columns - arithmetic, copies, casts of numbers - can be slower inside a stage than outside one. Do not raise spark.sql.codegen.maxFields to keep a wide projection in its stage; for a query like that, lowering it can help (section 2).
  3. Two lines of log4j2.properties show the give-ups Spark logs at INFO (section 3).
  4. To find all of this across a whole application's history, run one script over your event logs (the closing section).

How often does it happen? We counted every operator of every final plan in Spark's own query suites - the SQL golden files, TPC-DS, TPC-H and SSB, 22,284 plans and 93,787 operators. In the TPC suites almost nothing leaves its stage that is not meant to: 3.9% of TPC-DS's operators, most of them top-k sorts and windows, and not one operator for a function without generated code or for a schema that is too wide. In the golden files, which exercise everything Spark can do, it is 7.1%, and every one of them leaves without a line in the log.

Where the operators of Spark's own query suites runshare of every operator of every final plan, per suiteTPC-DS3.9%TPC-H1.2%SSBnonethe SQL golden files7.1%silent0%25%50%75%100%in a stagea columnar scan, by designleft out for a reason, silentlywhole-stage codegen switched off22,284 plans and 93,787 operators in all; the golden files run each query under several settings, some with whole-stage codegen off

Figure 1. The share of operators in a whole-stage codegen stage, outside one by design (a columnar scan, which a stage reads through a conversion), and outside one for a reason, per suite.

The benchmarks are not the queries people write, though. The ones that leave a stage in the golden files are the constructs a real query may well have: an object aggregate such as collect_list, a Python UDF, a window, a JSON function, a higher-order function over maps.

1. Silent: the operator that leaves its stage

In a physical plan, *(n) in front of an operator means it runs inside whole-stage codegen stage n - a codegen stage, not the Spark stage the UI numbers. An operator without it runs on its own: it reads rows, does its work, and writes rows. Nothing is logged when that happens, at any level, so the plan is the only place to see it. Run the query before you explain it: with adaptive execution on, the final plan exists only once the query has run.

A window. Window functions have no generated code; the operators around them do.

SQL> select id, sum(id) over (partition by g order by id) as running from t
+- *(2) Project [id#0L, running#17L]
   +- Window [sum(id#0L) windowspecdefinition(g#1L, id#0L ASC NULLS FIRST, ...
      +- *(1) Sort [g#1L ASC NULLS FIRST, id#0L ASC NULLS FIRST], false, 0

A top-k sort. ORDER BY ... LIMIT plans a TakeOrderedAndProject, which keeps the top rows of each partition in a heap and has no generated code. It is the most frequent give-up in TPC-DS.

SQL> select id, s from t order by s desc, id limit 5
TakeOrderedAndProject(limit=5, orderBy=[s#2 DESC NULLS LAST,id#0L ASC NULLS FIRST], ...
+- *(1) Project [id#0L, concat(row-, cast((id#0L % 100) as string)) AS s#2]
   +- *(1) Range (0, 1000, step=1, splits=1)

An aggregate that keeps an object per group. collect_list, collect_set and percentile_approx keep an object per group, not a row of fixed-width values, so Spark plans an ObjectHashAggregate, which is never in a stage; 1,239 of them run in the golden files.

SQL> select g, percentile_approx(id, 0.5) AS p50 from t group by g
+- ObjectHashAggregate(keys=[g#1L], functions=[percentile_approx(id#0L, 0.5, 10000, 0, 0)], ...
   +- ObjectHashAggregate(keys=[g#1L], functions=[partial_percentile_approx(id#0L, ...
      +- *(1) Project [id#0L, (id#0L % 10) AS g#1L]

A max of a string. A hash aggregate updates fixed-width values in place, and a string is not one, so max(s) per group sorts the rows and aggregates them in order. In Spark 4.2 a sort aggregate with grouping keys has no generated code; 4.3.0 adds it (section 6).

SQL> select g, max(s) AS last_s from t group by g
+- SortAggregate(key=[g#1L], functions=[max(s#2)], output=[g#1L, last_s#25])
   +- SortAggregate(key=[g#1L], functions=[partial_max(s#2)], output=[g#1L, max#28])
      +- *(1) Sort [g#1L ASC NULLS FIRST], false, 0

A function without generated code. One such function takes its whole operator out of the stage. from_json is the one you are likeliest to have; get_json_object over the same string generates code and keeps it:

SQL> select id + 1 AS a, from_json(js, 'a INT').a AS j from t
Project [(id#0L + 1) AS a#3L, from_json(StructField(a,IntegerType,true), concat(...
+- *(1) Range (0, 1000, step=1, splits=1)

SQL> select id + 1 AS a, get_json_object(js, '$.a') AS j from t
*(1) Project [(id#0L + 1) AS a#10L, get_json_object(concat({"a":, cast((id#0L % 10) as ...
+- *(1) Range (0, 1000, step=1, splits=1)

json_tuple and to_json generate code too. In Spark 4.2 no higher-order function does - not transform, filter, exists, forall or aggregate, and not the ones over maps or zip_with:

SQL> select id, transform(array(id, id + 1), x -> x * 2) AS twice from t
Project [id#0L, transform(array(id#0L, (id#0L + 1)), lambdafunction((lambda x#43L * 2), ...
+- *(1) Range (0, 1000, step=1, splits=1)

From 4.3.0 the five over arrays generate code (section 6). The others that leave a stage silently are a Python UDF's evaluation and a projection with more than a hundred output columns, the maxFields limit of section 2.

What it costs. Less than you might fear. An operator outside a stage still compiles its own expressions; what it loses is the stage, the passing of values in locals, so it pays to write a row and read it back at its boundary. We priced that for from_json in a projection of 16, 32 and 48 date expressions over a million cached rows: next to the same projection with get_json_object instead, which keeps the stage, it costs at most 135 to 224 nanoseconds a row, the same at every width - and the JSON function itself costs about a microsecond. The give-up is the smaller part of the bill.

What not to do. The obvious rewrite - compute everything else first, then from_json in a projection of its own - does not survive the optimizer. CollapseProject merges the two projections back into one, and over an aggregate it merges the projection into the aggregate, which then leaves its stage too:

+- HashAggregate(keys=[_groupingexpression#14L], functions=[sum(id#0L)])
   +- *(1) HashAggregate(keys=[_groupingexpression#14L], functions=[partial_sum(id#0L)])

What to do. Where a function that generates code does the job, use it: with get_json_object in place of from_json every operator above stays in its stage. Otherwise, leave it; a fifth of a microsecond a row is rarely the problem.

See it yourself: examples.scala, silent.scala and rewrite.scala print every plan in this section.

2. In a stage, and slower

This is the one that surprised us. A stage is supposed to be the fast path, and for most operators it is: an aggregate of 40 or 60 sums runs 1.6 to 2.1 times faster in a stage than without one. But a wide projection of cheap columns runs slower inside a stage than outside it, under the default settings, once everything is compiled:

projection AMD EPYC 7763 Intel Xeon 8573C laptop, Zen 5
50 cheap columns (id + k) 1.41 times slower 1.28 times slower 1.16 times slower
99 cheap columns 1.54 times slower 1.65 times slower 1.25 times slower
50 mixed columns 7% faster 12% faster
99 mixed columns 15% faster 16% faster

The mixed projection cycles through six kinds of column: an addition, date_add, a string concat, a division over a cast, a null test over a nullable column and substr.

A wide projection, in a stage and out of ittime a row in a whole-stage codegen stage over the time without one, defaults, after C20.81.01.21.41.6the same speedslower in a stagefaster in a stage1.41xrunner1.16xlaptop50 cheap columns1.54xrunner1.25xlaptop99 cheap columns-7%runner-12%laptop50 mixed columns-15%runner-16%laptop99 mixed columnsa column that does work of its owncosts more than the call C2leaves behindrunner: AMD EPYC 7763 64-Core Processor; laptop: AMD Ryzen AI 9 HX PRO 370; JDK 25

Figure 2. Time per row in a stage and outside one, for cheap and mixed projections of 50 and 99 columns, on a GitHub-hosted runner and a laptop. Above the dashed line the stage is slower.

On stock Spark 4.2.0 the same 99 cheap columns are 1.30, 1.45 and 1.37 times slower in a stage with JDK 17, 21 and 25. Nothing in the plan says so. It is one stage, as it should be, and explain("codegen") puts the stage's largest method at 3,670 bytes, under every limit the first post described:

SQL> select id + 1 AS c1, id + 2 AS c2, ..., id + 99 AS c99 from t
*(1) Project [(id#0L + 1) AS c1#49L, (id#0L + 2) AS c2#50L, (id#0L + 3) AS c3#51L, ...
+- *(1) Range (0, 1000, step=1, splits=1)

stage 1: maxMethodCodeSize 3670 bytes

Why. The JIT compiler that makes Java fast, C2, inlines small methods into the method that calls them, and it has a budget for how much: about 8000 bytes, counting the caller's own bytecode. A stage's method writes each output column through a small helper. With 25 columns the method is under a thousand bytes and C2 inlines every write. With 50 it is 1,861 bytes, and C2's log says, fifty times, failed to inline: size > DesiredMethodLimit: the budget is spent, and every row pays fifty real calls. Outside a stage, Spark splits the same projection into small methods, each of which has room to inline its writes. A column that does real work - builds a string, shifts a date - costs more than the call it leaves behind, which is why the mixed projection still wins in a stage. How much the calls cost depends on the processor: on the Zen 5 chips we measured the loss is small, on Zen 3 and Zen 4 servers it is half again.

What to do. For a query that is a wide projection of cheap columns, lower spark.sql.codegen.maxFields below the projection's width for that query. The projection leaves its stage, and the rest of the plan stays in one:

SQL> set spark.sql.codegen.maxFields=98;
SQL> select id + 1 AS c1, id + 2 AS c2, ..., id + 99 AS c99 from t
Project [(id#0L + 1) AS c1#446L, (id#0L + 2) AS c2#447L, (id#0L + 3) AS c3#448L, ...
+- *(1) Range (0, 1000, step=1, splits=1)

On the runner that took 99 cheap columns from 800 to 521 nanoseconds a row. Do not do it for a projection that computes; there the stage wins.

And do not raise maxFields to keep a wide projection in. The first post said to raise it with care and check the method's size. The method's size is not the risk. A 150-column projection of cheap columns, let into one stage by maxFields=200, has a method of 5,747 bytes, well under any limit - and for its first 12 to 22 seconds it runs two and a half to three and a half times slower than outside a stage while C2 compiles it, on every JDK, and then stays 1.2 to 1.8 times slower on the runners (on the laptop, 0.7 to 1.1 times, depending on the run). Every executor pays those seconds again for every stage it compiles.

The first minute of a new wide stageeach query's time, back to back: 150 cheap columns, in a stage and out of it500 ms1000 ms1500 ms2000 ms0 s10 s20 s30 s40 s50 s60 sseconds since the stage's first queryrunner, JDK 17, in a stagerunner, JDK 17, outsidelaptop, JDK 25, in a stagelaptop, JDK 25, outsideC2's code arrives 20 s in on the runnerand 13 s in on the laptop; untilthen the stage runs C1's coderunner: AMD EPYC 7763 64-Core Processor, stock Spark 4.2.0; laptop: AMD Ryzen AI 9 HX PRO 370, a build of Spark master

Figure 3. Each query's time, back to back for a minute, for the 150-column projection in a stage and outside one.

See it yourself: compile_wait.scala prints every query's time and when it settled.

3. Hidden in the log

Some give-ups are logged, at INFO, and spark-shell shows WARN and above. The method too long to be JIT compiled, from the first post. A common subexpression or an aggregate whose split functions would need more than the JVM's 255 parameter slots, so Spark keeps them in one method:

INFO CodegenContext: Failed to split subexpression code into small functions because the
  parameter length of at least one split function went over the JVM limit: 255

And the aggregate whose fast hash map Spark did not generate, because a key such as a struct is one it does not support. To see them, add two loggers to conf/log4j2.properties:

logger.codegen.name = org.apache.spark.sql.catalyst.expressions.codegen
logger.codegen.level = info
logger.hashagg.name = org.apache.spark.sql.execution.aggregate.HashAggregateExec
logger.hashagg.level = info

The shell then shows those lines and no other INFO line of Spark's, apart from one Code generated in ... ms for each class it compiles. From Spark 4.3.0 the method-too-long line comes from CodeCompiler rather than CodeGenerator, in the same package, so the setting above covers both; from 4.4.0 the first such method is a warning that names the remedy.

See it yourself: log_level.scala with its two lines.

4. Logged, with a fallback

Past 64 KB a method does not compile at all, and Spark has two fallbacks. A stage that fails logs an ERROR, then

WARN WholeStageCodegenExec: Whole-stage codegen disabled for plan (id=1):

and runs its operators one by one, each still compiled - what the first post measured as an eighth to a fifth slower on its shape. A CASE WHEN of 3,000 branches does it with the default code generation settings. Do not look for it in the plan, which still shows the stage it did not run:

SQL> select CASE WHEN id = 1 THEN id * 1 ... (3000 branches) ELSE 0 END AS v from range(0, 10)
*(1) Project [CASE WHEN (id#745L = 1) THEN id#745L WHEN (id#745L = 2) THEN (id#745L * 2) ...
+- *(1) Range (0, 10, step=1, splits=1)

EXPLAIN shows the plan Spark made, not how it ran it; here the WARN line is the only sign. A projection outside a stage that fails is evaluated by the interpreter instead:

WARN UnsafeProjection: Expr codegen error and falling back to interpreter mode

That costs more: 3.1 to 4.4 times the compiled projection's time on the runner, 3.1 to 4.2 on the laptop, from 100 to 1000 columns. You are unlikely to hit it with the defaults, since Spark splits a projection's code into methods well before 64 KB; it takes one enormous expression, or method splitting turned off.

5. Failing the query

About thirty places in Spark call a code generator directly, with no fallback, and there a compile failure is the query's. ORDER BY ... LIMIT is one: the top-k operator generates its ordering itself. An ordering over a 1,200-branch CASE WHEN, with method splitting off, fails:

org.codehaus.commons.compiler.InternalCompilerException: ... Code grows beyond 64 KB

Again this takes an expression far larger than most queries have. The other case on our list, a compile on the executor with no fallback, we could not provoke at all; we mention it because it exists.

6. Which release changes what

Read from the tracker on 2 October 2026.

ticket what changes for you state
SPARK-37019 transform, filter, exists, forall and aggregate generate code, so they keep their operator in its stage fixed in 4.3.0
SPARK-32750 a sort aggregate with grouping keys, such as a max of a string per group, runs in a stage fixed in 4.3.0
SPARK-59774 the method too long to be JIT compiled becomes a warning, once, naming the remedy fixed in 4.4.0
SPARK-33301 a large CASE WHEN in a stage is split into methods instead of passing 64 KB (section 4) open; in review as apache/spark#59069

How you'd know, across your own history

A missing *(n) is easy to see in one query and impossible to see across thousands. But every Spark event log (spark.eventLog.enabled=true) records each query's final plan, its operators' names and strings, and the settings the session changed, which is all it takes to classify every operator the way Spark did, offline. A short script does it:

python3 varka_codegen_report.py /path/to/spark-events/

It prints, for all the applications in the logs, how many operators ran in a stage and, for the rest, why, with the operators and functions most often to blame - the table this post's census is built from. Over the event log of the golden-file suites it matched the census within 2% on every reason. It is dev/varka_codegen_report.py, plain Python with no dependencies.


None of this is a bug in the usual sense: each case is a decision Spark makes on purpose, mostly the right one. What is missing is the telling. A per-operator mark in the SQL metrics, saying an operator ran outside whole-stage codegen and why, would turn section 1 into something the UI shows; a count of the fallbacks beside the existing codegen metrics would do the same for section 4. Neither exists yet.

How this was measured. Every number here is from a results file committed beside the post. The census and the classifier ran on the golden-file, TPC-DS, TPC-H and SSB suites of a build of Spark master from October 2026, which has the 4.3.0 changes of section 6; on 4.2 the higher-order functions and the keyed sort aggregates would add to its silent share. The costs are from four Spark-style benchmarks on the same build, run on GitHub-hosted runners (an AMD EPYC 7763 and 9V74 and an Intel Xeon 8573C) and on a laptop with an AMD Ryzen AI 9 HX PRO 370, with JDK 25; the snippets' outputs are from stock Spark 4.2.0 with JDK 17, 21 and 25, on runners that drew an EPYC 7763, 9V74, 9V45 and a Xeon 8573C. Every comparison is between numbers from the same run. The prose rounds; the results files and the snippets' outputs carry the exact values.