In Spark, what determines where the DAG scheduler places a stage boundary?
answer
- not every operator earns its own segment
- ask where each output row must travel
- one parent partition or many?
- look for Exchange in the physical plan
- materialize to disk, then fetch
basics
~10 sSpark cuts a stage at every wide dependency: any point where a row's destination depends on its key, so data must be redistributed. Narrow operations are pipelined into the current stage instead.
solid answer
~50 sThe `DAGScheduler` inside the Spark driver walks the plan backwards from the action and starts a new stage wherever it meets a **wide dependency** — a step where a child partition can depend on many parent partitions, so records must be re-routed across the cluster. `groupBy`, `reduceByKey`, `distinct`, `orderBy`, `repartition` and any sort-merge join on a key that is not already co-partitioned all create one. **Narrow** steps — `filter`, `map`, `select`, `withColumn`, `union`, a broadcast join — keep each output partition dependent on exactly one input partition, so Spark fuses them into the current stage and runs them as a single per-record function. The visible marker in a DataFrame plan is `Exchange`: every `Exchange` node in `explain()` is one boundary, so counting them tells you the stage count. A boundary is also a materialization point — the upstream stage writes its output to disk before the downstream stage can fetch it — which is why boundaries, not operators, dominate runtime.
code
python · 6 lines(spark.read.parquet("/data/events")
.filter("v > 0")
.select("k", "v")
.repartition(8) # boundary 1: round-robin exchange
.groupBy("k").count() # boundary 2: hash exchange
.write.parquet("/out"))go deeper
Know that a shuffle is what ends a stage, and be able to sort common operators into shuffling (groupBy, join, distinct, orderBy) versus non-shuffling (filter, map, select).
Explain the wide-versus-narrow dependency rule precisely, point at Exchange in an explain() plan as the observable boundary, and say why the upstream stage must materialize before the next one starts.
Demonstrate boundary removal in real plans: broadcast a small side, group once instead of twice, prefer coalesce to repartition when narrowing, and exploit pre-bucketed inputs to skip an exchange.
Frame boundaries as the unit of cost in a pipeline's design: the storage layout, partitioning and bucketing choices you standardize on determine how many exchanges every downstream job must pay for.
## The rule in one line Spark cuts a stage wherever it cannot keep computing a partition locally. Formally, a boundary appears at every **wide dependency** (`ShuffleDependency`); everything else is fused into the current stage. ## Narrow versus wide A **narrow dependency** means each output partition draws from exactly one input partition. `map`, `filter`, `select`, `withColumn`, `flatMap`, `mapPartitions`, `union` and a broadcast (map-side) join are all narrow: the executor already holds everything a given output partition needs. Spark therefore never materializes the intermediate result; it composes the operators into one function evaluated record by record inside a single task. Ten chained narrow transformations still cost one stage. A **wide dependency** means an output partition can draw from many input partitions, because the row's destination is decided by its key rather than by where it already sits. Grouping, deduplication, key-based joins, global ordering and explicit repartitioning are all wide. To satisfy them, every executor must write records out bucketed by destination and every downstream task must fetch its bucket from every upstream task. That redistribution cannot be pipelined, so Spark ends the stage there. ## What the boundary actually costs A stage boundary is not just a bookkeeping line. The upstream `ShuffleMapStage` **materializes**: each of its tasks writes its output to local disk as shuffle files before the downstream stage begins. That has three consequences interviewers care about. First, it serializes execution — the downstream stage cannot start until enough upstream tasks have written their output, so a boundary is a barrier in the timeline. Second, it converts an in-memory pipeline into disk I/O plus network transfer, which is usually the dominant cost in a Spark job. Third, it creates a recovery point: because the map output exists on disk, a failed downstream task can simply re-fetch rather than recompute the whole lineage — and, conversely, losing an executor that holds map output forces Spark to recompute the parent stage. ## Reading boundaries out of a plan For DataFrames and Spark SQL, `df.explain()` is authoritative. Each `Exchange` node in the physical plan is exactly one stage boundary, and its partitioning tells you which kind: `hashpartitioning(key, n)` for a grouping or key join, `RoundRobinPartitioning(n)` for a bare `repartition(n)`, `rangepartitioning` for a global `orderBy`. Count the `Exchange` nodes, add one, and you have the stage count for that job. For RDDs the equivalent signal is `toDebugString`, where each indentation shift marks a `ShuffledRDD`. ## Boundaries you can remove Because boundaries dominate cost, the useful engineering move is deleting them rather than tuning them. A join whose smaller side fits in memory can be turned into a broadcast join, which is narrow — the small table is shipped to every executor and no exchange is needed at all. Two aggregations on the same key can share one exchange if you group once instead of twice. `coalesce(n)` narrows partition count without a shuffle, while `repartition(n)` always adds one, so the wrong choice at the end of a job inserts a boundary purely to write files. Reading data already bucketed on the join key can let Spark skip the exchange entirely. And an accidental `distinct()` or a `dropDuplicates()` inside a loop is a boundary someone added without noticing. ## What is not a boundary A source scan is not a boundary — partitioning at the read is decided by input splits, not by a shuffle. A `cache()`/`persist()` is not a boundary either; it changes where a later job reads from but does not cut the current stage. Adding executors does not add or remove stages. And an action does not create a boundary, it creates a job whose last stage is the `ResultStage`. ## Where adaptive execution fits From Spark 3.2 onwards adaptive query execution is enabled by default, and it re-plans *at* these boundaries: after a map stage completes, Spark has real output statistics and may coalesce the downstream partition count or switch a join strategy before launching the next stage. That means the physical plan you printed before running (`isFinalPlan=false`) can differ from what actually ran — but the boundary positions themselves still come from the same wide-dependency rule. ## The interview answer Say the rule (wide dependency), name three wide and three narrow operators, point at `Exchange` in the plan as the observable marker, and close with the consequence: boundaries force materialization to disk plus a network fetch, so the fastest job is the one with the fewest of them.
- How can you count a Spark job's stages before running it?Print the physical plan with `df.explain()` and count the `Exchange` nodes: stages equal exchanges plus one. The partitioning shown on each `Exchange` also tells you why it is there — `hashpartitioning` for a group or key join, `RoundRobinPartitioning` for `repartition(n)`, `rangepartitioning` for a global `orderBy`.
- Why does repartition(50) add a stage boundary while coalesce(50) usually does not?`repartition` redistributes every row across the target partitions, which is a wide dependency and needs a full shuffle. `coalesce` only merges existing partitions into fewer, so each output partition reads whole parent partitions — a narrow dependency, no exchange, and no boundary. The price is that `coalesce` cannot balance uneven partitions and cannot increase the count.
- Does converting a join to a broadcast join remove a stage boundary?Yes. A broadcast (map-side) join ships the small side to every executor, so each output partition depends on exactly one partition of the large side — a narrow dependency. The exchange on the large table disappears, and the join fuses into the existing stage. It only works while the broadcast side fits comfortably in executor and driver memory.
saying these in an interview costs you the question
- Says every transformation creates a new stage
- Thinks a cache() or persist() call cuts a stage
- Believes adding executors changes the number of stages
- Cannot name a single wide versus narrow operator
- Claims a broadcast join still requires an exchange