Menu

Apache Spark course · Lesson 6 of 9

Catalyst, Tungsten and Code Generation: How Spark SQL Optimises Queries

How Spark SQL optimises a query: Catalyst's plan phases, Tungsten memory and whole-stage code generation, cost-based optimisation, pushdown, pruning and UDF costs.

  • Advanced
  • 17 min read
  • Updated Oct 2026
On this page
  1. Sample data
  2. The Catalyst optimiser
  3. What it is
  4. Seeing the phases
  5. Other explain modes
  6. Pitfalls
  7. In interviews
  8. The Tungsten execution engine
  9. What it is
  10. Where you see Tungsten
  11. In interviews
  12. Whole-stage code generation
  13. What it is
  14. Looking at the generated code
  15. When code generation is skipped
  16. In interviews
  17. Cost-based optimisation
  18. What it is
  19. Collecting and using statistics
  20. Statistics versus AQE
  21. Pitfalls
  22. In interviews
  23. Predicate and projection pushdown
  24. What it is
  25. Reading the scan node
  26. What blocks pushdown
  27. In interviews
  28. Column pruning
  29. What it is
  30. What defeats it
  31. In interviews
  32. SQL versus DataFrame performance
  33. Same optimiser, same plan
  34. What actually changes performance
  35. In interviews
  36. Practice questions
  37. Key takeaways

When you write df.filter(...).groupBy(...) or the equivalent SQL, you describe what you want. Spark SQL’s optimiser, Catalyst, decides how: it rewrites the query, picks physical operators, and hands the result to Tungsten, the execution layer that generates compact Java code for it. Knowing what these layers do tells you which code styles Spark can optimise, which it cannot, and how to prove from a plan that an optimisation happened.

Sample data

Orders written as Parquet partitioned by date, read back as a DataFrame. show_plan prints a plan with the temporary path shortened.

import contextlib, io, os, tempfile
from pyspark.sql import SparkSession, functions as F

spark = (SparkSession.builder.master("local[2]").appName("catalyst-tungsten")
         .config("spark.sql.shuffle.partitions", "4")
         .config("spark.sql.warehouse.dir", tempfile.mkdtemp())
         .getOrCreate())

base = tempfile.mkdtemp()
path = os.path.join(base, "orders")

(spark.range(0, 100_000, numPartitions=4)
     .select(F.col("id").alias("order_id"),
             (F.col("id") % 1000).alias("customer_id"),
             ((F.col("id") * 7) % 500).cast("double").alias("amount"),
             F.when(F.col("id") % 4 == 0, "web").otherwise("store").alias("channel"),
             F.expr("date_add(DATE'2026-01-01', cast(id % 30 as int))").alias("order_date"))
     .write.partitionBy("order_date").parquet(path))

orders = spark.read.parquet(path)

def show_plan(df, mode=None):
    buf = io.StringIO()
    with contextlib.redirect_stdout(buf):
        df.explain(mode) if mode else df.explain()
    print(buf.getvalue().replace("file:" + base, "<data>").rstrip())

The Catalyst optimiser

What it is

Catalyst is the query optimiser inside Spark SQL. Every DataFrame, Dataset and SQL query becomes a tree of operators, and Catalyst transforms that tree in four phases:

Phase Input → output What happens
Parsing SQL text or DataFrame calls → unresolved logical plan Syntax only; table and column names are not checked yet
Analysis → analysed logical plan Names are resolved against the catalog and the DataFrame schemas; types are checked and implicit casts added. Misspelt columns fail here
Logical optimisation → optimised logical plan Rule-based rewrites: predicate pushdown, column pruning, constant folding, null-filter inference, combining filters and projections, simplifying expressions, rewriting subqueries as joins
Physical planning → physical plan Chooses operators (broadcast hash join versus sort-merge join, hash versus sort aggregate), inserts Exchange and Sort nodes where the data distribution requires them, using size estimates

After physical planning, Tungsten generates code, and Adaptive Query Execution can revise the physical plan at run time.

