基于 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.runJob(SparkContext.scala:2481)最终调用 dagScheduler.runJob(DAGScheduler.scala:1043),它又调用 submitJob(DAGScheduler.scala:984)。
submitJob 最后:
eventProcessLoop.post(JobSubmitted(...))事件循环调度到 handleJobSubmitted。DAGScheduler.scala:1400:
private[scheduler] def handleJobSubmitted(jobId: Int, finalRDD: RDD[_], func: ...,partitions: Array[Int], callSite: ..., listener: ...,properties: Properties): Unit = {var finalStage: ResultStage = nullfinalStage = 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) = stagestage}
这里三个关键步:
getShuffleDependenciesAndResourceProfiles(rdd)— 找 rdd 的直接 shuffle 父依赖getOrCreateParentStages(shuffleDeps, jobId)— 为每个 shuffle 依赖创建或获取父 stagenew 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.rddval (shuffleDeps, resourceProfiles) = getShuffleDependenciesAndResourceProfiles(rdd)val parents = getOrCreateParentStages(shuffleDeps, firstJobId)val stage = new ShuffleMapStage(id, rdd, shuffleDep, parents, firstJobId, ...)stageIdToStage(stage.id) = stagestage.shuffleDep.shuffleIdToStage(stage.id) = stagestage}
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,输出MapStatusResultStage持有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 100, 4)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 里 pipelinegroupBy是 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
三个总结:
Stage 划分的唯一规则是 shuffle 边界。 窄依赖全 pipeline,一个 shuffle 切一个 Stage。
DAGScheduler 是单线程事件循环。 所有 state 修改在事件处理线程中,不需要额外同步。
ShuffleMapStage 输出可以复用。
shuffleIdToStage让相同 shuffle 的不需要重复计算。
每天花费5分钟学习spark,让你技术之路走得更稳、更快。
喜欢的点个关注。
夜雨聆风