Project Tungsten is a set of changes that made Spark's execution engine work closer to the hardware. Whole-stage code generation is the part that turns a chain of operators into one tight piece of Java code. These two are a big reason DataFrames run so much faster than RDDs of Python objects.
What Tungsten does
- Stores rows in a compact binary format, in memory that the JVM garbage collector does not have to track (often off-heap). Millions of rows do not become millions of small Java objects.
- Lays out and processes data with the CPU cache in mind.
- Sorts and hashes directly on the binary data, without turning it back into objects first.
Whole-stage code generation
Without it, each operator in a plan (scan, filter, project, aggregate) is a separate step that passes rows to the next through a function call per row. With it, Spark compiles filter, project and the first part of the aggregate into a single function with a simple loop. The JVM optimizes that loop well, and per-row call overhead disappears.
You can see it in the plan:
df.filter("amount > 0").select("order_id").explain()
# *(1) Project [order_id]
# +- *(1) Filter (amount > 0)
# +- *(1) FileScan parquet ...The *(1) means these three operators were fused into code-generation stage 1. If a step has no star, it was not fused.
Where it stops
A Python UDF cannot be compiled into that Java loop. Spark must hand rows to a Python process, so the plan has a break there, and the fused pipeline is cut in two. The same happens with RDD lambdas, which Spark treats as black boxes. That is one more reason to prefer built-in functions.
The interview-sized answer: Tungsten is the memory and CPU efficiency work, code generation is how operators get fused, and both only work when Spark understands your operations, which means DataFrames and SQL.