Spark 核心之 Stage 和 Task 原理剖析

发布时间:2026/9/27 13:22:42
Spark 核心之 Stage 和 Task 原理剖析 摘要如果说 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 · 大数据架构 · 数据工程实践

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询