乐于分享
好东西不私藏

【极速学习spark源码】窄依赖 vs 宽依赖:Dependency.scala 的设计

【极速学习spark源码】窄依赖 vs 宽依赖:Dependency.scala 的设计

基于 Spark 4.2,分支 branch-4.2

阅读时长约 10 分钟 · 入门到中级

背景

Spark 里经常听到两个词:

  • 窄依赖(narrow dependency)

  • 宽依赖(wide dependency)

很多解释会说:

窄依赖不会产生 shuffle,宽依赖会产生 shuffle。

这句话没错,但还不够。真正从源码看,窄依赖和宽依赖的区别体现在两个地方:

  1. Dependency.scala 怎么表达父子分区关系

  2. DAGScheduler 怎么根据 ShuffleDependency 切 Stage


一、Dependency 的最小抽象

Dependency.scala:41

abstract classDependency[TextendsSerializable{  def rdd: RDD[T]}

这就是所有 RDD 依赖的基类。它只做一件事:指向父 RDD。

RDD 自己通过 getDependencies 暴露依赖列表:

protected def getDependencies: Seq[Dependency[_]] = deps // RDD.scala:131final def dependencies: Seq[Dependency[_]] = ...         // RDD.scala:260

有了这些依赖,Spark 就能从最后一个 RDD 往前追 lineage。


二、窄依赖:子分区依赖少量父分区

Dependency.scala:52

abstract classNarrowDependency[T](_rddRDD[T]) extendsDependency[T{  def getParents(partitionId: Int): Seq[Int]}

窄依赖的关键不是“没有 shuffle”这四个字,而是这行方法:

def getParents(partitionId: Int): Seq[Int]

它回答一个问题:

子 RDD 的第 partitionId 个分区,需要父 RDD 的哪些分区?

典型实现有两个。

1、OneToOneDependency

Dependency.scala:266

classOneToOneDependency[T](rddRDD[T]) extendsNarrowDependency[T](rdd{  override def getParents(partitionId: Int): List[Int] = List(partitionId)}

这就是最常见的窄依赖:子分区 N 只依赖父分区 N。

mapfiltermapPartitions 大多是这种形态。

2、RangeDependency

Dependency.scala:280

classRangeDependency[T](    rddRDD[T],    inStartInt,    outStartInt,    lengthIntextends NarrowDependency[T](rdd) {

它描述的是一段子分区范围到一段父分区范围的映射,常见于 union / partition range 类场景。


三、宽依赖:ShuffleDependency

Dependency.scala:84

class ShuffleDependency[KClassTagV: ClassTagC: ClassTag](    @transient private val _rdd: RDD[_ <: Product2[K, V]],    val partitioner: Partitioner,    val serializer: Serializer = SparkEnv.get.serializer,    val keyOrdering: Option[Ordering[K]] = None,    val aggregator: Option[Aggregator[K, V, C]] = None,    val mapSideCombine: Boolean = false,    val shuffleWriterProcessor: ShuffleWriteProcessor = new ShuffleWriteProcessor)  extends Dependency[Product2[K, V]]

它比 NarrowDependency 重得多,里面带着:

  • partitioner:shuffle 后有多少个 reduce 分区,key 怎么分配

  • serializer:shuffle 数据怎么序列化

  • aggregator:是否聚合

  • mapSideCombine:是否 map 端聚合

  • shuffleId:唯一标识这次 shuffle(Dependency.scala:129

这就是宽依赖的核心:它不是一个简单的父子分区映射,而是一条 shuffle 边。


四、map / filter 为什么是窄依赖

RDD.scala:424

def map[UClassTag](fT => U): RDD[U] = withScope {  val cleanF = sc.clean(f)  new MapPartitionsRDD[U, T](this(_, _, iter) => iter.map(cleanF))}

RDD.scala:444

def filter(fT => Boolean): RDD[T] = withScope {  val cleanF = sc.clean(f)  new MapPartitionsRDD[T, T](this(_, _, iter) => iter.filter(cleanF), preservesPartitioning = true)}

两者都创建 MapPartitionsRDD

再看 MapPartitionsRDD.scala:50

private[spark] class MapPartitionsRDD[UClassTagT: ClassTag](    var prev: RDD[T],    f: (TaskContext, Int, Iterator[T]) => Iterator[U],    preservesPartitioning: Boolean = false,    ...) extends RDD[U](prev)

extends RDD[U](prev) 会走 RDD 单父构造器,自动创建 OneToOneDependency。这个构造器在 RDD.scala:102-104

def this(@transient oneParent: RDD[_]) =  this(oneParent.context, List(new OneToOneDependency(oneParent)))

所以:

map / filter
  → MapPartitionsRDD
    → RDD 单父构造器
      → OneToOneDependency

这就是窄依赖。


五、reduceByKey / groupByKey 为什么是宽依赖

reduceByKey 是 key-value RDD 的典型 shuffle 操作。

PairRDDFunctions.scala:305

def reduceByKey(partitionerPartitionerfunc(V, V) => V): RDD[(K, V)] = self.withScope {  combineByKeyWithClassTag[V]((v: V) => v, func, func, partitioner)}

groupByKey 也会进入 combineByKeyWithClassTag,入口在 PairRDDFunctions.scala:497

最终会创建 ShuffledRDDShuffledRDD.scala:78

override def getDependencies: Seq[Dependency[_]] = {  val serializer = userSpecifiedSerializer.getOrElse {    val serializerManager = SparkEnv.get.serializerManager    if (mapSideCombine) {      serializerManager.getSerializer(implicitly[ClassTag[K]], implicitly[ClassTag[C]])    } else {      serializerManager.getSerializer(implicitly[ClassTag[K]], implicitly[ClassTag[V]])    }  }  List(new ShuffleDependency(prev, part, serializer, keyOrdering, aggregator, mapSideCombine))}

一句话:

reduceByKey / groupByKey
  → combineByKeyWithClassTag
    → ShuffledRDD
      → ShuffleDependency

这就是宽依赖。


六、join 为什么有时窄、有时宽

join 不是无脑 shuffle。它先走 cogroup

PairRDDFunctions.scala:544

def join[W](other: RDD[(K, W)], partitioner: Partitioner): RDD[(K, (V, W))] = self.withScope {  this.cogroup(other, partitioner).flatMapValues(...)}

cogroup 创建 CoGroupedRDDPairRDDFunctions.scala:797

new CoGroupedRDD[K](Seq(self, other), partitioner)

真正判断是否 shuffle 的地方在 CoGroupedRDD.scala:98

override def getDependencies: Seq[Dependency[_]] = {  rdds.map { rdd =>    if (rdd.partitioner == Some(part)) {      new OneToOneDependency(rdd)    } else {      new ShuffleDependency[K, Any, CoGroupCombiner](rdd.asInstanceOf[RDD[_ <: Product2[K, _]]], part, serializer)    }  }}

这段特别重要:

  • 如果父 RDD 已经有相同的 partitioner:OneToOneDependency

  • 如果 partitioner 不匹配:ShuffleDependency

所以 join 的性能不只看数据量,还看两边是否已经按同一个 partitioner 分区。


七、DAGScheduler 如何用 ShuffleDependency 切 Stage

宽依赖之所以重要,是因为它会切 Stage。

DAGScheduler 里有几个关键方法:

  • DAGScheduler.scala:704createResultStage

  • DAGScheduler.scala:801getShuffleDependenciesAndResourceProfiles

  • DAGScheduler.scala:837getMissingParentStages

  • DAGScheduler.scala:528getOrCreateShuffleMapStage

  • DAGScheduler.scala:574createShuffleMapStage

核心逻辑可以简化成:

// 简化后的伪代码while (traversing RDD graph) {  dependency match {    case shuffleDep: ShuffleDependency[_, _, _] =>      getOrCreateShuffleMapStage(shuffleDep, firstJobId)    case narrowDep: NarrowDependency[_] =>      visit(narrowDep.rdd)  }}

也就是说:

  • 遇到 NarrowDependency:继续往父 RDD 走,仍然属于当前 Stage

  • 遇到 ShuffleDependency:停下来,创建或获取一个父 ShuffleMapStage

这就是 Stage 切分的根本规则。


八、区别和联系

对比项窄依赖宽依赖
源码类型NarrowDependencyShuffleDependency
父子分区关系子分区依赖少量父分区多个父分区重分布到多个子分区
是否 shuffle
是否切 Stage不切
是否可 pipeline可以不可以跨 shuffle pipeline
典型操作map / filter / mapPartitionsreduceByKey / groupByKey / repartition
join 特例partitioner 相同可窄依赖partitioner 不同走 shuffle

九、实战:观察依赖类型

val rdd = sc.parallelize(1 to 102)val mapped = rdd.map(_ + 1)val pairs = sc.parallelize(Seq(("a"1), ("b"2), ("a"3)), 2)val reduced = pairs.reduceByKey(_ + _)println("mapped dependencies:")mapped.dependencies.foreach(d => println(d.getClass.getSimpleName))println("reduced dependencies:")reduced.dependencies.foreach(d => println(d.getClass.getSimpleName))println("reduced lineage:")println(reduced.toDebugString)

输出:

mapped dependencies:
OneToOneDependency
reduced dependencies:
ShuffleDependency
reduced lineage:
(2) ShuffledRDD[3] at reduceByKey at ShuffleTest.scala:44 []
 +-(2) ParallelCollectionRDD[2] at parallelize at ShuffleTest.scala:43 []

十、总结

窄依赖和宽依赖不是两个抽象概念,而是两个非常具体的源码对象:

NarrowDependency
  └─ getParents(partitionId): Seq[Int]

ShuffleDependency
  ├─ partitioner
  ├─ serializer
  ├─ aggregator
  ├─ mapSideCombine
  └─ shuffleId

调度层只认一个关键事实:

遇到 NarrowDependency 继续往上追;遇到 ShuffleDependency 就切 Stage。

这就是为什么 map.filter.map 可以在一个 Stage 里 pipeline,而 reduceByKey 会天然制造新的 Stage。


每天花费10分钟学习spark,让你技术之路走得更稳、更快。

喜欢的点个关注。