Catalyst and physical plans
Read explain() output: parsed, analyzed and optimized plans, and the physical plan Spark runs.
On this page
Show code in
Every code block on the page follows this.
You will learn
- The four plans Catalyst produces, and what each one is for
- The optimisations you can see in explain() output
- How to read a physical plan from the bottom up
- What whole-stage code generation is
Read first
- Transformations vs actions · 3 min read
- Narrow vs wide transformations · 3 min read
Comfortable with these? Read on.
explain("formatted") shows the result. Read physical plans bottom up, look for Exchange (shuffles), the join strategy, and what was pushed into the scan.From code to plan
-
Parsed
Unresolved logical plan: names not checked yet
-
Analysed
Names and types resolved against the catalog
-
Optimised
Rules rewrite the logical plan
-
Physical
Strategies choose how to execute; then code is generated
DataFrame code and SQL end up in the same pipeline, which is why they perform the same. A Python DataFrame call does not run Python per row; it builds this plan.
What the optimiser does
| Rule | Example |
|---|---|
| Predicate pushdown | A filter after a join is moved below it, and into the file scan where possible. |
| Column pruning | Columns that are never used are not read. |
| Constant folding | col("x") > 1 + 2 becomes x > 3. |
| Combining projections and filters | Several withColumn and filter calls collapse into one Project and one Filter. |
| Null propagation, boolean simplification | x AND true becomes x. |
| Join reordering | Only with the cost-based optimiser on (spark.sql.cbo.enabled, off by default) and table statistics collected. |
Using explain
df.explain() # physical plan only df.explain("extended") # all four plans df.explain("formatted") # numbered operators with details: easiest to read df.explain("cost") # with statistics, if available
EXPLAIN FORMATTED SELECT country, SUM(amount) FROM orders WHERE status = 'PAID' GROUP BY country;
Reading a physical plan
Data flows from the bottom (the scan) up to the top (the result). Each indentation level is an input to the line above it.
explain() of an aggregation over a join
*(5) HashAggregate(keys=[country], functions=[sum(amount)])
+- Exchange hashpartitioning(country, 200)
+- *(4) HashAggregate(keys=[country], functions=[partial_sum(amount)])
+- *(4) Project [amount, country]
+- *(4) BroadcastHashJoin [customer_id], [id], Inner, BuildRight
:- *(4) Filter (status = PAID)
: +- *(4) ColumnarToRow
: +- FileScan parquet orders [customer_id,amount,status]
: PushedFilters: [IsNotNull(status), EqualTo(status,PAID)]
+- BroadcastExchange HashedRelationBroadcastMode
+- FileScan parquet customers [id,country]- 1Scans: only the needed columns are read, and the status filter is pushed into the Parquet reader.
- 2BroadcastHashJoin ... BuildRight: customers is small and broadcast; orders is not shuffled.
- 3partial_sum then Exchange then sum: the aggregation pre-aggregates, shuffles by country, and finishes.
- 4*(4): the star and number mark operators fused into one generated function, called whole-stage codegen stage 4.
Whole-stage code generation
Instead of calling a function per operator per row, Spark generates one Java function for a chain of operators (scan, filter, project, partial aggregate) and compiles it at runtime. Rows stay in CPU registers instead of passing through virtual calls. Operators without a star, like Python UDFs or some generators, break the chain. explain("codegen") shows the generated code.
Common mistakes
Reading plans top down
Assuming the pre-execution plan is final
Expecting join reordering by default
Key takeaways
- Catalyst produces parsed, analysed, optimised and physical plans.
- Read physical plans bottom up; Exchange means shuffle.
- Check FileScan for PushedFilters and PartitionFilters.
- Starred operators are fused by whole-stage codegen; UDFs break optimisation.
Check yourself
3 questions1. Which node indicates a shuffle?
Show the answer
Exchange. Exchange nodes redistribute data between stages.
2. What does *(3) in front of an operator mean?
Show the answer
It is part of whole-stage codegen stage 3. The star and id mark operators fused into one generated function.
3. Why can a Python UDF stop filter pushdown?
Show the answer
Catalyst cannot see inside it, so it cannot reason about moving filters through it. UDFs are black boxes to the optimiser.
Go deeper
Transformations vs actionsSpark internals
Dynamic partition pruningSpark internals
Adaptive Query Execution
Primary sources: EXPLAIN · Performance tuning