Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Project Tungsten and whole-stage codegen

PySpark · Execution Model

Project Tungsten and whole-stage codegen

Hardpyspark-51
tungstenwhole-stage-codegencatalystperformance

Question

What is Project Tungsten, and what does whole-stage code generation do?

Solution

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.

🎯 Put this concept into practice

Solidify this answer with real hands-on interview drills in the browser studio.

Open related drill →
PreviousNext