Every Spark job is a graph of datasets, and every edge in that graph is a dependency of one of two kinds. A narrow dependency lets a task compute its output partition from a small, fixed set of parent partitions it can read on the spot. A wide dependency means each output partition needs data from many or all parent partitions, so the data must be redistributed across the cluster first. That single distinction decides where stages begin and end, what gets written to disk, how much crosses the network, and what has to be recomputed when an executor dies.

The mechanics of the shuffle itself, map-side sort, spill files and fetches, are in the Spark shuffle article, and the scheduler that cuts stages is in the DAG scheduler article. This page is about the dependency model: how to see it, predict it, and change it.

Advertisement

The definition, stated precisely

The original RDD paper defines it from the parent's side: a dependency is narrow when each partition of the parent is used by at most one partition of the child, and wide when several child partitions may depend on it. Spark's source states it from the child's side: a narrow dependency is one where each child partition depends on a small number of parent partitions, and the child can be computed by pipelining. The practical test is simpler than either: can a task compute its partition by iterating over parent partitions it can address directly, without someone first regrouping the data by key? If yes, the dependency is narrow.

Narrow dependencies chain. A textFile followed by map, filter and map runs as one task per partition, each record flowing through all three functions before the next record is read, never materialized in between. A wide dependency breaks that chain. Upstream tasks write their output partitioned by the downstream partitioner into shuffle files on local disk; downstream tasks start only when those files exist and then fetch their slice from every upstream task.

Narrow: a child reads a few fixed parents. Wide: it reads a slice of every parent.Narrow (map, filter, union, coalesce)P0P1P2C0C1C2same task, same executor, no disk, no networkWide (groupByKey, reduceByKey, join, distinct)P0P1P2C0C1C2map side writes shuffle files; reduce side fetches themStage boundary: noneoperators pipeline into one task per partitionStage boundary: yesShuffleDependency ends the map stageRecovery: a lost narrow partition is recomputed from its parents; lost shuffle output re-runs map tasks of the parent stage.
Narrow dependencies keep each partition inside one task. A wide dependency connects every parent partition to every child partition through shuffle files, and that is where Spark ends a stage.

The classes behind the words

In the RDD API a dependency is a real object, and every RDD exposes its list through rdd.dependencies. The base classes are short enough to read. A narrow dependency answers one question, which parent partitions a given child partition reads. A shuffle dependency instead carries what the shuffle needs: the partitioner that decides which reducer gets each key, the serializer, optional key ordering, and an optional aggregator for map-side combining.

// Simplified from org.apache.spark.Dependency
abstract class Dependency[T] { def rdd: RDD[T] }

abstract class NarrowDependency[T](rdd: RDD[T]) extends Dependency[T] {
  def getParents(partitionId: Int): Seq[Int]   // which parent partitions this child reads
}

class OneToOneDependency[T](rdd: RDD[T]) extends NarrowDependency[T](rdd) {
  override def getParents(partitionId: Int) = List(partitionId)
}

class RangeDependency[T](rdd: RDD[T], inStart: Int, outStart: Int, length: Int)
    extends NarrowDependency[T](rdd) {          // used by union: a block of child ids maps to one parent
  override def getParents(partitionId: Int) =
    if (partitionId >= outStart && partitionId < outStart + length)
      List(partitionId - outStart + inStart) else Nil
}

class ShuffleDependency[K, V, C](rdd: RDD[_ <: Product2[K, V]], partitioner: Partitioner, ...)
  // plus serializer, key ordering, optional aggregator and map-side combine flag,
  // and a shuffleId registered with the shuffle manager

Most operators use OneToOneDependency: child partition i reads parent partition i. union uses RangeDependency, mapping a contiguous block of child partitions onto each parent, so the union of a 4-partition and a 6-partition RDD has 10 partitions and no shuffle. coalesce without shuffle uses a narrow dependency whose getParents returns several parent partitions, grouped with locality in mind. cartesian is an edge case worth knowing: it declares narrow dependencies on both parents, so it does not shuffle, but each parent partition is read by many tasks and the output grows as the product of the inputs.

You can check any RDD yourself. toDebugString prints the lineage with indentation that changes at each shuffle boundary.

val sales = sc.textFile("s3a://bucket/sales/*.csv", 8)       // 8 partitions
  .map(_.split(","))
  .map(f => (f(0), f(2).toDouble))                            // (storeId, amount)

