skip to content

In Spark, what does Catalyst do between your DataFrame code and the physical plan?

level: middleimportance: must knowfreq 70%

answer

  1. a compiler, not a black box
  2. four trees, printed in order
  3. names resolve before rules rewrite
  4. logical is what, physical is how
  5. the shuffle is inserted, never requested

basics

~10 s

Catalyst runs four stages: analysis resolves names and types against the catalog, rule-based logical optimization rewrites the plan, physical planning picks operators and inserts shuffles, and Tungsten generates JVM code for fused operator chains.

solid answer

~50 s

Your code first becomes an **unresolved logical plan** — a tree of operators where table and column names are just strings. The **Analyzer** resolves them against the catalog, attaches types and coerces them, producing the analyzed plan; a bad column name fails here with `AnalysisException`, before any task runs. The **Optimizer** then applies batches of rewrite rules to a fixed point: predicate pushdown, column pruning, constant folding, boolean simplification, limit pushdown. The result is the optimized logical plan — still *what*, not *how*. **Physical planning** turns each logical operator into an executable one (this is where `JoinSelection` chooses a broadcast or sort-merge join), then preparation rules insert `Exchange` and `Sort` nodes wherever an operator's required distribution or ordering is not already satisfied. Finally **whole-stage code generation** compiles fused operator chains into JVM bytecode. `df.explain(True)` prints all four trees.

code

text · 24 lines
text
== Parsed Logical Plan ==
'Aggregate ['dept], ['dept, sum('salary) AS total#22]
+- 'Filter ('salary > 100000)
   +- 'UnresolvedRelation [emp]

