Skip to content
Reliable Data Engineering
Lesson
Open in the interactive app

Reading Spark Execution Plans Like a Pro

Senior engineers debug Spark by reading plans, not by guessing configs. This module is a field guide: how to produce a plan, how to read it, what every common operator means, and what to look for.


1. From query to physical plan

flowchart LR
    Q["SQL / DataFrame code"] --> P["Parsed (unresolved) logical plan"]
    P -->|"Analyzer: resolve tables, columns, types (catalog)"| A["Analyzed logical plan"]
    A -->|"Optimizer (Catalyst rules): pushdown, pruning,<br/>constant folding, join reorder (CBO), subquery rewrites"| O["Optimized logical plan"]
    O -->|"Planner: choose physical operators<br/>(join strategies, aggregation types)"| PH["Physical plan"]
    PH -->|"Whole-stage codegen,<br/>AQE re-optimisation at runtime"| EX["Executed plan (RDD DAG of stages/tasks)"]
df.explain()                 # physical plan only
df.explain("extended")       # parsed, analyzed, optimized logical + physical
df.explain("formatted")      # compact operator tree + numbered details: the most readable
df.explain("cost")           # optimized logical plan with size/row estimates (stats)
df.explain("codegen")        # generated Java code per codegen stage
EXPLAIN FORMATTED SELECT ...;   EXPLAIN COST SELECT ...;

explain() shows the plan before execution. With AQE, the plan actually executed can differ. Check the SQL / DataFrame tab of the Spark UI after the query runs (isFinalPlan=true).


2. How to read a plan

The five-question checklist

  1. How much is read? Scans: PartitionFilters, PushedFilters, ReadSchema, data-skipping indicators, files/bytes read (SQL tab).
  2. How many shuffles, and of what size? Count the Exchange nodes; check shuffle bytes in the SQL tab.
  3. Which join strategies? BHJ vs SMJ vs SHJ vs nested loop.
  4. Is anything leaving the JVM or blocking optimisation? BatchEvalPython/ArrowEvalPython, missing codegen.
  5. Did AQE change the plan? AQEShuffleRead (coalesced / skewed), strategy switches.

3. Operator field guide

Reading data

OperatorMeaningLook for
FileScan parquet/delta … / BatchScan (V2, e.g. Iceberg)Reads filesPartitionFilters (directory pruning), PushedFilters (row-group skipping), ReadSchema (column pruning), DataFilters
ColumnarToRowConverts columnar batches from the vectorized reader into rowsNormal above scans; many of them can signal unsupported vectorization
InMemoryTableScan / InMemoryRelationReads a cached DataFrameConfirms the cache is used
LocalTableScanData embedded in the plan (small literal data)
Scan ExistingRDDDataFrame created from an RDDLost optimisation opportunities

Row-level operators

OperatorMeaning
FilterPredicate evaluation (check whether it could have been pushed into the scan)
ProjectColumn selection / expressions
Generateexplode, posexplode, inline: multiplies rows; watch output row counts
ExpandDuplicates each row once per grouping set: used for ROLLUP, CUBE, GROUPING SETS and multiple count(DISTINCT …) in one query (rows × number of distinct expressions!)
UnionConcatenation of children
SampleSampling

Exchanges (shuffles and broadcasts)

OperatorMeaning
Exchange hashpartitioning(keys, n)Shuffle by key (joins, aggregations, windows, repartition(col))
Exchange rangepartitioning(keys, n)Range shuffle for global sort (orderBy), with sampling first
Exchange RoundRobinPartitioning(n)repartition(n) without keys
Exchange SinglePartitionEverything into one task (global aggregates without keys, some window functions without partitionBy, limit on large data): a scalability red flag
BroadcastExchangeCollect and broadcast a side for BHJ/BNLJ
ReusedExchangeSpark reused an identical exchange computed elsewhere in the plan (good: computed once)
AQEShuffleRead (formerly CustomShuffleReader)AQE’s reading of shuffle output: coalesced, skewed (split partitions) or local (after converting to broadcast)

Aggregations

OperatorMeaning
HashAggregate (partial) → Exchange → HashAggregate (final)Two-phase aggregation: partial aggregation before the shuffle shrinks the data moved
ObjectHashAggregateHash aggregation for aggregates with object state (e.g. collect_list, some UDAFs)
SortAggregateFallback when hash aggregation isn’t possible (unsupported types or memory): sorts first, so it’s slower
Aggregates with distinctOften rewritten into extra aggregation levels (and Expand for multiple distincts)

Joins

BroadcastHashJoin (with BuildLeft/BuildRight), ShuffledHashJoin, SortMergeJoin (preceded by Sort on both sides), BroadcastNestedLoopJoin, CartesianProduct. See join strategies.

Sorting, windows, limits

OperatorMeaning
Sort [k ASC], true/falsetrue = global sort (after range exchange); false = sort within partitions (for SMJ, windows)
Window [...]Window function evaluation, after Exchange hashpartitioning(partition keys) + Sort
WindowGroupLimit (Spark 3.5)Optimisation for “top-N per group” (row_number() <= N) that filters early
TakeOrderedAndProjectorderBy(...).limit(n) done as per-partition top-n then merge: efficient
CollectLimit / GlobalLimit / LocalLimitLimits; CollectLimit at the root of show()/take()

Leaving the JVM / blocking optimisation

OperatorMeaning
BatchEvalPythonClassic Python UDF evaluation (pickled rows, per row)
ArrowEvalPython / MapInPandas / FlatMapGroupsInPandasPandas/Arrow UDFs and applyInPandas (vectorized, still outside the JVM)
DeserializeToObject / SerializeFromObject / MapElementsTyped Dataset lambdas (object conversion per row)

Subqueries and AQE

