乐于分享
好东西不私藏

【极速学习spark源码】DAGScheduler:Stage 是如何划分的

【极速学习spark源码】DAGScheduler:Stage 是如何划分的

基于 Spark 4.2,分支 branch-4.2

阅读时长约 12 分钟 · 进阶


背景

前几篇文章反复提到“Stage”。大多数 Spark 用户知道 DAGScheduler 负责把 Job 切分成 Stage,但不知道怎么切

规则只有一条:

遇到 ShuffleDependency 就切。窄依赖的全部算子放在同一个 Stage 里。


一、DAGScheduler 类定义和核心字段

DAGScheduler.scala:124

private[spark] class DAGScheduler(    private[scheduler] val sc: SparkContext,    private[scheduler] val taskScheduler: TaskScheduler,    listenerBus: LiveListenerBus,    mapOutputTracker: MapOutputTrackerMaster,    blockManagerMaster: BlockManagerMaster,    env: SparkEnv,    clock: Clock = new Clock())

它是基于事件循环的单线程架构:

eventProcessLoop.post(JobSubmitted(...))

所有操作通过 post event 异步执行,避免了同步调度时的并发问题。


二、从 runJob 到 handleJobSubmitted

SparkContext.runJobSparkContext.scala:2481)最终调用 dagScheduler.runJobDAGScheduler.scala:1043),它又调用 submitJobDAGScheduler.scala:984)。

submitJob 最后:

eventProcessLoop.post(JobSubmitted(...))

事件循环调度到 handleJobSubmittedDAGScheduler.scala:1400

private[scheduler] def handleJobSubmitted(    jobId: Int, finalRDD: RDD[_], func: ...,    partitions: Array[Int], callSite: ..., listener: ...,    properties: Properties): Unit = {  var finalStage: ResultStage = null  finalStage = createResultStage(finalRDD, func, partitions, jobId, callSite)  val job = new ActiveJob(jobId, finalStage, callSite, listener, artifacts, properties)  submitStage(finalStage)  submitWaitingStages()}

三、createResultStage:从最终 RDD 反向推

DAGScheduler.scala:704

private def createResultStage(    rdd: RDD[_], func: (TaskContext, Iterator[_]) => _,    partitions: Array[Int], jobId: Int, callSite: CallSite): ResultStage = {  val (shuffleDeps, resourceProfiles) = getShuffleDependenciesAndResourceProfiles(rdd)  val parents = getOrCreateParentStages(shuffleDeps, jobId)  val stage = new ResultStage(id, rdd, func, parents, ...)  stageIdToStage(stage.id) = stage  stage}

这里三个关键步:

  1. getShuffleDependenciesAndResourceProfiles(rdd) — 找 rdd 的直接 shuffle 父依赖

  2. getOrCreateParentStages(shuffleDeps, jobId) — 为每个 shuffle 依赖创建或获取父 stage

  3. new ResultStage(id, rdd, func, parents, ...) — 创建 ResultStage


四、getShuffleDependenciesAndResourceProfiles:找 shuffle 边界

DAGScheduler.scala:801

private def getShuffleDependenciesAndResourceProfiles(    rdd: RDD[_]): (HashSet[ShuffleDependency[_, _, _]], HashSet[ResourceProfile]) = {  val parents = new HashSet[ShuffleDependency[_, _, _]]  val resourceProfiles = new HashSet[ResourceProfile]  val visited = new HashSet[RDD[_]]  traverseParentRDDsWithinStage(rdd, visited, parents, resourceProfiles)  (parents, resourceProfiles)}

这个方法遍历 RDD 依赖图,但:

  • 遇到 ShuffleDependency 就停下,把它加入 parent stage 集合

  • 遇到 NarrowDependency 就继续往父 RDD 走

这就保证了:同一个 Stage 里只有窄依赖,shuffle 边界被隔离到 Stage 之间。


五、createShuffleMapStage:为 shuffle 创建中间 stage

DAGScheduler.scala:574

def createShuffleMapStage[K, V, C](    shuffleDep: ShuffleDependency[K, V, C],    firstJobId: Int): ShuffleMapStage = {  val rdd = shuffleDep.rdd  val (shuffleDeps, resourceProfiles) = getShuffleDependenciesAndResourceProfiles(rdd)  val parents = getOrCreateParentStages(shuffleDeps, firstJobId)  val stage = new ShuffleMapStage(id, rdd, shuffleDep, parents, firstJobId, ...)  stageIdToStage(stage.id) = stage  stage.shuffleDep.shuffleIdToStage(stage.id) = stage  stage}

getOrCreateShuffleMapStage 在 DAGScheduler.scala:528:先查 shuffleIdToStage 看是否已有对应的 stage,没有再创建。

