skip to content

In Spark, what makes a transformation narrow or wide, and why does the difference matter?

level: middleimportance: should knowfreq 62%

answer

  1. how many parents does one child need?
  2. one of them can be pipelined in place
  3. the other writes files to disk first
  4. joins can be either one
  5. count the Exchange nodes

basics

~20 s

A transformation is narrow when each output partition draws from a bounded set of input partitions — map, filter, mapPartitions — so it runs in place. It is wide when an output partition may need rows from every input partition, which forces a shuffle: groupBy, join, distinct, repartition.

solid answer

~40 s

The distinction is about **partition dependencies**. In a narrow dependency, each child partition depends on one (or a few fixed) parent partitions, so Spark can pipeline the whole chain inside a single task with no data movement: `map`, `filter`, `flatMap`, `mapPartitions`, `union`, `coalesce`. In a wide dependency, a child partition may need rows from *every* parent partition — `groupByKey`, `reduceByKey`, `join`, `distinct`, `sortByKey`, `repartition` — so Spark must write shuffle files and fetch them across the network. Wide transformations are where the stage boundaries, the network cost and most of the failure modes live. They also reset the partition count to `spark.sql.shuffle.partitions`. Narrow chains are cheap to recover too: losing one partition means recomputing one parent partition, whereas a lost shuffle output can force re-running an entire upstream stage.

code

