Ch.4: How Spark Tungsten Turns Plans Into Generated Code

Outline

Transcript

0:00 Spark can make the same cluster run faster without adding machines. Not by moving the work somewhere else, but by making the loop inside each executor cheaper. Spark 2.0's release notes reported two-to-10-times speedups in common SQL and DataFrame operator benchmarks around whole-stage code generation. The surprising part is that this was not a bigger-cluster story. Same cluster, cheaper loop? Exactly. That is the Tungsten story people miss. Once Catalyst has picked a decent plan, the expensive part isn't only the cluster shape.

0:27 It's the row-processing loop inside each task. OK. So less buying horsepower, more asking what one core stops wasting. The easy mental model is: Spark is slow, so add machines. More executors, more cores, bigger nodes. Yeah. We've all opened that tab before checking the plan. Classic. Right. But after Catalyst has chosen a reasonable physical plan, a different bottleneck shows up. One executor core is reading rows, checking predicates, computing expressions, building hash tables, and updating aggregates millions or billions of times.

0:59 And if that inner loop is full of object allocation, virtual calls, pointer chasing, and garbage collection pressure, more machines can hide the pain, but they don't remove it. It's basically the infrastructure equivalent of trying to solve a traffic jam by making the cars bigger. Tungsten was Spark's answer to that problem. The previous chapter was about Catalyst. It asked: what equivalent plan should Spark run? Tungsten asks a lower-level question: once we know the plan, how do we run it efficiently on real hardware?

1:28 So this is the handoff. Catalyst gives Spark an optimized logical plan and then a physical plan. Tungsten and whole-stage codegen make that physical plan feel less like an interpreted tree of operators and more like a compiled query pipeline. Catalyst decides what; Tungsten decides how. Same query, same answer, much less waste in the hot path. OK, walk me through the cost. Where does the JVM actually start hurting analytical work? The old danger isn't that the JVM is bad. The danger is that analytical processing punishes generic object shapes.

2:00 Meaning? Imagine every row as a little graph of objects. A row object points to field objects. Strings live somewhere else. Nullable values need wrappers. Every operation follows pointers, checks types, allocates temporary objects, and leaves cleanup work for the garbage collector. And then you do that for a billion rows. Oh, wow. Exactly. Honestly, none of those costs is dramatic once. But in a query engine, the tiny cost is the workload. I didn't realize this when I first learned it. The cheap annoyance becomes the bill because it's one tiny cost repeated until it dominates the loop.

2:34 Yeah. So how does Tungsten stop that bill from running up? The most concrete Tungsten idea is the binary row format. Spark's internal `UnsafeRow` stores a row in a compact memory layout: null bits, fixed-width values, and a variable-length region for things like strings. The name isn't exactly comforting. No. It sounds like something you find during an incident review. But the point is precise: Spark can treat a row as bytes in a known layout instead of a pile of ordinary JVM objects. That framing works. I think of generic JVM objects as a poorly organized warehouse.

3:10 The CPU is the worker running around to assemble a product aisle 4 for a string, aisle 12 for an integer, aisle 50 for a Boolean. I love that image. Exhausting and slow. Right. UnsafeRow is like ditching the warehouse entirely. Perfectly packed kits on a predictable conveyor belt that feeds the CPU cache hierarchy more cleanly. Yeah. Less object overhead, better cache behavior, fewer tiny allocations, and less garbage collector work. It doesn't make memory free. It makes the common row path more predictable.

3:42 That predictability matters because Spark jobs already have enough memory pressure. Joins build hash tables. Aggregations keep state. Sorts spill. Shuffles materialize data. Right. Tungsten doesn't erase those structures it stops the row loop from piling extra trash on top of them. That distinction matters. A slow join may still be slow because the data is huge. The engine just shouldn't make every row heavier than it needs to be. Yeah. The win isn't glamorous. It's thousands of small inefficiencies that stop happening inside the loop.

4:09 Right. That's where query engines make money: they stop making the ordinary path pay a hidden tax on every row. There is another piece that often gets mixed into this story: vectorized reads and columnar batches. Vectorized just means the operator processes a batch of column values at a time instead of one row at a time. Same arithmetic, fewer trips through the dispatcher. This is the Parquet path, right? Parquet is the common example column batches arrive friendlier to CPU caches and compression.

4:39 So the scan path is already trying to feed the engine in a shape the CPU likes. Exactly. The chapter isn't "rows versus columns, pick one." Spark uses different internal shapes at different points. The important idea is that execution speed depends on data layout, not just data volume. Oh, right. The shape of the bytes is already part of the optimization. Before whole-stage codegen, many query engines used what people call the volcano model. Each operator exposes a method like `next`, and the parent asks the child for one row at a time.

5:11 Filter asks scan for a row. Project asks filter for a row. Aggregate asks project for a row, and— —somewhere in there, an actual computation happens. Correct and beautifully modular. Also chatty. Very chatty. Every row pays the abstraction tax twice once on the way in, once on the way out. The CPU is doing useful work, but it's also paying for the shape of the abstraction. That is the uncomfortable engineering lesson here. The abstraction that makes the engine clean can become the thing the hardware keeps tripping over.