sales.dependencies.foreach(d => println(d.getClass.getSimpleName))
// OneToOneDependency

val totals = sales.reduceByKey(_ + _)
totals.dependencies.foreach(d => println(d.getClass.getSimpleName))
// ShuffleDependency

println(totals.toDebugString)
// Representative output, trimmed; 8 is a minimum partition count, not a guarantee
// (8) ShuffledRDD[4] at reduceByKey ...
//  +-(8) MapPartitionsRDD[3] at map ...
//     |  MapPartitionsRDD[2] at map ...
//     |  s3a://bucket/sales/*.csv MapPartitionsRDD[1] at textFile ...
//     |  s3a://bucket/sales/*.csv HadoopRDD[0] at textFile ...

Advertisement

Which operators create which

NarrowWideDepends on partitioning
map, flatMap, filter, mapPartitionsgroupByKey, reduceByKey, aggregateByKey on unpartitioned inputreduceByKey or join when the input already has the target partitioner
mapValues, flatMapValues (keep the partitioner)distinct, repartition, partitionBycogroup and join of co-partitioned parents
union, zipPartitions, samplesortByKey (range partitioning, plus a sampling job)DataFrame joins on bucketed tables with matching buckets
coalesce(n) without shufflejoins whose sides are partitioned differentlyAggregations after a shuffle on the same keys

The third column is where the leverage is. A key-based operation needs all values for a key in one partition. If the data is already arranged that way, and Spark knows it, no shuffle is required. In the RDD API Spark knows it when the parent's partitioner equals the partitioner the operation would use: reduceByKey on an RDD already hash-partitioned the same way runs as a narrow per-partition combine.

Worked example: making a join narrow

A recommendation job joins a 50 GB event stream against a 5 GB user table every hour, several times, by userId. Each join shuffles both sides by default: two wide dependencies, roughly 55 GB of shuffle write and read per join. Partition both sides once with the same partitioner and persist them, and later joins find matching partitioners and need no new shuffle.

import org.apache.spark.HashPartitioner
val part = new HashPartitioner(200)

// Partition both sides once, with the same partitioner, and keep them
val users  = sc.parallelize(userRows).partitionBy(part).persist()
val events = rawEvents.map(e => (e.userId, e)).partitionBy(part).persist()

val joined = users.join(events)     // both parents already hash-partitioned by part
println(joined.toDebugString)
// No new ShuffledRDD above the two partitionBy shuffles: the CoGroupedRDD inside
// join() takes a OneToOneDependency on each parent whose partitioner matches

// mapValues keeps the partitioner; map throws it away
users.mapValues(_.name).partitioner        // Some(HashPartitioner)
users.map { case (k, v) => (k, v.name) }.partitioner   // None

The cost moves rather than disappears: partitionBy is itself a shuffle, paid once instead of per join. It pays off when the partitioned dataset is reused several times, and it fails silently if any step in between throws the partitioner away. map does exactly that, because Spark cannot know that your function kept the keys unchanged, so prefer mapValues and flatMapValues on partitioned pair RDDs. Check .partitioner in a test; it is the cheapest assertion you can write about performance.

The same model in DataFrames and SQL

DataFrames hide the dependency objects but not the idea. The physical plan states what distribution each operator requires from its children, and when a child does not provide it, the planner inserts an Exchange node. Every Exchange hashpartitioning, Exchange rangepartitioning or Exchange SinglePartition in explain() is a wide dependency and a stage boundary. BroadcastExchange is different: the small side is collected and sent to every executor, and the large side is joined in place without being shuffled.

# PySpark, DataFrame API
orders = spark.read.parquet("s3a://bucket/orders")
custs  = spark.read.parquet("s3a://bucket/customers")

q = (orders.filter("amount > 100")
           .join(custs, "customer_id")
           .groupBy("region").sum("amount"))
q.explain()

# Representative output, trimmed. Each Exchange is a wide dependency.
# HashAggregate(keys=[region], functions=[sum(amount)])
# +- Exchange hashpartitioning(region, 200)
#    +- HashAggregate(keys=[region], functions=[partial_sum(amount)])
#       +- SortMergeJoin [customer_id], [customer_id], Inner
#          :- Sort ... +- Exchange hashpartitioning(customer_id, 200)
#          :              +- Filter (amount > 100) +- FileScan parquet orders
#          +- Sort ... +- Exchange hashpartitioning(customer_id, 200)
#                         +- FileScan parquet customers