text · 6 lines
text
== Physical Plan ==
*(2) HashAggregate(keys=[user_id#12], functions=[count(1)])
+- Exchange hashpartitioning(user_id#12, 200), ENSURE_REQUIREMENTS
   +- *(1) HashAggregate(keys=[user_id#12], functions=[partial_count(1)])
      +- *(1) Filter (country#14 = PT)
         +- FileScan parquet [user_id#12,country#14]

go deeper

for a junior

Be able to sort the common operations into the two buckets — map and filter on one side, groupBy and join on the other — and say that the second kind moves data over the network.

for a middle

Define the difference in terms of partition dependencies, and explain the four consequences: pipelining, network cost, partition count reset, and recovery cost.

for a senior

Demonstrate reading Exchange nodes out of a physical plan and reducing shuffle count in a real job through broadcast joins, reused partitioning, and map-side combines.

for a principal

Own the pipeline-shape argument: how many exchanges a workload can afford, when to persist a partitioning in storage so downstream jobs inherit it, and the cluster-wide cost of shuffle-heavy designs.

## Dependencies, not operation names The words *narrow* and *wide* describe a relationship between a child RDD's partitions and its parent's partitions, not an intrinsic property of a function name. - **Narrow dependency:** each partition of the child depends on at most a small, fixed number of parent partitions. In the common case it is exactly one. - **Wide (shuffle) dependency:** a partition of the child may depend on rows from *all* parent partitions, because the rows it needs are scattered by whatever criterion the operation groups on. That single distinction drives most of Spark's execution behaviour. ## Which operations fall where Narrow: `map`, `mapPartitions`, `flatMap`, `filter`, `sample` (without replacement bookkeeping), `union`, `coalesce`, and — importantly — a `join` whose two sides are **already co-partitioned** by the join key with the same partitioner. Wide: `repartition`, `repartitionByRange`, `groupByKey`, `groupBy`, `reduceByKey`, `aggregateByKey`, `distinct`, `sortByKey`, `orderBy`, and any `join` or set operation whose inputs are not already partitioned on the key. The conditional cases are the interesting ones. `join` is wide *by default* and narrow *when Spark can prove co-partitioning* — the same reason `equals` on a custom `Partitioner` matters. A broadcast join is a third thing again: it avoids the shuffle by sending the small side to every executor, making the join narrow with respect to the large side. ## Why it matters — four consequences **1. Pipelining and stage boundaries.** Spark fuses consecutive narrow transformations into one stage and runs them within a single task, one row (or one column batch) at a time, without materialising intermediate results. `df.filter(...).select(...).withColumn(...)` costs one pass. A wide transformation cannot be fused: the parent side must finish and write its output before the child side can fetch it, so it becomes a stage boundary. **2. Cost.** A narrow chain touches no network and no shuffle disk. A wide transformation serialises rows, writes them to local disk, and pulls them across the network — usually the dominant cost of a job and the thing you tune first. **3. Partition count.** Narrow transformations preserve the partition count (except `coalesce`, which merges, and `union`, which sums). Wide transformations set it to `spark.sql.shuffle.partitions`, 200 by default, subject to adaptive coalescing at runtime. **4. Failure recovery.** Lineage recovery is cheap across narrow dependencies: lose one partition, recompute exactly one parent partition. Across a wide dependency it is expensive — the lost partition's inputs came from every parent, so unless the shuffle files survive, Spark must re-run the whole upstream stage. This is precisely why Spark materialises shuffle output to disk and why an external shuffle service exists. ## The `coalesce` special case `coalesce` is narrow, which sounds like unambiguously good news and is actually the source of its famous trap. Because it introduces no stage boundary, its reduced parallelism applies to everything fused with it — `read.filter().coalesce(1).write` runs the read and the filter with one task. Narrowness is what makes it cheap and what makes it dangerous. ## Reducing the shuffle rather than avoiding it You cannot make a genuine regrouping narrow, but you can make it cheaper. `reduceByKey` and `groupByKey` are both wide, yet `reduceByKey` applies the reduction **map-side** before writing shuffle files, so far fewer bytes cross the network; `groupByKey` ships every value. The DataFrame API does this automatically for standard aggregations, which is one concrete reason to prefer it over hand-written RDD code. Broadcasting a small dimension table removes the shuffle on the fact side entirely. Reusing an existing partitioning — for instance, joining twice on the same key without an intervening repartition — lets Spark skip the second exchange. ## Reading it off a plan In a physical plan, every wide dependency appears as an `Exchange` node (`Exchange hashpartitioning(user_id, 200)`), and in the Spark UI each `Exchange` is a stage boundary with shuffle read/write byte counts attached. Counting `Exchange` nodes in `df.explain()` is the fastest way to answer "how many shuffles does my job do", and cutting that count is usually the highest-leverage optimisation available. ## Where the vocabulary trips people A Spark *stage* is a set of tasks with no shuffle between them; the shuffle boundary defines it. That is unrelated to a MapReduce job's fixed map and reduce phases, even though a single Spark shuffle resembles one map/reduce pair — Spark simply chains as many as the query needs, keeping intermediate results in memory where it can, rather than round-tripping through a distributed filesystem between every pair.

  • Can a join ever be narrow in Spark?
    Yes, in two ways. If both sides are already partitioned by the join key with the same partitioner and partition count, Spark can join partition-to-partition with no exchange. A broadcast join is the other case: the small side is shipped whole to every executor, so the large side is read in place with no shuffle. Otherwise a join is wide.
  • Both reduceByKey and groupByKey are wide — why prefer reduceByKey?
    `reduceByKey` combines values map-side before writing shuffle files, so only one partial result per key per partition crosses the network. `groupByKey` ships every individual value, which inflates shuffle bytes and can make a single hot key's value list too large for one executor. The shuffle still happens; it just moves far less data.
  • How do you count the shuffles in a Spark query without running it?
    Call `df.explain()` and count the `Exchange` nodes in the physical plan — each one is a shuffle and a stage boundary, and the node text names the partitioning, for example `Exchange hashpartitioning(user_id, 200)`. In the UI the same boundaries appear as separate stages with shuffle read and write byte counts.

Narrow is everyone tidying their own desk in parallel; wide is everyone dumping their papers into a shared pile so they can be re-sorted by topic before work continues.

saying these in an interview costs you the question

  • Says any transformation that changes row count is wide
  • Claims a join is always a shuffle in Spark
  • Thinks reduceByKey avoids the shuffle entirely
  • Believes narrow transformations each get their own stage
  • Says coalesce is wide because it changes the partition count

context