Seeing the phases

explain("extended") prints all four. Watch what changes between them:

q = (orders.filter((F.col("amount") > 100 + 50) & (F.col("channel") == "web"))
           .select("order_id", "amount"))
show_plan(q, "extended")
== Parsed Logical Plan ==
'Project ['order_id, 'amount]
+- Filter ((amount#8 > cast(150 as double)) AND (channel#9 = web))
   +- Relation [order_id#6L,customer_id#7L,amount#8,channel#9,order_date#10] parquet

== Analyzed Logical Plan ==
order_id: bigint, amount: double
Project [order_id#6L, amount#8]
+- Filter ((amount#8 > cast(150 as double)) AND (channel#9 = web))
   +- Relation [order_id#6L,customer_id#7L,amount#8,channel#9,order_date#10] parquet

== Optimized Logical Plan ==
Project [order_id#6L, amount#8]
+- Filter ((isnotnull(amount#8) AND isnotnull(channel#9)) AND ((amount#8 > 150.0) AND (channel#9 = web)))
   +- Relation [order_id#6L,customer_id#7L,amount#8,channel#9,order_date#10] parquet

== Physical Plan ==
*(1) Project [order_id#6L, amount#8]
+- *(1) Filter (((isnotnull(amount#8) AND isnotnull(channel#9)) AND (amount#8 > 150.0)) AND (channel#9 = web))
   +- *(1) ColumnarToRow
      +- FileScan parquet [order_id#6L,amount#8,channel#9,order_date#10] Batched: true, DataFilters: [isnotnull(amount#8), isnotnull(channel#9), (amount#8 > 150.0), (channel#9 = web)], Format: Parquet, Location: InMemoryFileIndex(1 paths)[<data>/orders], PartitionFilters: [], PushedFilters: [IsNotNull(amount), IsNotNull(channel), GreaterThan(amount,150.0), EqualTo(channel,web)], ReadSchema: struct<order_id:bigint,amount:double,channel:string>

Between the analysed and optimised plans, Catalyst folded the constant cast(150 as double) into 150.0 and inferred isnotnull filters, since a NULL can never satisfy > 150. In the physical plan, the filters appear in PushedFilters (sent to the Parquet reader) and ReadSchema lists only three of the five columns. The *(1) marks operators fused by whole-stage code generation.

Other explain modes

Mode Shows
explain() or "simple" Physical plan only
"extended" All four plans
"formatted" Physical plan outline plus a numbered details section per operator (easiest to read for big plans)
"cost" Optimised logical plan with size and row-count statistics
"codegen" The generated Java source per code-generation stage

Pitfalls

  • Catalyst cannot see inside a Python UDF, a Python lambda on an RDD, or an opaque Scala closure in a typed Dataset.map. Filters after those cannot be pushed below them.
  • Analysis errors appear immediately (a wrong column name fails at the select), but runtime errors (a bad cast under ANSI mode) appear only at the action.
  • Plans with thousands of operators, such as from long loops of withColumn, make analysis and optimisation slow. Build columns in one select, or checkpoint long iterative plans.

In interviews

“What does Catalyst do?” Name the four phases, give two or three rule examples (predicate pushdown, column pruning, constant folding), and say that physical planning chooses join strategies using statistics. Bonus: mention that Catalyst rules are pattern-matching functions over immutable trees, applied in batches until the plan stops changing.

The Tungsten execution engine

What it is

Project Tungsten is the set of execution-layer improvements that make Spark use CPU and memory efficiently. Its three main ideas:

  1. Explicit binary memory management. Rows are stored in a compact binary format (UnsafeRow) rather than as Java objects, and Spark manages memory pages itself, on or off the JVM heap. This cuts object overhead and garbage collection, and lets sorting and hashing work directly on bytes.
  2. Cache-aware algorithms. Sort and hash-aggregation structures are laid out so that the CPU cache is used well, for example sorting on a compact key prefix stored next to a pointer.
  3. Code generation. Instead of interpreting a tree of expression objects row by row, Spark generates Java code specialised to the query and compiles it at run time (with the Janino compiler).

These are why DataFrame operations often beat hand-written RDD code: an RDD of Python or Java objects gets none of this.

Where you see Tungsten

  • Generated-code stages marked *(n) in plans.
  • The UnsafeRowWriter and column vector classes in explain("codegen") output.
  • Off-heap memory settings: spark.memory.offHeap.enabled (default false) and spark.memory.offHeap.size let Tungsten allocate execution and storage memory outside the JVM heap, which reduces GC work for large executors.
  • Columnar in-memory caching and vectorised Parquet and ORC readers (ColumnarToRow in plans), which read batches of column values rather than one row at a time.

In interviews

“What is Tungsten?” Say: Spark’s execution-engine layer for CPU and memory efficiency, with binary row formats and self-managed memory (less GC), cache-aware sorting and hashing, and whole-stage code generation. Contrast with Catalyst: Catalyst decides what plan to run; Tungsten makes that plan run fast.

Whole-stage code generation

What it is

Whole-stage code generation fuses a chain of operators inside one stage (scan, filter, project, partial aggregate) into a single generated Java function with a tight loop. Instead of each operator calling the next one per row through virtual function calls (the “Volcano” iterator model), the fused code keeps values in local variables and CPU registers, the way hand-written code would.

In plans, each fused group gets a number: *(1), *(2). Operators without the star (an Exchange, a Python UDF evaluation, some complex expressions) break the chain.

Looking at the generated code

explain("codegen") prints the source of each fused stage. Trimmed, the inner loop of a filter on amount > 400 reads like hand-written Java:

show_plan(orders.filter("amount > 400").select("order_id", "amount"), "codegen")
Found 1 WholeStageCodegen subtrees.
== Subtree 1 / 1 (maxMethodCodeSize:338; maxConstantPoolSize:173(0.26% used); numInnerClasses:0) ==
*(1) Project [order_id#6L, amount#8]
+- *(1) Filter (isnotnull(amount#8) AND (amount#8 > 400.0))
   +- *(1) ColumnarToRow
      +- FileScan parquet [order_id#6L,amount#8,order_date#10] Batched: true, DataFilters: [isnotnull(amount#8), (amount#8 > 400.0)], Format: Parquet, Location: InMemoryFileIndex(1 paths)[<data>/orders], PartitionFilters: [], PushedFilters: [IsNotNull(amount), GreaterThan(amount,400.0)], ReadSchema: struct<order_id:bigint,amount:double>

Generated code:
...
/* 031 */   protected void processNext() throws java.io.IOException {
/* 032 */     if (columnartorow_mutableStateArray_1[0] == null) {
/* 033 */       columnartorow_nextBatch_0();
/* 034 */     }
/* 035 */     while ( columnartorow_mutableStateArray_1[0] != null) {
/* 036 */       int columnartorow_numRows_0 = columnartorow_mutableStateArray_1[0].numRows();
/* 037 */       int columnartorow_localEnd_0 = columnartorow_numRows_0 - columnartorow_batchIdx_0;
/* 038 */       for (int columnartorow_localIdx_0 = 0; columnartorow_localIdx_0 < columnartorow_localEnd_0; columnartorow_localIdx_0++) {
/* 039 */         int columnartorow_rowIdx_0 = columnartorow_batchIdx_0 + columnartorow_localIdx_0;
/* 040 */         do {
/* 041 */           boolean columnartorow_isNull_1 = columnartorow_mutableStateArray_2[1].isNullAt(columnartorow_rowIdx_0);
/* 042 */           double columnartorow_value_1 = columnartorow_isNull_1 ? -1.0 : (columnartorow_mutableStateArray_2[1].getDouble(columnartorow_rowIdx_0));
/* 043 */
/* 044 */           boolean filter_value_2 = !columnartorow_isNull_1;
/* 045 */           if (!filter_value_2) continue;
/* 046 */
/* 047 */           boolean filter_value_3 = false;
/* 048 */           filter_value_3 = org.apache.spark.sql.catalyst.util.SQLOrderingUtil.compareDoubles(columnartorow_value_1, 400.0D) > 0;
/* 049 */           if (!filter_value_3) continue;
/* 050 */
/* 051 */           ((org.apache.spark.sql.execution.metric.SQLMetric) references[2] /* numOutputRows */).add(1);
...

(Trimmed to the summary, the fused plan and the inner loop; the full output is a few hundred lines.) The filter is two if (...) continue; checks inside a for loop over a column batch, with no per-row function calls.

When code generation is skipped

  • An operator does not support it (some generators, Python UDFs, certain aggregate functions). The plan shows that operator without a *.
  • The generated method would be too large. Spark falls back to interpreted execution when a method exceeds spark.sql.codegen.hugeMethodLimit, and stops fusing very wide schemas above spark.sql.codegen.maxFields (100 by default). Queries over hundreds of columns can lose codegen.
  • You disable it: spark.sql.codegen.wholeStage=false (default true), useful only for debugging.

In interviews

“What is whole-stage code generation?” Explain operator fusion into one generated function per stage, the *(n) marker in plans, and why it is faster (no virtual calls per row, data stays in registers, compiler-friendly loops). Mention that Python UDFs and very wide schemas break or disable it.

Cost-based optimisation

What it is

Catalyst’s rules do not need statistics. Cost-based optimisation (CBO) adds decisions that do: estimating how many rows each operator produces and choosing join order and join strategy accordingly. Statistics come from ANALYZE TABLE and are stored in the catalog.

Setting Default in 4.2 Effect
spark.sql.cbo.enabled false Use column statistics for row-count estimates
spark.sql.cbo.joinReorder.enabled false Reorder multi-way inner joins by estimated cost
spark.sql.statistics.histogram.enabled false Collect histograms with column statistics for better selectivity estimates
spark.sql.cbo.starSchemaDetection false Detect fact and dimension tables to order star joins

Even with CBO off, Spark uses size in bytes (from file sizes or ANALYZE TABLE ... COMPUTE STATISTICS) to decide broadcast joins against spark.sql.autoBroadcastJoinThreshold.

Collecting and using statistics

orders.write.saveAsTable("orders_t")
query = "SELECT * FROM orders_t WHERE amount > 400"

def stats_lines(df):
    buf = io.StringIO()
    with contextlib.redirect_stdout(buf):
        df.explain("cost")
    return [line.split("Statistics")[1] for line in buf.getvalue().splitlines() if "Statistics(" in line]

print("no stats:   ", stats_lines(spark.sql(query)))
spark.sql("ANALYZE TABLE orders_t COMPUTE STATISTICS FOR COLUMNS amount, customer_id")
spark.conf.set("spark.sql.cbo.enabled", "true")
print("with stats: ", stats_lines(spark.sql(query)))

spark.sql("DESCRIBE EXTENDED orders_t amount").show()
spark.conf.set("spark.sql.cbo.enabled", "false")
no stats:    ['(sizeInBytes=466.5 KiB)', '(sizeInBytes=466.5 KiB)']
with stats:  ['(sizeInBytes=1085.0 KiB, rowCount=1.98E+4)', '(sizeInBytes=5.3 MiB, rowCount=1.00E+5)']
+--------------+----------+
|     info_name|info_value|
+--------------+----------+
|      col_name|    amount|
|     data_type|    double|
|       comment|      NULL|
|           min|       0.0|
|           max|     499.0|
|     num_nulls|         0|
|distinct_count|       510|
|   avg_col_len|         8|
|   max_col_len|         8|
|     histogram|      NULL|
+--------------+----------+

Without statistics, the filter’s estimate is just the file size (466.5 KiB here). With column statistics and CBO on, Spark knows amount ranges from 0 to 499 and estimates that amount > 400 keeps about a fifth of 100,000 rows: 19,800, which is exactly right here because the values are uniform. Note that the byte sizes change too: with row counts available, Spark estimates size as rows times an estimated row width (5.3 MiB for the table) instead of using the compressed file size, so the numbers are estimates for planning, not measurements. On skewed data, histograms improve such estimates.

Statistics versus AQE

AQE (on by default since 3.2) measures real sizes at each shuffle and can switch to a broadcast join or coalesce partitions at run time, which covers many cases CBO was designed for. CBO still matters for decisions made before the first shuffle, chiefly the order of multi-way joins, which AQE does not change.

Pitfalls

  • Stale statistics are worse than none: a table analysed when small and since grown can be broadcast and run out of memory. Re-analyse after big loads, or rely on AQE.
  • Statistics live in the catalog, so they only exist for catalog tables, not for paths read directly with spark.read.parquet(path).
  • Collecting column statistics scans the table; schedule it.

In interviews

“Does Spark have a cost-based optimiser?” Yes, off by default; it uses ANALYZE TABLE statistics for cardinality estimates, join reordering and join strategy. Explain how AQE overlaps with it, and the risk of stale statistics.

Predicate and projection pushdown

What it is

Pushdown moves work into the data source so less data reaches Spark at all.

  • Predicate pushdown: filters are handed to the reader. For Parquet and ORC, the reader uses per-row-group min and max statistics to skip whole row groups; for JDBC, the filter becomes a SQL WHERE clause run by the database.
  • Projection pushdown: only the needed columns are requested from the source (the source-side half of column pruning).
  • Partition pruning: filters on partition columns skip entire directories before any file is opened.

Reading the scan node

A filter on the partition column and on a data column, side by side:

q2 = orders.filter("order_date = DATE'2026-01-05' AND amount > 400").select("customer_id")
show_plan(q2)
== Physical Plan ==
*(1) Project [customer_id#7L]
+- *(1) Filter (isnotnull(amount#8) AND (amount#8 > 400.0))
   +- *(1) ColumnarToRow
      +- FileScan parquet [customer_id#7L,amount#8,order_date#10] Batched: true, DataFilters: [isnotnull(amount#8), (amount#8 > 400.0)], Format: Parquet, Location: InMemoryFileIndex(1 paths)[<data>/orders], PartitionFilters: [isnotnull(order_date#10), (order_date#10 = 2026-01-05)], PushedFilters: [IsNotNull(amount), GreaterThan(amount,400.0)], ReadSchema: struct<customer_id:bigint,amount:double>
Field Meaning here
PartitionFilters order_date = 2026-01-05: only one of 30 date folders is listed and read
PushedFilters GreaterThan(amount,400.0): Parquet can skip row groups whose max amount is at most 400
DataFilters Filters Spark still applies row by row after reading (pushed filters are advisory, so Spark re-checks them)
ReadSchema Only customer_id and amount are read from the files

What blocks pushdown

  • Functions on the column: F.year(col("ts")) == 2026 or upper(name) = 'X' usually cannot be pushed; rewrite as a range on the raw column (ts >= '2026-01-01' AND ts < '2027-01-01').
  • Python UDFs in the filter: Spark cannot translate them.
  • Casts that change the column’s type, such as comparing a string column with a number, can prevent pushdown and, under ANSI mode (the default since Spark 4.0), can also raise errors.
  • Sources without support: CSV and JSON readers support fewer pushdowns than Parquet and ORC; JDBC pushes filters but not all functions.
  • spark.sql.parquet.filterPushdown turned off (it defaults to true).

In interviews

“How do you check that a filter was pushed down?” Look at the scan node of the physical plan for PushedFilters and PartitionFilters, and at the SQL tab’s scan metrics (files and bytes read). Explain why year(ts) = 2026 may not push down and how to rewrite it.

Column pruning

What it is

Column pruning removes columns a query never uses, at every level of the plan, so they are never read, shuffled, cached or serialized. With columnar formats (Parquet, ORC, Delta) unread columns cost almost nothing, because each column is stored separately.

Compare the two ReadSchema values above: the first query read three columns, the second only two, even though both started from the same five-column DataFrame. Catalyst also prunes nested fields: selecting address.city from a struct reads only that field from Parquet (spark.sql.optimizer.nestedSchemaPruning.enabled, on by default).

What defeats it

  • select("*") or df.collect() of whole rows, when only a few columns are needed.
  • Caching a wide DataFrame and querying a few columns: the cache stores every column.
  • Row-based formats (CSV, JSON): the reader must still parse every line, though it can skip converting unused fields.
  • A Python UDF that takes a whole row (struct("*")) forces every column through serialization.

In interviews

Column pruning is often asked together with pushdown: “Why is Parquet faster than CSV for analytics?” Columnar storage enables column pruning, row-group statistics enable predicate pushdown, and compression per column reduces I/O.

SQL versus DataFrame performance

Same optimiser, same plan

SQL strings and DataFrame calls are two front ends to the same Catalyst pipeline, so equivalent queries produce the same physical plan and the same performance:

orders.createOrReplaceTempView("orders")

sql_df = spark.sql("""
    SELECT customer_id, SUM(amount) AS total
    FROM orders WHERE channel = 'web'
    GROUP BY customer_id""")

api_df = (orders.filter(F.col("channel") == "web")
                .groupBy("customer_id").agg(F.sum("amount").alias("total")))

import re
def physical(df):
    plan = df._jdf.queryExecution().executedPlan().toString()
    return re.sub(r"#\d+L?|plan_id=\d+", "", plan)

print("identical physical plans:", physical(sql_df) == physical(api_df))
identical physical plans: True

(Expression ids and plan ids are stripped before comparing, since they differ by construction.) Choose between SQL and the DataFrame API for readability, testing and reuse, not speed.

What actually changes performance

Code style Optimised by Catalyst? Code generation? Cost
SQL or DataFrame built-in functions Yes Yes Baseline
Scala/Java typed Dataset with lambdas Partly: lambdas are opaque Partly Object serialization for each lambda
Python UDF (Arrow-optimised) Opaque to the optimiser No, runs in Python workers Batches through Arrow to Python and back
pandas UDF Opaque No Vectorised Arrow batches; usually much faster than row UDFs
RDD with Python lambdas Not at all No Every row pickled to and from Python

The Python UDF cost is visible in a plan. In this Spark 4.2 run, spark.sql.execution.pythonUDF.arrow.enabled reports true, so a plain @udf runs as ArrowEvalPython; a filter on the UDF result even evaluates the UDF twice:

@F.udf("double")
def add_vat(x):
    return x * 1.2

show_plan(orders.select((F.col("amount") * 1.2).alias("gross")))
print()
show_plan(orders.select(add_vat("amount").alias("gross")).filter("gross > 10"))
== Physical Plan ==
*(1) Project [(amount#8 * 1.2) AS gross#2594]
+- *(1) ColumnarToRow
   +- FileScan parquet [amount#8,order_date#10] Batched: true, DataFilters: [], Format: Parquet, Location: InMemoryFileIndex(1 paths)[<data>/orders], PartitionFilters: [], PushedFilters: [], ReadSchema: struct<amount:double>

== Physical Plan ==
*(3) Project [pythonUDF0#2598 AS gross#2596]
+- ArrowEvalPython [add_vat(amount#8)#2595], [pythonUDF0#2598], 101
   +- *(2) Project [amount#8]
      +- *(2) Filter (pythonUDF0#2597 > 10.0)
         +- ArrowEvalPython [add_vat(amount#8)#2595], [pythonUDF0#2597], 101
            +- *(1) Project [amount#8]
               +- *(1) ColumnarToRow
                  +- FileScan parquet [amount#8,order_date#10] Batched: true, DataFilters: [], Format: Parquet, Location: InMemoryFileIndex(1 paths)[<data>/orders], PartitionFilters: [], PushedFilters: [], ReadSchema: struct<amount:double>

The built-in expression is one fused codegen stage. The UDF version leaves the JVM for Python twice, once for the filter and once for the projection, and splits code generation into three stages. If a UDF is expensive or non-deterministic, compute it once in a column and cache or checkpoint, or better, replace it with built-in functions (see UDFs and safer alternatives).

In interviews

“Is Spark SQL faster than the DataFrame API?” No: both compile to the same optimised plan, which you can prove with explain(). The real performance gap is between built-in expressions and opaque code (Python UDFs, typed lambdas, RDDs), because the optimiser and code generator cannot see inside them.

Practice questions

Name the four phases of Catalyst and what can fail in each.

Parsing turns SQL or API calls into an unresolved logical plan (syntax errors fail here). Analysis resolves names and types against the catalog (unknown tables or columns, type mismatches fail here). Logical optimisation applies rule-based rewrites such as pushdown, pruning and constant folding (it should never change results). Physical planning chooses operators and join strategies and adds exchanges. Data-dependent errors, such as ANSI cast failures, appear only when the plan executes.

What is the difference between Catalyst and Tungsten?

Catalyst is the optimiser: it decides the plan, rewriting the query and choosing physical operators. Tungsten is the execution engine work that runs the plan efficiently: binary row formats and explicit memory management to reduce object overhead and GC, cache-aware sorting and hashing, and whole-stage code generation.

A query filters with year(event_ts) = 2026 on a large Parquet table and reads every row group. Why, and how do you fix it?

Wrapping the column in a function usually prevents the filter from being pushed to the Parquet reader as a simple comparison, so row-group statistics cannot be used. Rewrite as a range on the raw column: event_ts >= '2026-01-01' AND event_ts < '2027-01-01', then confirm the comparison appears in PushedFilters. If the table is partitioned by a date column, filter on that column too so PartitionFilters prunes directories.

When does the cost-based optimiser help, given that AQE exists?

CBO works before execution using catalog statistics, so it can choose the order of multi-way joins and a join strategy before any shuffle runs, which AQE cannot change. AQE corrects decisions at shuffle boundaries using measured sizes. CBO is off by default; enabling it requires ANALYZE TABLE ... FOR COLUMNS and keeping statistics fresh.

How can you tell from a plan that whole-stage code generation is used, and what breaks it?

Operators inside a fused stage are prefixed with *(n), and explain("codegen") lists the generated subtrees. Operators without the star break the chain, such as Python UDF evaluation (ArrowEvalPython or BatchEvalPython) and exchanges; very wide schemas beyond spark.sql.codegen.maxFields or methods beyond spark.sql.codegen.hugeMethodLimit fall back to non-fused execution.

A teammate rewrites a DataFrame pipeline in SQL to make it faster. What do you tell them?

SQL and the DataFrame API go through the same Catalyst optimiser and produce the same physical plan, so the rewrite alone will not change speed; compare explain() output to prove it. Look instead for opaque code (Python UDFs, RDD conversions), missing pushdown or pruning, unnecessary shuffles and skew.

Key takeaways

  • Catalyst turns queries into parsed, analysed, optimised and physical plans; explain("extended") shows all four.
  • Tungsten runs plans efficiently with binary rows, managed memory and whole-stage code generation, marked *(n) in plans.
  • The cost-based optimiser is off by default; it needs ANALYZE TABLE statistics and mainly helps join ordering.
  • Confirm pushdown and pruning in the scan node: PartitionFilters, PushedFilters and ReadSchema.
  • SQL and DataFrame code give the same plan; Python UDFs, typed lambdas and RDDs are what the optimiser cannot see into.

By DataDank Editorial · Last reviewed Oct 2026 · All examples run on PySpark 4.2.0 in local mode (local[2]) against small local Parquet files and an in-memory catalog. Temporary file paths in plans are replaced with <data> for readability; expression ids such as #8 differ between runs.

Progress is saved in this browser only. No account needed.

Search
Filter by type