这保证了 shuffle 结果可以复用:如果两个 Job 依赖同一个 shuffle,就不用再算一次。


六、submitStage:递归提交父 Stage

DAGScheduler.scala:1540

private def submitStage(stage: Stage): Unit = {  val missing = getMissingParentStages(stage).sortBy(_.id)  if (missing.isEmpty) {    submitMissingTasks(stage, jobId.get)  } else {    for (parent <- missing) {      submitStage(parent)    }    waitingStages += stage  }}

getMissingParentStages 在 DAGScheduler.scala:837,它比 getShuffleDependenciesAndResourceProfiles 更直接:

  • 遍历当前 stage 的所有 RDD 的依赖

  • 遇到 ShuffleDependency → 获取对应的 ShuffleMapStage,如果不可用则加入 missing parents

  • 遇到 NarrowDependency → 继续往上游遍历


七、Stage / ResultStage / ShuffleMapStage

Stage.scala:56

private[spark] abstract class Stage(    val id: Int,    val rdd: RDD[_],    val numTasks: Int,    val parents: List[Stage],    val firstJobId: Int,    val callSite: CallSite)

ResultStage.scala:30

class ResultStage(    id: Int, rdd: RDD[_], val func: (TaskContext, Iterator[_]) => _,    parents: List[Stage], ...)

ShuffleMapStage.scala:37

class ShuffleMapStage(    id: Int, rdd: RDD[_], val shuffleDep: ShuffleDependency[_, _, _],    parents: List[Stage], ...)

区别:

  • ShuffleMapStage 持有 shuffleDep,输出 MapStatus

  • ResultStage 持有 func,输出用户结果


八、handleTaskCompletion:完成后推进下游

DAGScheduler.scala:2203

Task 完成时,事件循环处理 handleTaskCompletion

ShuffleMapTask 完成(DAGScheduler.scala:2330):

  • 从 pendingPartitions 移除该 partition

  • 注册 map output 到 MapOutputTracker

  • 当 stage 全部完成后调 processShuffleMapStageCompletion(shuffleStage)DAGScheduler.scala:2955

  • 最终调用 submitWaitingChildStages(shuffleStage)DAGScheduler.scala:1263

ResultTask 完成(DAGScheduler.scala:2275):

  • 标记对应 output id 完成

  • 全部完成后 job 结束


九、实战:观察 Stage 划分

val rdd1 = sc.parallelize(1 to 1004)val rdd2 = rdd1.map(_ * 2)          // 窄依赖val rdd3 = rdd2.map(_ + 1)          // 窄依赖val rdd4 = rdd3.groupBy(x => x % 2// 宽依赖val rdd5 = rdd4.mapValues(_.sum)    // 窄依赖val rdd6 = rdd5.count()             // action// 查看完整的血缘println(rdd5.toDebugString)

输出:

(4) MapPartitionsRDD[5] at mapValues at RDDtest03.scala:30 []
 |  ShuffledRDD[4] at groupBy at RDDtest03.scala:29 []
 +-(4) MapPartitionsRDD[3] at groupBy at RDDtest03.scala:29 []
    |  MapPartitionsRDD[2] at map at RDDtest03.scala:28 []
    |  MapPartitionsRDD[1] at map at RDDtest03.scala:27 []
    |  ParallelCollectionRDD[0] at parallelize at RDDtest03.scala:26 []
  • map / mapValues 在同一个 Stage 里 pipeline

  • groupBy 是 ShuffleDependency,产生 Stage 边界

  • rdd5.toDebugString 里缩进表示 Stage 划分

  • Web UI 会显示对应的 Stage 边界


十、总结

Stage 划分的完整流程:

rdd.action
  → SparkContext.runJob
    → DAGScheduler.submitJob
      → eventProcessLoop.post(JobSubmitted)
        → handleJobSubmitted
          → createResultStage(rdd)
            → getShuffleDependenciesAndResourceProfiles
              → traverseParentRDDsWithinStage
                → ShuffleDependency → 停,创建 ShuffleMapStage
                → NarrowDependency → 继续
          → submitStage
            → getMissingParentStages
            → submitStage(parent)   # 递归
  → submitMissingTasks
    → taskScheduler.submitTasks

三个总结:

  1. Stage 划分的唯一规则是 shuffle 边界。 窄依赖全 pipeline,一个 shuffle 切一个 Stage。

  2. DAGScheduler 是单线程事件循环。 所有 state 修改在事件处理线程中,不需要额外同步。

  3. ShuffleMapStage 输出可以复用。shuffleIdToStage 让相同 shuffle 的不需要重复计算。


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

喜欢的点个关注。