Apache Spark courseLesson 6 of 9
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.
On this page
- Sample data
- The Catalyst optimiser
- What it is
- Seeing the phases
- Other explain modes
- Pitfalls
- In interviews
- The Tungsten execution engine
- What it is
- Where you see Tungsten
- In interviews
- Whole-stage code generation
- What it is
- Looking at the generated code
- When code generation is skipped
- In interviews
- Cost-based optimisation
- What it is
- Collecting and using statistics
- Statistics versus AQE
- Pitfalls
- In interviews
- Predicate and projection pushdown
- What it is
- Reading the scan node
- What blocks pushdown
- In interviews
- Column pruning
- What it is
- What defeats it
- In interviews
- SQL versus DataFrame performance
- Same optimiser, same plan
- What actually changes performance
- In interviews
- Practice questions
- 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 oneselect, 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:
- 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. - 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.
- 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
UnsafeRowWriterand column vector classes inexplain("codegen")output. - Off-heap memory settings:
spark.memory.offHeap.enabled(defaultfalse) andspark.memory.offHeap.sizelet 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 (
ColumnarToRowin 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 abovespark.sql.codegen.maxFields(100 by default). Queries over hundreds of columns can lose codegen. - You disable it:
spark.sql.codegen.wholeStage=false(defaulttrue), 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
WHEREclause 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")) == 2026orupper(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.filterPushdownturned off (it defaults totrue).
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("*")ordf.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 TABLEstatistics and mainly helps join ordering. - Confirm pushdown and pruning in the scan node:
PartitionFilters,PushedFiltersandReadSchema. - SQL and DataFrame code give the same plan; Python UDFs, typed lambdas and RDDs are what the optimiser cannot see into.
Progress is saved in this browser only. No account needed.

