基于 Spark 4.2,分支
branch-4.2阅读时长约 10 分钟 · 入门到中级
背景
Spark 里经常听到两个词:
窄依赖(narrow dependency)
宽依赖(wide dependency)
很多解释会说:
窄依赖不会产生 shuffle,宽依赖会产生 shuffle。
这句话没错,但还不够。真正从源码看,窄依赖和宽依赖的区别体现在两个地方:
Dependency.scala怎么表达父子分区关系DAGScheduler怎么根据ShuffleDependency切 Stage
一、Dependency 的最小抽象
Dependency.scala:41:
abstract classDependency[T] extendsSerializable{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](_rdd: RDD[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](rdd: RDD[T]) extendsNarrowDependency[T](rdd) {override def getParents(partitionId: Int): List[Int] = List(partitionId)}
这就是最常见的窄依赖:子分区 N 只依赖父分区 N。
map、filter、mapPartitions 大多是这种形态。
2、RangeDependency
Dependency.scala:280:
classRangeDependency[T](rdd: RDD[T],inStart: Int,outStart: Int,length: Int) extends NarrowDependency[T](rdd) {
它描述的是一段子分区范围到一段父分区范围的映射,常见于 union / partition range 类场景。
三、宽依赖:ShuffleDependency
Dependency.scala:84:
class ShuffleDependency[K: ClassTag, V: ClassTag, C: 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[U: ClassTag](f: T => U): RDD[U] = withScope {val cleanF = sc.clean(f)new MapPartitionsRDD[U, T](this, (_, _, iter) => iter.map(cleanF))}
RDD.scala:444:
def filter(f: T => 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[U: ClassTag, T: 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(partitioner: Partitioner, func: (V, V) => V): RDD[(K, V)] = self.withScope {combineByKeyWithClassTag[V]((v: V) => v, func, func, partitioner)}
groupByKey 也会进入 combineByKeyWithClassTag,入口在 PairRDDFunctions.scala:497。
最终会创建 ShuffledRDD。ShuffledRDD.scala:78:
override def getDependencies: Seq[Dependency[_]] = {val serializer = userSpecifiedSerializer.getOrElse {val serializerManager = SparkEnv.get.serializerManagerif (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 创建 CoGroupedRDD。PairRDDFunctions.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:704:createResultStageDAGScheduler.scala:801:getShuffleDependenciesAndResourceProfilesDAGScheduler.scala:837:getMissingParentStagesDAGScheduler.scala:528:getOrCreateShuffleMapStageDAGScheduler.scala:574:createShuffleMapStage
核心逻辑可以简化成:
// 简化后的伪代码while (traversing RDD graph) {dependency match {case shuffleDep: ShuffleDependency[_, _, _] =>getOrCreateShuffleMapStage(shuffleDep, firstJobId)case narrowDep: NarrowDependency[_] =>visit(narrowDep.rdd)}}
也就是说:
遇到
NarrowDependency:继续往父 RDD 走,仍然属于当前 Stage遇到
ShuffleDependency:停下来,创建或获取一个父ShuffleMapStage
这就是 Stage 切分的根本规则。
八、区别和联系
| 对比项 | 窄依赖 | 宽依赖 |
|---|---|---|
| 源码类型 | NarrowDependency | ShuffleDependency |
| 父子分区关系 | 子分区依赖少量父分区 | 多个父分区重分布到多个子分区 |
| 是否 shuffle | 否 | 是 |
| 是否切 Stage | 不切 | 切 |
| 是否可 pipeline | 可以 | 不可以跨 shuffle pipeline |
| 典型操作 | map / filter / mapPartitions | reduceByKey / groupByKey / repartition |
| join 特例 | partitioner 相同可窄依赖 | partitioner 不同走 shuffle |
九、实战:观察依赖类型
val rdd = sc.parallelize(1 to 10, 2)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,让你技术之路走得更稳、更快。
喜欢的点个关注。
夜雨聆风