尧图网络 高端网站定制 · 原创设计
免费咨询热线
400-888-6620
免费获取方案
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 · 大数据架构 · 数据工程实践
RELATED

相关推荐

开源项目零文档上手指南:从“大同生日快乐”到实战评估方法论

开源项目零文档上手指南:从“大同生日快乐”到实战评估方法论

1. 先搞清楚“大同生日快乐”到底在说什么 看到“大同生日快乐”这个标题,很多人第一反应可能是某个城市、某个品牌或者某个人的生日祝福。但在技术博客的语境下,它更可能指向一个特定的项目、一个代码库、一个数据集,或者一个与“大同”相关…

📅 2026/9/16 3:41:54
SQL数据更新操作详解与性能优化实战

SQL数据更新操作详解与性能优化实战

1. SQL数据更新操作的核心价值与场景 数据库的更新操作是每个开发者必须掌握的日常技能。在实际业务中,我们经常遇到需要批量修改用户状态、调整商品价格、修复错误数据等场景。以电商平台为例,双十一活动结束后,需要将数百万商品的"促销…

📅 2026/9/15 14:54:42
60.ABAP MARA 物料主数据 ALV 报表完整实现(带源码 + 字段目录)

60.ABAP MARA 物料主数据 ALV 报表完整实现(带源码 + 字段目录)

摘要 本文面向具备一定编程基础、希望系统掌握SAP ABAP开发的技术人员,从SAP系统架构、ABAP语言特性、数据字典对象、Open SQL规范、ALV报表开发、模块化编程到性能调优与代码规范,构建一条从基础到精通的完整学习路径。全文以工程化落地为导向,提供可直接运行的ABAP代码示…

📅 2026/9/4 21:16:09
MORE NEWS

更多资讯

📰

AI增强卫星动力学仿真:C语言环境下轻量级误差补偿模型

1. 为什么卫星动力学仿真突然需要AI介入——从轨道预报误差说起我第一次在航天院所做轨道预报验证时,被一组数据震住了:用经典二体J2摄动模型跑7天轨道,位置误差就突破800米;换成高精度数值积分加30阶地球引力场模型,计…

📰

向日葵CLI实战攻略:用MCP打通自动化排障与标准化操作,TaoToken统一Key接入

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

📰

MiniMax-01技术报告解读(三)预训练:从数据配比到TaoToken统一API的工程复现

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

📰

智能数字版权保护系统架构设计:AI应用架构师如何用TaoToken统一Key打通多模型版权识别链路

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

📰

用 Telegram 远程操控本地 OpenCode:opencode-telegram-bot 实战指南(TaoToken 配置篇)

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

📰

GPT-Image 2.5实测:12种玩法让朋友圈素材全包圆

假期第一天,朋友圈里已经开始有人晒定位、晒日出、晒航空公司餐盒了。我本来没打算出门凑热闹,却在电脑前折腾了一整晚GPT-Image 2.5,硬是把假期朋友圈的“素材”提前包圆了。你别说,这版图片生成模型跟之前的工具完全不是一个手感…

TODAY

今日更新

THIS WEEK

本周精选

THIS MONTH

本月热门

读完文章,想聊聊您的网站?

告诉我们您的行业与需求,资深顾问一对一梳理方案与报价,全程免费。

📞 💬