== Analyzed Logical Plan ==
dept: string, total: bigint
Aggregate [dept#10], [dept#10, sum(salary#12) AS total#22L]
+- Filter (salary#12 > 100000)
   +- Relation[id#9,dept#10,name#11,salary#12] parquet

== Optimized Logical Plan ==
Aggregate [dept#10], [dept#10, sum(salary#12) AS total#22L]
+- Project [dept#10, salary#12]
   +- Filter (isnotnull(salary#12) AND (salary#12 > 100000))
      +- Relation[id#9,dept#10,name#11,salary#12] parquet

== Physical Plan ==
*(2) HashAggregate(keys=[dept#10], functions=[sum(salary#12)])
+- Exchange hashpartitioning(dept#10, 200), ENSURE_REQUIREMENTS
   +- *(1) HashAggregate(keys=[dept#10], functions=[partial_sum(salary#12)])
      +- *(1) Filter (isnotnull(salary#12) AND (salary#12 > 100000))
         +- *(1) ColumnarToRow
            +- FileScan parquet [dept#10,salary#12] PushedFilters: [IsNotNull(salary), GreaterThan(salary,100000)]

go deeper

for a junior

Be able to name the phases in order — analyze, optimize, plan physically, generate code — and to run df.explain(True) and point at the four printed trees.

for a middle

Explain what each phase consumes and produces, name concrete rules such as predicate pushdown and column pruning, and say which phase reports a missing column versus which one inserts a shuffle.

for a senior

Use the phase boundaries as a diagnosis tool: locate whether a regression is a resolution issue, a missing rewrite, a bad physical choice or a runtime re-plan, and know that AQE makes the printed plan provisional.

for a principal

Own the implications for the platform: where statistics come from, which optimizer behaviours you standardize across jobs, and how plan-level review becomes part of how expensive pipelines get accepted.

## Why there are phases at all Catalyst is Spark SQL's query compiler. It takes a declarative description of a result and produces executable code, and like any compiler it does this in stages, each with a narrow job and a well-defined input and output tree. Knowing the stage boundaries is what lets you say *where* a problem lives: a name that will not resolve is an analysis problem, a filter that was not pushed to the scan is an optimizer problem, an unwanted shuffle is a physical-planning problem. ## Stage 0 — building the unresolved plan A SQL string is parsed by an ANTLR-based parser; DataFrame API calls build the same tree directly, one node per call. Either way the output is an **unresolved logical plan**: operators such as `Project`, `Filter`, `Aggregate`, `Join` and an `UnresolvedRelation` where a table should be. Nothing has been checked. In `explain(True)` output this appears as `== Parsed Logical Plan ==`, with a leading apostrophe on identifiers Spark has not yet resolved (`'salary`). ## Stage 1 — analysis The **Analyzer** walks the tree with resolution rules and a **catalog** (the session catalog, or an external one such as a Hive metastore or a Unity/REST catalog). It replaces `UnresolvedRelation` with a real relation, binds each column reference to an attribute with an id and a data type (`salary#12`), resolves function names, expands `*`, and applies type coercion so that comparing a `DECIMAL` to an `INT` has defined semantics. Aggregate expressions are validated — grouping by one column and selecting another without an aggregate fails here. This stage is where `AnalysisException` comes from, and it is a genuinely useful property: a typo in a column name fails the job in the driver in milliseconds, not two hours into a cluster run. ## Stage 2 — logical optimization The **Optimizer** is rule-based: batches of tree-rewrite rules, each applied repeatedly until the tree stops changing (a fixed point) or a batch's iteration limit is hit. The rules that come up in interviews are: - **Predicate pushdown** — move `Filter` nodes below joins and projections and, where the source supports it, into the scan itself, so fewer rows are ever materialized. - **Column pruning** — drop columns no downstream operator reads, which for a columnar source such as Parquet means whole column chunks are never read from disk. - **Constant folding** — evaluate constant subexpressions at plan time. - **Boolean simplification, null propagation, limit pushdown, subquery decorrelation, combining adjacent filters and projections.** Cost-based join reordering also lives here but is off by default (`spark.sql.cbo.enabled` is `false`); everything above is purely rule-based and always on. The output is the **optimized logical plan**, which still describes *what* to compute, not how. ## Stage 3 — physical planning The **SparkPlanner** applies strategies that convert logical operators into physical ones. This is where implementation choices happen: `JoinSelection` picks `BroadcastHashJoin`, `SortMergeJoin`, `ShuffledHashJoin`, `BroadcastNestedLoopJoin` or `CartesianProduct`; an aggregate becomes `HashAggregate`, `ObjectHashAggregate` or `SortAggregate` depending on whether the buffer types are mutable and fixed-width. A second, easily missed step follows: **preparation rules** run over the chosen physical plan. `EnsureRequirements` compares each operator's *required* child distribution and ordering against what the child actually delivers, and inserts an `Exchange` (a shuffle) or a `Sort` when they do not match. This is the mechanism behind every shuffle you see in a plan — you never asked for one; an operator demanded that its input be partitioned by the join or grouping key. `CollapseCodegenStages` then wraps runs of code-gen-capable operators in a whole-stage codegen node. ## Stage 4 — code generation Tungsten's **whole-stage code generation** fuses each such run of operators into a single Java method — a tight loop over rows with no virtual calls between operators — and compiles it at runtime. In the printed plan these operators carry a `*(n)` prefix, where `n` is the codegen stage id. Operators that cannot be generated (a Python UDF evaluation, some sources) break the chain and appear without the marker. ## And then it changes again With Adaptive Query Execution on — the default since Spark 3.2 — the physical plan is not final. The top of the printed plan reads `AdaptiveSparkPlan isFinalPlan=false`, and after each completed shuffle Spark re-plans the remainder using the statistics it just measured. So `explain()` before execution shows the *initial* plan; the plan that actually ran is the one in the SQL tab of the Spark UI. ## Reading it back `df.explain(True)` prints the parsed, analyzed, optimized and physical plans in order, which is the fastest way to answer "at which stage did my expectation break?" — the column resolved but the filter is not in the optimized plan's scan, the filter is there but a second `Exchange` appeared anyway, and so on.

  • Why does a misspelled column name fail immediately while a data type error may not?
    Column resolution happens in the analysis phase, in the driver, against the catalog schema — so a name that is not in the schema fails before any task is scheduled. A value-level problem, such as a cast that overflows or a corrupt record in one file, is only discovered when an executor actually reads that row, which can be minutes into the job.
  • Who decides that a shuffle is needed, and at which stage?
    The EnsureRequirements preparation rule, after physical operators are chosen. Each physical operator declares the distribution and ordering it requires of its children — a sort-merge join requires both sides hash-partitioned on the join keys and sorted. If the child does not already satisfy that, an Exchange (and possibly a Sort) is inserted. Nothing in your query text asks for a shuffle directly.
  • What is the difference between the optimized logical plan and the physical plan?
    The logical plan says what to compute — a join between two relations on a key, an aggregate by a column — with no commitment to how. The physical plan names implementations: which join algorithm, which aggregate strategy, where data is exchanged, which operators are fused for code generation. One logical plan can have many valid physical plans with very different costs.
  • Does the printed physical plan always match what ran?
    Not with Adaptive Query Execution enabled, which is the default from Spark 3.2. explain() prints the initial plan under an AdaptiveSparkPlan node marked isFinalPlan=false; Spark re-plans the remainder after each shuffle using measured statistics. To see what actually executed, open the query in the SQL tab of the Spark UI, or call explain() on a DataFrame after an action has run it.

A translator first checks that every word in a sentence exists in the dictionary, then rephrases for concision, then chooses the target grammar, and only then speaks. Each pass assumes the previous one succeeded.

saying these in an interview costs you the question

  • Says the optimizer chooses the join algorithm during logical optimization
  • Thinks Catalyst rewrites the SQL text rather than a plan tree
  • Claims analysis and optimization happen on the executors
  • Believes the printed physical plan is always what executed
  • Confuses whole-stage codegen stage ids with scheduler stage ids

context