In Spark, how do you read the physical plan that df.explain() prints?
answer
- read it upside down
- leaves first, root last
- count the shuffles
- the scan line confesses what it read
- the star marks the fused loop
basics
~20 sRead it bottom-up: leaves are scans, indentation shows the tree. Each Exchange is a shuffle and a stage boundary, a *(n) prefix marks whole-stage code generation, and the scan line's PushedFilters and ReadSchema show what was pushed into the source.
solid answer
~50 sThe plan is a tree printed with the root at the top, so **execution flows upward from the leaves**. The bottom lines are `FileScan` or similar sources; `+-` and `:-` mark children, and a node with two children (a join) shows one branch with `:-` and the other with `+-`. Four things carry most of the diagnostic value: **`Exchange`** — a shuffle, and therefore a stage boundary; count them, because each one is a full write-read of the data. **`*(n)`** — this operator is fused into whole-stage codegen stage `n`; a gap in the markers shows where fusion broke. **`PushedFilters` and `ReadSchema`** on the scan — proof that predicate pushdown and column pruning happened. **The join node's name** — `BroadcastHashJoin` versus `SortMergeJoin`, or the alarming `CartesianProduct`. Use `explain("formatted")` for a numbered node list with inputs and outputs, and `explain("cost")` to see the row and size estimates the planner used.
code
text · 10 lines== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- BroadcastHashJoin [cust_id#12], [cust_id#31], Inner, BuildRight
:- Filter (isnotnull(cust_id#12) AND (amount#14 > 100.0))
: +- FileScan parquet db.orders[cust_id#12,amount#14,dt#15]
: PartitionFilters: [isnotnull(dt#15), (dt#15 = 2026-08-01)]
: PushedFilters: [IsNotNull(cust_id), GreaterThan(amount,100.0)]
: ReadSchema: struct<cust_id:string,amount:double>
+- BroadcastExchange HashedRelationBroadcastMode(...)
+- FileScan parquet db.customers[cust_id#31,name#32]go deeper
Know that explain() prints a tree read from the bottom up, and be able to find the scan at the bottom and the join or aggregate above it.
Explain what Exchange, the *(n) markers, PushedFilters and ReadSchema each tell you, and use the explain modes deliberately rather than always calling the default.
Diagnose from a plan: locate the missing pushdown, the surprise shuffle or the wrong join operator, and cross-check against the SQL tab's runtime metrics rather than the estimate.
Make plan review routine for expensive pipelines — captured plans in code review, alerts on plan shape changes, and a shared vocabulary so teams can discuss a regression without re-running the job.
## The shape of the output `df.explain()` prints a tree with the **root at the top and the leaves at the bottom**, using `+-` for a child and `:-` for the first of two children. Data flows the other way: the bottom-most nodes read files, and each row travels upward through the operators above it. Reading bottom-up is the single habit that makes plans legible. ``` SortMergeJoin [cust_id#12], [cust_id#31], Inner :- Sort [cust_id#12 ASC NULLS FIRST], false, 0 : +- Exchange hashpartitioning(cust_id#12, 200), ENSURE_REQUIREMENTS : +- Filter isnotnull(cust_id#12) : +- FileScan parquet [cust_id#12,amount#14] PushedFilters: [IsNotNull(cust_id)] +- Sort [cust_id#31 ASC NULLS FIRST], false, 0 +- Exchange hashpartitioning(cust_id#31, 200), ENSURE_REQUIREMENTS +- FileScan parquet [cust_id#31,name#32] ``` Here both join inputs are shuffled by `cust_id` and sorted, then merged. The `:` column belongs to the left branch; the `+-` at the same indentation starts the right branch. ## The five things to look for **1. `Exchange` — every shuffle you have.** An `Exchange hashpartitioning(key, 200)` means rows are redistributed by hash of `key` into that many partitions; `ENSURE_REQUIREMENTS` is the preparation rule that inserted it because the operator above demanded that distribution. `Exchange rangepartitioning(...)` comes from a global sort, `Exchange SinglePartition` from an unpartitioned window or a global aggregate — the last one is worth a wince, because it funnels everything through one task. Counting Exchanges tells you how many stage boundaries the query has, which is usually the dominant cost. **2. `*(n)` — whole-stage code generation.** Operators fused into one compiled loop share a codegen stage id, printed as `*(1)`, `*(2)` and so on. Operators without a marker are not fused: a Python UDF evaluation, a source that does not support it, or an operator excluded from codegen. The ids are codegen stages, **not** scheduler stages — do not read `*(3)` as "the third stage in the Spark UI". **3. `PushedFilters` and `ReadSchema` on the scan.** `PushedFilters: [IsNotNull(cust_id), GreaterThan(amount,100.0)]` means the data source itself applies those predicates while reading — for Parquet that includes skipping row groups by their min/max statistics. `ReadSchema: struct<cust_id:string,amount:double>` shows column pruning: only those columns are read off disk. If a filter you wrote is **absent** from `PushedFilters` and instead appears as a separate `Filter` node above the scan, Spark read every row and discarded them in memory. Common reasons: the predicate wraps the column in a function or a cast, it involves a UDF, the source does not support that filter type, or it references two columns of the same table. Partition pruning appears separately as `PartitionFilters`, and that one is the difference between listing a few directories and listing all of them. **4. The join operator's name.** `BroadcastHashJoin` (with a `BroadcastExchange` beneath one side), `SortMergeJoin`, `ShuffledHashJoin`, `BroadcastNestedLoopJoin`, `CartesianProduct`. The last two mean your join condition was not an equality, which is fine for a tiny side and catastrophic otherwise. **5. `AdaptiveSparkPlan isFinalPlan=false` at the top.** With Adaptive Query Execution enabled — the default since Spark 3.2 — what `explain()` prints before execution is the *initial* plan. Spark re-plans after each shuffle using measured statistics, so the operators that actually ran may differ. To see the real one, run the query and open it in the **SQL / DataFrame** tab of the Spark UI, which shows the final plan with runtime metrics per node. ## The explain modes - `explain()` — physical plan only. - `explain(True)` or `explain("extended")` — parsed, analyzed, optimized logical plans plus the physical plan. - `explain("formatted")` — a compact tree of numbered nodes followed by a detail block per node listing inputs, outputs, filter conditions and pushed filters. Best for wide plans, because the tree stops being a wall of nested text. - `explain("cost")` — annotates the optimized logical plan with the statistics the planner used (`sizeInBytes`, `rowCount` when available). This is how you discover that a broadcast decision rested on a wild size estimate. - `explain("codegen")` — prints the generated Java source, which matters when you are chasing a codegen fallback. In SQL, `EXPLAIN FORMATTED SELECT ...` gives the same output. ## A working checklist For a slow query, walk the plan bottom-up and ask: does the scan show my filters and only my columns? How many Exchanges are there, and is any of them `SinglePartition`? Which join algorithm was chosen, and is it what I expected from the input sizes? Where do the `*(n)` markers stop, and why? Then compare with the Spark UI's SQL tab, where the same tree carries actual row counts per node — the gap between the plan's expectation and the measured rows is usually the whole story.
- A filter you wrote appears as its own Filter node above the scan instead of in PushedFilters. What does that cost, and why did it happen?Every row is read from storage and discarded in memory, so you pay full I/O and decoding. Typical causes: the predicate wraps the column in a function or cast, it calls a UDF, it compares two columns of the same table, or the source cannot express that filter. Rewriting the predicate so the bare column faces a literal usually restores pushdown.
- How do you see the plan that actually executed rather than the initial one?Run the query, then open the SQL / DataFrame tab in the Spark UI and click the query. It renders the final plan with per-node runtime metrics — rows output, spill, shuffle bytes. With AQE enabled the printed explain() output is the pre-execution plan and can differ, for example where a sort-merge join was converted after real shuffle sizes were measured.
- What does an Exchange SinglePartition node in a plan tell you?All rows are being funnelled into one partition, so one task does that work regardless of cluster size. It comes from a global aggregate with no grouping key, a window without a PARTITION BY, or a global sort collecting to order. It is fine for a small final result and a scalability wall for anything large.
It reads like a family tree drawn with the ancestor at the bottom: the newest generation sits at the top, but everything begins with the roots.
saying these in an interview costs you the question
- Reads the plan top-down as the order of execution
- Cannot say what an Exchange node means
- Treats the *(n) codegen id as the Spark UI stage number
- Assumes a filter in the query text was necessarily pushed down
- Trusts the printed plan as what ran with AQE enabled