资讯动态

Spark 核心之 Stage 和 Task 原理剖析

发布时间:2026/8/7 10:32:37 来源:尧图企业网站定制
摘要如果说 Job 是 Spark 的任务单Stage 就是施工阶段Task 就是每个工人的具体活。一个 Job 被 DAGScheduler 沿 Shuffle 边界切分为多个 Stage——前面的全是 ShuffleMapStage最后一个必须是 ResultStage。每个 Stage 的 Partition 数决定了 Task 数量ShuffleMapStage 产生 ShuffleMapTask写 Shuffle 文件ResultStage 产生 ResultTask直接返回结果。本文从 Stage 类型体系、DAG → Stage 切分源码、Task 生成与序列化、两种 Task 执行差异四个维度配合 1 张原创深色架构图 完整源码分析带你彻底看懂 Spark 最核心的执行引擎。关键词Spark Stage, ShuffleMapStage, ResultStage, ShuffleMapTask, ResultTask, DAGScheduler, Task 序列化, MapOutputTracker一、开篇Stage 和 Task 是什么关系先说结论Job 用户的一个 Action 操作 ├── Stage 0: ShuffleMapStage → 2 个 ShuffleMapTask └── Stage 1: ResultStage → 3 个 ResultTask概念定义数量StageShuffle 边界切分的计算阶段每个 Job 可有多个Task处理一个 Partition 的最小计算单元每个 Stage 可有多个ShuffleMapStage输出 Shuffle 中间文件的 StageJob 中除最后一个外的所有ResultStage输出最终结果的 Stage每个 Job 有且仅有一个二、Stage 与 Task 全景图三、Stage 切分从 RDD DAG 到 Stage3.1 核心源码// 源码DAGScheduler.scala - 创建 ResultStageprivatedefcreateResultStage(finalRDD:RDD[_],func:(TaskContext,Iterator[_])_,partitions:Array[Int],jobId:Int,callSite:CallSite):ResultStage{// 从 finalRDD 回溯 → 遇到 ShuffleDep → 创建 ShuffleMapStagevalparentsgetOrCreateParentStages(finalRDD,jobId)validnextStageId.getAndIncrement()newResultStage(id,finalRDD,func,partitions,parents,jobId,callSite)}// 递归获取父 StageprivatedefgetOrCreateParentStages(rdd:RDD[_],firstJobId:Int):List[Stage]{rdd.dependencies.flatMap{caseshufDep:ShuffleDependency[_,_,_]getOrCreateShuffleMapStage(shufDep,firstJobId)::Nilcase_Nil// NarrowDep 不切分}.toList}3.2 Stage 提交顺序// 递归提交先父后子privatedefsubmitStage(stage:Stage):Unit{valmissinggetMissingParentStages(stage).sortBy(_.id)if(missing.isEmpty){submitMissingTasks(stage,jobId.get)// 无缺失父 Stage → 执行}else{for(parent-missing)submitStage(parent)// 递归提交父 Stage}}四、Task 生成从 Stage 到 TaskSet// 源码DAGScheduler.scala - submitMissingTasks()privatedefsubmitMissingTasks(stage:Stage,jobId:Int):Unit{// 计算需要计算的 Partition跳过已完成的valpartitionsToComputestage.findMissingPartitions()// 为每个 Partition 创建一个 Taskvaltasks:Seq[Task[_]]stagematch{casestage:ShuffleMapStagepartitionsToCompute.map{idnewShuffleMapTask(stage.id,stage.rdd,stage.shuffleDep,...)}casestage:ResultStagepartitionsToCompute.map{idnewResultTask(stage.id,stage.rdd,stage.func,id,...)}}// 封装为 TaskSet提交给 TaskSchedulertaskScheduler.submitTasks(newTaskSet(tasks.toArray,stage.id,...))}Task 数量 Stage 最后一个 RDD 的 Partition 数量。五、两种 Stage 与两种 Task 对比5.1 ShuffleMapStage ShuffleMapTask// ShuffleMapTask.runTask() — 执行逻辑overridedefrunTask(context:TaskContext):MapStatus{valwriternewShuffleWriter(partition,shuffleDep)// ① 执行 RDD 算子链map/flatMap/filter...valiterrdd.iterator(partition,context)// ② 将结果写入 Shuffle 文件writer.write(iter)// ③ 返回 MapStatus文件位置 分区长度writer.stop(successtrue).get}5.2 ResultStage ResultTask// ResultTask.runTask() — 执行逻辑overridedefrunTask(context:TaskContext):U{// ① 执行 RDD 算子链valiterrdd.iterator(partition,context)// ② 将最终结果应用 func如 collect 的收集逻辑func(context,iter)// ③ 序列化结果 → StatusUpdate → Driver}5.3 对比表维度ShuffleMapStageResultStageTask 类型ShuffleMapTaskResultTask输出Shuffle 中间文件最终计算结果返回类型MapStatusU (泛型)一个 Job 中的数量0~N1唯一六、Task 序列化# 推荐 Kryo 序列化比 Java 快 10 倍--confspark.serializerorg.apache.spark.serializer.KryoSerializer--confspark.kryo.registrationRequiredtrue# 强制注册// 代码中注册 Kryo 类valconfnewSparkConf().set(spark.serializer,org.apache.spark.serializer.KryoSerializer).registerKryoClasses(Array(classOf[MyDataClass],classOf[MyModel]))为什么需要序列化Driver 端的 Task 对象包含 RDD 算子闭包需要跨网络发送到 Executor必须序列化为字节流。七、总结要点总结Stage 切分遇到 ShuffleDependency 即切分递归提交先父后子Task 生成每个 Partition → 一个 Task类型由 Stage 决定两种 StageShuffleMapStage写 Shuffle ResultStage返回结果序列化Task 闭包必须可序列化推荐 Kryo金句Stage 是 Spark 的流水线工位Task 是每个工位上的工人。Shuffle 就是工位之间的传送带——上一个工位写完下一个工位才能开始。作者starzy | AI Data Engineer / 大数据技术实践者博客blog.starzy.cn | GitHubstarzy1990.github.io专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价