OperatorMeaning
Subquery / ScalarSubquery / SubqueryBroadcastUncorrelated subqueries evaluated separately (e.g. DPP’s key list)
dynamicpruningexpression(...) in PartitionFiltersDynamic partition pruning applied
AdaptiveSparkPlan isFinalPlan=false/trueAQE-managed plan; true once executed

On Databricks with Photon, operators appear as PhotonScan, PhotonShuffleExchangeSink, PhotonGroupingAgg, etc., with PhotonResultStage / ColumnarToRow transitions. Operators not supported by Photon fall back to regular Spark (a transition in the plan).


4. Whole-stage code generation

Spark fuses chains of operators within a stage (scan → filter → project → partial aggregate) into one generated Java method that processes rows in a tight loop, avoiding virtual function calls and intermediate objects (the “Volcano” iterator model).


5. Confirm with the SQL tab

The plan says what Spark intends; the SQL/DataFrame tab says what happened, with metrics per operator:

MetricWhereTells you
number of files read, size of files read, files pruned / partitions prunedScanWhether pruning and data skipping worked
number of output rowsEvery operatorSelectivity of filters; join explosions (output ≫ inputs)
shuffle bytes written / records, data sizeExchangeShuffle volume; skew via min/med/max
spill size, peak memorySort, aggregate, joinMemory pressure per operator
build side size / timeBroadcast and hash joinsBroadcast cost
partitions coalesced / skewed partitions splitAQEShuffleReadWhat AQE did
time (sort, aggregate, build, scan)Most operatorsWhere time goes

Click the stages from the SQL DAG to see task-level distributions (max vs median) for skew.


6. Worked examples

Example 1: "Why does this aggregation shuffle 900 GB?"
HashAggregate(keys=[user_id], functions=[count(distinct page), count(distinct session_id)])
+- Exchange hashpartitioning(user_id, 2000)
   +- HashAggregate(keys=[user_id, page, session_id, gid], functions=[])
      +- Expand [[user_id, page, null, 1], [user_id, null, session_id, 2]]
         +- FileScan parquet events[user_id, page, session_id]

Two count(DISTINCT …) in one aggregation trigger Expand, which doubles every row before the shuffle, and the partial aggregate can only deduplicate, not reduce. Fixes: compute each distinct count in a separate aggregation and join the results (or two-level aggregation: first groupBy(user_id, page), then count), or use approx_count_distinct if exactness isn’t required.

Example 2: "The join is slow and spills"
SortMergeJoin [customer_id], [customer_id], Inner
:- Sort [customer_id ASC], false
:  +- Exchange hashpartitioning(customer_id, 200)
:     +- Filter isnotnull(customer_id)
:        +- FileScan parquet orders[...] PartitionFilters: [], PushedFilters: [IsNotNull(customer_id)]
+- Sort [customer_id ASC], false
   +- Exchange hashpartitioning(customer_id, 200)
      +- Filter ((segment = enterprise) AND isnotnull(customer_id))
         +- FileScan parquet customers[...]

Readings:

  • PartitionFilters: [] on orders means no pruning; the query probably filters on a derived date, so filter the partition column instead.
  • The filtered customers side is likely small, but Spark estimated it above 10 MB, so add broadcast() or fresh stats, or let AQE convert at runtime.
  • 200 shuffle partitions for a large orders table means each task handles GBs and spills, so raise partitions or the AQE advisory size.

The SQL tab confirms it with the bytes read, shuffle size per partition and spill metrics.

Example 3: "Everything runs in one task"
Window [row_number() windowspecdefinition(ts ASC ...) AS rn]
+- Sort [ts ASC], false
   +- Exchange SinglePartition

A window function without PARTITION BY forces all rows into a single partition (one task, no parallelism, and likely OOM). Add a partition key (e.g. per user or per day), or compute global rankings differently: approximate quantiles, or a two-step approach with per-partition ranks plus offsets.


Interview questions

What's the difference between the logical and physical plan, and what does each stage of Catalyst do?

The analyzer resolves names and types against the catalog (analyzed logical plan); the optimizer applies rule-based (and cost-based) rewrites such as predicate pushdown, column pruning, constant folding, subquery rewriting and join reordering (optimized logical plan); the planner picks physical operators such as join and aggregation algorithms (physical plan); then whole-stage codegen and AQE shape the executed plan. Logical plans describe what, physical plans describe how.

How do you tell from a plan whether partition pruning and predicate pushdown worked?

On the scan node: PartitionFilters lists the predicates used to skip directories (it should include your partition column filter, or a dynamicpruningexpression for DPP); PushedFilters lists predicates pushed to the file reader for row-group/stripe skipping; ReadSchema shows the pruned columns. In the SQL tab, compare files/bytes read with the table size.

What is Expand and why can it make queries slow?

Expand replicates each input row once per grouping set. It implements ROLLUP/CUBE/GROUPING SETS and multiple DISTINCT aggregates in one query. It multiplies the data before the shuffle by the number of grouping sets or distinct expressions. Splitting distinct counts into separate aggregations, or using approximate algorithms, avoids it.

Why can explain() be misleading with AQE enabled?

explain() shows the plan before execution. AQE re-optimises at runtime from actual shuffle statistics: coalescing partitions, switching join strategies, splitting skewed partitions. The executed plan in the SQL tab (isFinalPlan=true, AQEShuffleRead nodes, strategy changes) is the source of truth.

What breaks whole-stage code generation and why does it matter?

Python UDFs (evaluated outside the JVM), some unsupported expressions and aggregates, and huge generated methods (very wide or deeply nested expressions over the codegen size limit). Without codegen Spark uses the slower iterator model with per-row virtual calls and object creation, so CPU-heavy stages slow down noticeably.