5:42 So the fix has to remove the handoffs. Let's hear it. Whole-stage codegen is Spark saying: if a run of physical operators can be fused, generate one Java function for the whole run. So instead of scan calls filter calls project calls aggregate, Spark writes code that does the scan, predicate, expression, and update in one tight block. Sure. The operators still exist in the plan. But at runtime, the compatible stretch becomes generated Java, compiled by the JVM, and executed like ordinary code. Fewer virtual calls. Fewer row wrappers.

6:13 More local variables. Less interpreter-shaped overhead. The axis is interpreter versus compiled. It is the difference between a relay race where every runner stops to hand off a baton and one runner sprinting the whole leg. The speed comes from removing handoffs, not from a faster runner. That explains why the Spark 2.0 release notes could tie benchmark speedups like that to whole-stage code generation. Hold on. Local execution changes, not more machines? Across what kinds of operators? Common SQL patterns: filters, projections, hash aggregates. The cluster didn't magically get bigger.

6:46 The inner loop got cheaper. The generated code isn't mysterious if you squint at it the right way. It still looks like a printer had a bad day. Absolutely. It isn't written for humans. But conceptually it's simple: read a value, check null, evaluate the predicate, compute the expression, update the buffer, next row. That tracks. The payoff is that those steps sit next to each other in one compiled function. It's ugly because it's specialized, not because it's mysterious. Spark isn't asking each operator to interpret a generic row independently.

7:19 It's emitting the boring, specialized code that a human would never want to maintain by hand. The important limit is that whole-stage codegen fuses compatible local operators. It doesn't fuse through everything. Exchange boundaries break the pipeline because a shuffle materializes data between stages. Sorts can break it. Very complex generated methods hit bytecode limits and Spark falls back to interpreted execution. Oh, right. The boundary. Python UDFs are their own boundary. Right. Once you leave Spark's optimized SQL expression world, Spark can't, you know, just look inside that code and generate one neat Java loop around it.

7:56 So whole-stage codegen is powerful, but it isn't a spell that crosses every boundary in the job. OK. So how do you actually see whole-stage codegen happening in practice? The nice part is that Spark leaves fingerprints. In the physical plan and the SQL tab, `WholeStageCodegen` shows up as a wrapper around fused operators. And in Spark's plan views, the markers are the clue, not the decoration. When you see codegen IDs or starred nodes, they are pointing at operators participating in generated-code regions.

8:29 A run sharing the same ID is one fused stage. You can also find generated-code metrics when, The UI exposes them. The point isn't that you should read generated Java every day. Please don't make that your hobby. I have, twice. Wouldn't recommend it. The point is that you can answer a production question: did this query stay on the generated fast path, or did it fall back into a more fragmented execution shape? That one question is often enough to explain why two queries with similar logical plans feel completely different.

9:01 Right. This is the part I've shipped wrong more times than I want to admit. I've had a data scientist say, "It's one cleanup function. How expensive can it be?" And honestly, that sounds reasonable. This is why UDFs can hurt more than their source code suggests. The function itself may be simple. The boundary is expensive. Because Spark can't reason about the expression anymore? Exactly. A built-in SQL expression is visible to Catalyst and to codegen. Spark can push it, fold it, inline it, null-check it, and place it inside the generated loop.

9:34 A UDF is opaque. Yeah. It says: call this box. Maybe that box is fast. Maybe not. But the engine has lost a lot of its ability to reorganize work around it. And if the UDF crosses into Python, now serialization and process boundaries join the party. Right. Spark moves data to a Python worker, waits for the result, and then brings it back into its internal format. A tiny function can turn a clean expression into a serialization nightmare. It's worth saying the quiet part clearly. Generated code isn't, honestly, pure upside.

10:07 Right. I've heard platform engineers push back here: "Great, now the thing failing in production is code no one wrote." Fair complaint. Spark has to generate the code, compile it, cache it, and fall back when the generated method gets too large or an expression can't be supported. Debugging can also get weird, because the code that ran isn't the code you wrote. It's Java Spark emitted from a physical plan that came from an optimized plan that came from your query. Exactly. Reasonable architecture, terrible outage sentence. So the tradeoff is control versus speed: Spark spends complexity inside the engine so your query can stay declarative while the executor runs something specialized.

10:48 The practical habit is simple. Most of us reach for more cores first. When a Spark SQL job is slow, don't only ask how many executors it has. Ask what the executor is doing. Is it running a fused generated pipeline? Is it stuck behind a shuffle? Is it falling through a Python boundary? Exactly. That's the cheaper-loop question. Those are different problems. More cores might help one of them. It might do almost nothing for another. This is the same reason the previous chapter cared so much about plan shape.

11:17 A better plan gives Spark fewer bad choices. Tungsten and codegen make the good choices cheaper to execute. Spark only feels fast when the distributed layer and the local loop are both doing their job. So the short version is: Catalyst provides the map; Tungsten builds the race car. Man, that's a pretty elegant division of labor. Spark performance isn't just parallelism. It's memory layout, generated code, and whether the engine can stay on its optimized path. Coming next, we move from the local loop to one of the most expensive distributed taxes in Spark: shuffle costs.

11:50 Thanks for listening to Learning Podcasts.