Skip to content
Great engineers know 3 min read

Catalyst and physical plans

Read explain() output: parsed, analyzed and optimized plans, and the physical plan Spark runs.

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

Comfortable with these? Read on.

TL;DR Catalyst turns DataFrame and SQL code into a plan in stages: parsed, analysed, optimised, and physical. 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

  1. Parsed

    Unresolved logical plan: names not checked yet

  2. Analysed

    Names and types resolved against the catalog

  3. Optimised

    Rules rewrite the logical plan

  4. 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

RuleExample
Predicate pushdownA filter after a join is moved below it, and into the file scan where possible.
Column pruningColumns that are never used are not read.
Constant foldingcol("x") > 1 + 2 becomes x > 3.
Combining projections and filtersSeveral withColumn and filter calls collapse into one Project and one Filter.
Null propagation, boolean simplificationx AND true becomes x.
Join reorderingOnly with the cost-based optimiser on (spark.sql.cbo.enabled, off by default) and table statistics collected.

Using explain

PySparkSpark SQL
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]
  1. 1
    Scans: only the needed columns are read, and the status filter is pushed into the Parquet reader.
  2. 2
    BroadcastHashJoin ... BuildRight: customers is small and broadcast; orders is not shuffled.
  3. 3
    partial_sum then Exchange then sum: the aggregation pre-aggregates, shuffles by country, and finishes.
  4. 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.

Why this matters: Python UDFs are opaque to Catalyst: it cannot push filters through them, prune columns inside them, or generate code for them. Prefer built-in functions; when you must use Python, pandas UDFs (vectorised with Arrow) are much faster than row-at-a-time UDFs.

Common mistakes

Reading plans top down

Execution starts at the bottom with the scans.

Assuming the pre-execution plan is final

With AQE, joins and partition counts change at runtime; check the SQL tab after execution.

Expecting join reordering by default

The cost-based optimiser is off by default and needs statistics (ANALYZE TABLE).

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 questions

1. 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

Primary sources: EXPLAIN · Performance tuning