Reading the plan above: three shuffles. Both join inputs are hash-partitioned by customer_id, then the aggregate needs data hash-partitioned by region. Note the partial aggregate before the last Exchange, the map-side combine that shrinks shuffle data. You can remove the join shuffles by broadcasting a small customers table, or by writing both tables bucketed by customer_id into the same number of buckets, which lets the planner treat them as already co-partitioned. Adaptive execution can also convert a sort-merge join to a broadcast join at runtime and merge small shuffle partitions; see Spark AQE.

Seeing dependencies in the Spark UI

The web UI shows the same model without any code. On a job's page the DAG visualization draws one box per stage and the RDDs or operators inside it; every edge between boxes is a wide dependency, and everything inside a box was pipelined through narrow ones. The stages table then tells you what each boundary cost: a stage that ends at a shuffle reports shuffle write, and the stage after it reports shuffle read of roughly the same size.

Use those two columns as your budget. In the event-join example above, before pre-partitioning, each hourly run showed three stages per join, two map stages and one reduce stage, with about 55 GB of shuffle write. After the change, the first run showed the partitionBy shuffles and every later join in the same application showed only a short stage with no shuffle read beyond the cached partitions. If the numbers do not drop, a partitioner was lost somewhere; look for a map, a flatMap, or a DataFrame round trip between the partitioning step and the join.

In SQL jobs the SQL tab shows the physical plan with metrics per node. Each Exchange node reports the data size it moved, which is the quickest way to rank shuffles by cost before deciding which to remove.

What a dependency costs on failure

Lineage is Spark's fault-tolerance mechanism, and the dependency type sets the price of using it. If an executor dies holding a partition produced through narrow dependencies only, Spark recomputes that one partition by re-running its chain from the parents it needs. Nothing else is touched.

Behind a wide dependency the picture changes. Reducers read shuffle files written by map tasks on other executors. If those files are lost with an executor, a fetch fails, and the scheduler marks the map output lost, resubmits the parent stage to regenerate the missing map outputs, and retries the reducers. One lost machine can re-run a slice of an expensive upstream stage. The lineage article covers checkpointing, which cuts this chain for long iterative jobs.

Shuffle files have an upside too: they are an implicit cache. A later job that needs the same shuffle output, still on disk, shows that stage as skipped in the UI and does not recompute it.

Traps and trade-offs

The most common trap is coalesce(1) before a write. Because coalesce without shuffle is narrow, it does not run after the upstream work; it joins the upstream stage. That whole stage, every scan, parse and filter, now runs as a single task on one core. repartition(1) adds a shuffle, which sounds worse, but it lets the upstream stage keep its full parallelism and moves only the final, filtered data through one task. The coalesce versus repartition article works through the numbers.

  • Wide is not always bad. A shuffle is how Spark rebalances skewed or badly sized input; a well-placed repartition can make everything after it faster.
  • Narrow is not always cheap. A very long narrow chain means a long lineage, so recovery and plan analysis get slower; checkpoint or persist periodically in iterative jobs.
  • groupByKey and reduceByKey are both wide, but reduceByKey combines on the map side, so far less data crosses the shuffle.
  • Pre-partitioning trades one shuffle and memory for persistence against many later shuffles. Do it for data reused several times, not once.
  • Skew turns a wide dependency into a straggler. One hot key lands on one reducer however many partitions you choose.

What to do next

  1. Run explain() or toDebugString on your most expensive job and count its Exchanges or shuffle boundaries.
  2. For each one, write down why it is there: aggregation, join, sort or explicit repartition.
  3. Replace groupByKey with reduceByKey or aggregateByKey wherever the values are combined.
  4. Broadcast small join sides, and bucket or pre-partition large tables joined repeatedly on the same key.
  5. Replace coalesce(1) before writes with repartition(1) unless the upstream stage is trivial.
  6. Add a test asserting the partitioner on reused pair RDDs, and compare shuffle read and write in the stage view before and after each change.
Key takeaway: Narrow dependencies are pipelined work inside one task; wide dependencies are where data is regrouped across the cluster, written to disk, and where stages end. The type of every edge decides parallelism, network traffic and recovery cost. Learn to read it from toDebugString and explain(), keep partitioners alive with mapValues, co-partition or broadcast to remove repeated shuffles, and remember that coalesce is narrow, which is exactly why it can quietly serialise a whole stage.