Spark 核心之 Stage 和 Task 原理剖析

发布时间:2026/8/7 10:32:44

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 · 大数据架构 · 数据工程实践

相关新闻

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

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

2026/8/7 10:32:44

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

大气层整合包:Switch破解新手的终极一站式解决方案

大气层整合包:Switch破解新手的终极一站式解决方案

2026/8/7 10:32:43

大气层整合包:Switch破解新手的终极一站式解决方案 【免费下载链接】Atmosphere-stable 大气层整合包系统稳定版 项目地址: https://gitcode.com/gh_mirrors/at/Atmosphere-stable 还在为Switch破解的繁琐步骤和兼容性问题头疼吗?大气层整合包系统…

Java软件开发面试题小结(一)

Java软件开发面试题小结(一)

2026/8/7 10:22:43

1.mysql为什要用B树?主要核心归结为两点:“减少磁盘 I/O” 和 “高效范围查询”。简要理由如下:磁盘读写代价低:B树存的数据地址,不是数据本身,数据多层级少,B 树的节点大小固定且与磁盘页(Page…

免费解锁9大网盘直链下载:本地化工具实现高速下载新体验

免费解锁9大网盘直链下载:本地化工具实现高速下载新体验

2026/8/7 12:12:48

免费解锁9大网盘直链下载:本地化工具实现高速下载新体验 【免费下载链接】Online-disk-direct-link-download-assistant 一个基于 JavaScript 的网盘文件下载地址获取工具。基于【网盘直链下载助手】修改 ,支持 百度网盘 / 阿里云盘 / 中国移动云盘 / 天…

【多模态】13-基于Gemini的半结构化图像检索系统分析

【多模态】13-基于Gemini的半结构化图像检索系统分析

2026/8/7 12:12:48

1. 案例目标本案例演示如何使用Gemini Pro Vision模型从图像中提取结构化信息,并结合向量数据库实现语义搜索和元数据过滤的自动检索系统。具体目标包括:从收据图像中提取结构化信息(公司、日期、地址、总额等)将提取的结构化信息…

如何用Python一键完整下载任何网站到本地?终极离线浏览解决方案

如何用Python一键完整下载任何网站到本地?终极离线浏览解决方案

2026/8/7 12:12:48

如何用Python一键完整下载任何网站到本地?终极离线浏览解决方案 【免费下载链接】WebSite-Downloader A website downloader written with Python 项目地址: https://gitcode.com/gh_mirrors/web/WebSite-Downloader 你是否曾遇到过网络不稳定却急需访问重要…

基于差动磁场阵列的矿热炉电极毫米级高精度定位技术解析

基于差动磁场阵列的矿热炉电极毫米级高精度定位技术解析

2026/8/7 12:12:48

1. 项目概述与核心价值 在矿热炉这种大型工业电炉的生产现场,电极位置的精确控制一直是个老大难问题。炉内是上千度的高温熔池,电极在高温、强腐蚀、粉尘弥漫的环境下工作,传统的机械式或光学式检测手段在这里基本失灵。电极位置哪怕偏差几厘…

Python libvirt API开发实战:探索KVM虚拟化的编程之路

Python libvirt API开发实战:探索KVM虚拟化的编程之路

2026/8/7 12:12:48

Python libvirt API开发实战:探索KVM虚拟化的编程之路 在云计算与虚拟化技术蓬勃发展的当下,KVM(Kernel-based Virtual Machine)凭借其高效、稳定且开源的特性,成为了众多企业和开发者构建虚拟化环境的优选方案。而Pyt…

网络基础2(二)

网络基础2(二)

2026/8/7 12:02:47

1.HTTP协议下面进入HTTP协议部分,先来谈谈简单的预备知识:浏览器中输入的东西我们一般把它称为域名。根据我们目前学到的知识,客户端想访问服务端在技术上只需要知道IP和端口号就可以访问服务。实际上日常生活中我们并不使用IP地址&#xff0…

ncmdumpGUI:一键解锁网易云音乐ncm文件的终极解决方案

ncmdumpGUI:一键解锁网易云音乐ncm文件的终极解决方案

2026/8/6 19:19:00

ncmdumpGUI:一键解锁网易云音乐ncm文件的终极解决方案 【免费下载链接】ncmdumpGUI C#版本网易云音乐ncm文件格式转换,Windows图形界面版本 项目地址: https://gitcode.com/gh_mirrors/nc/ncmdumpGUI 你是否曾经从网易云音乐下载了心爱的歌曲&am…

分布式配置中心选型实战:Nacos与Consul在创业场景下的对比

分布式配置中心选型实战:Nacos与Consul在创业场景下的对比

2026/8/5 6:02:27

分布式配置中心选型实战:Nacos与Consul在创业场景下的对比工程导读:本文深入讨论 分布式配置中心选型实战:Nacos与Consul在创业场景下的对比 在生产工程实践中的核心落地方案。基于 分布式架构与微服务设计 视角,剖析实际痛点、架…

MoneyPrinterPlus实战指南:AI视频批量生成与自动化发布完整解决方案

MoneyPrinterPlus实战指南:AI视频批量生成与自动化发布完整解决方案

2026/8/5 8:19:55

MoneyPrinterPlus实战指南:AI视频批量生成与自动化发布完整解决方案 【免费下载链接】MoneyPrinterPlus AI一键批量生成各类短视频,自动批量混剪短视频,自动把视频发布到抖音,快手,小红书,视频号上,赚钱从来没有这么容易过! 支持本地语音模型chatTTS,fasterwhisper,…

CAD图库管理:从文件归档到设计资产管理的效率革命

CAD图库管理:从文件归档到设计资产管理的效率革命

2026/8/7 0:02:15

你肯定遇到过这种情况:打开一个老项目,想找某个特定的图块——比如一个标准的门、一个特定的设备符号,或者一个公司logo。你记得它就在某个DWG文件里,或者曾经从某个同事那里拷来过。于是,你开始在一堆命名混乱的文件夹…

5分钟掌握Wand-Enhancer:2026年终极WeMod专业版免费解锁指南

5分钟掌握Wand-Enhancer:2026年终极WeMod专业版免费解锁指南

2026/8/7 0:02:15

5分钟掌握Wand-Enhancer:2026年终极WeMod专业版免费解锁指南 【免费下载链接】Wand-Enhancer Advanced UX and interoperability extension for Wand (WeMod) app 项目地址: https://gitcode.com/GitHub_Trending/we/Wand-Enhancer Wand-Enhancer是一款功能强…

“Quality Control(质量控制)”在软件工程中通常指通过一系列活动确保软件产品符合预定的质量标准和用户需求

“Quality Control(质量控制)”在软件工程中通常指通过一系列活动确保软件产品符合预定的质量标准和用户需求

2026/8/7 0:02:15

“Quality Control(质量控制)”在软件工程中通常指通过一系列活动确保软件产品符合预定的质量标准和用户需求。而“软件测试”是质量控制的关键手段之一,属于QC范畴下的具体实践,其目标是发现缺陷、验证功能正确性、评估软件质量属…

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

2026/8/6 5:43:30

一天写完毕业论文在2026年已不再是天方夜谭。2026年最炸裂、实测能大幅提速的AI论文写作工具,覆盖选题构思、文献整理、内容生成、格式排版等核心场景,真正帮你高效搞定论文难题。 一、全流程王者:一站式搞定论文全链路(一天定稿首…

导师推荐!2026最新AI论文工具测评与实用推荐

导师推荐!2026最新AI论文工具测评与实用推荐

2026/8/7 8:02:42

2026年真正好用的AI论文工具,核心看生成的论文质量、低AI味、格式正确、学术适配四大指标。综合实测,千笔AI、ThouPen、豆包、DeepSeek、Grammarly 是当前最值得推荐的梯队,覆盖从免费到付费、从中文到英文、从文科到理工的全场景需求。 一、…

告别游戏崩溃:XCOM 2模组管理器的智能革命

告别游戏崩溃:XCOM 2模组管理器的智能革命

2026/8/4 15:11:03

告别游戏崩溃:XCOM 2模组管理器的智能革命 【免费下载链接】xcom2-launcher The Alternative Mod Launcher (AML) is a replacement for the default game launchers from XCOM 2 and XCOM Chimera Squad. 项目地址: https://gitcode.com/gh_mirrors/xc/xcom2-lau…