当前位置: 首页 > news >正文

Spark 核心之 Application 和 Job 原理剖析

摘要:你是否清楚一个 Spark 应用中到底有几个 Job?为什么.collect()会触发 Job 而.map()不会?Application、Job、Stage、Task 之间的关系到底是什么?本文从四层执行层级全景图、Action 触发 Job 的源码链路、DAGScheduler 的 Stage 切分规则、多 Job 应用实战示例四个维度,配合 1 张原创深色架构图 + 完整源码分析,带你彻底理清 Spark 的执行层级体系。

关键词:Spark Application, Job, Stage, Task, Action, DAGScheduler, SparkContext, 执行层级


一、开篇:你写的一个 Application 到底有几个 Job?

先看一段代码,你能准确说出它会产生几个 Job 吗?

vallines=sc.textFile("hdfs:///data/words.txt")// Transformationvalwords=lines.flatMap(_.split(" "))// Transformationvalpairs=words.map((_,1))// Transformationvalcounts=pairs.reduceByKey(_+_)// Transformationcounts.collect()// Action 1 → Job 0counts.count()// Action 2 → Job 1counts.saveAsTextFile("hdfs:///output")// Action 3 → Job 2

答案:3 个 Job。因为每个 Action 算子都会触发一个新的 Job。而textFileflatMapmapreduceByKey都是 Transformation——它们只是构建 DAG,不触发任何计算。


二、执行层级全景图

2.1 四层模型

Application (SparkContext) │ ├── Job-0 (调用 collect() 触发) │ ├── Stage 0 (ShuffleMapStage: 2 Tasks) │ └── Stage 1 (ResultStage: 3 Tasks) │ ├── Job-1 (调用 count() 触发) │ ├── Stage 2 (ShuffleMapStage: 2 Tasks) │ └── Stage 3 (ResultStage: 3 Tasks) │ └── Job-2 (调用 saveAsTextFile() 触发) ├── Stage 4 (ShuffleMapStage: 2 Tasks) └── Stage 5 (ResultStage: 3 Tasks)

三、Action 触发 Job 的源码链路 🔥

// 源码:RDD.scala - collect()defcollect():Array[T]=withScope{valresults=sc.runJob(this,(iter:Iterator[T])=>iter.toArray)results.flatten}// 源码:RDD.scala - count()defcount():Long=sc.runJob(this,Utils.getIteratorSize _).sum// 源码:RDD.scala - saveAsTextFile()defsaveAsTextFile(path:String):Unit={// ... 内部最终调用 sc.runJob()}// 核心链路// Action → sc.runJob() → DAGScheduler.runJob()// → DAGScheduler.handleJobSubmitted()// → 创建 ActiveJob → Stage 切分 → submitStage()

3.1 哪些算子是 Action?

算子返回类型说明
collect()Array[T]拉取所有数据到 Driver
count()Long计数
take(n)Array[T]取前 n 个
reduce(f)T聚合
foreach(f)Unit遍历
saveAsTextFile()Unit保存到文件
first()T取第一个

反直觉点reduceByKey不是 Action!它是 Transformation,触发 Shuffle 但不触发 Job。


四、Stage 切分规则

// 源码:DAGScheduler.scalaprivatedefgetMissingParentStages(stage:Stage):List[Stage]={stage.rdd.dependencies.flatMap{caseshufDep:ShuffleDependency[_,_,_]=>// WideDep → 切分 → 创建父 ShuffleMapStagegetOrCreateShuffleMapStage(shufDep,stage.firstJobId)case_=>Nil// NarrowDep → 不切分,保持在同一 Stage}.toList}

一句话规则:遇到 ShuffleDependency (Wide Dependency) 即切分 Stage。


五、多 Job 实战示例

valrdd=sc.parallelize(1to1000,4)// 4 Partitions// Job 1: countprintln(s"Count:${rdd.count()}")// Action → Job-1// Job 2: collectvalarr=rdd.collect()// Action → Job-2 (无 Shuffle, 1 Stage)// Job 3: saverdd.saveAsTextFile("hdfs:///output")// Action → Job-3// 总计: 3 个 Job, 3×1=3 个 ResultStage
// 带 Shuffle 的场景valrdd=sc.parallelize(1to1000,4).map(x=>(x%10,x))// Narrow.groupByKey()// Wide! Stage 边界rdd.count()// Action → Job-1: Stage 0 (ShuffleMapStage) + Stage 1 (ResultStage)

六、Application/Job/Stage/Task 对比表

层级定义触发条件数量
ApplicationSparkContext 实例spark-submit1
Job一个 Action 的完整计算Action 算子1~N
StageShuffle 边界切分的计算阶段ShuffleDependency每个 Job 1~N
Task处理一个 Partition 的最小单元Stage 内 Partition 数每个 Stage 1~N

七、总结

要点总结
层级关系1 App = N Jobs = N×M Stages = N×M×P Tasks
Job 触发每个 Action 算子调用 sc.runJob() 创建新 Job
Stage 切分遇到 ShuffleDependency 即切分
Task 生成Stage 最后一个 RDD 的 Partition 数量 = Task 数

金句:Transformation 是"画图纸"(构建 DAG),Action 是"按下启动键"(触发 Job)。一个 Application 可以画无数张图纸,但只有按下的启动键才算数。


作者:starzy | AI Data Engineer / 大数据技术实践者
博客:blog.starzy.cn | GitHub:starzy1990.github.io
专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践

http://www.jsqmd.com/news/1332421/

相关文章:

  • 【网管运维助手】批量Ping与端口扫描工具,网络管理员排查故障的得力助手
  • Java面试实战:Spring Boot与Docker核心技术解析
  • H3C交换机开局配置:Console密码与SSH安全访问全解析
  • 学 Simulink—— 混合励磁同步电机(HESM)励磁电流与转矩协调控制仿真
  • Meta-Orchestrator:用事件驱动与动态DAG构建智能协同编程代理系统
  • Android进程被杀全链路分析:从am_proc_died到内存泄漏排查实战
  • 高端瓷砖十大品牌权威解读:金丝玉玛以K金工艺领跑行业新高度
  • iFakeLocation终极指南:免费iOS虚拟定位工具完整使用教程
  • 找无天价违约金合同透明的抖音直播公会新人避坑攻略 - 甄选测评馆
  • 基于CNN与Python的人脸情绪识别:从数据预处理到实时系统搭建
  • 大模型稳定输出JSON的工程化解决方案:从提示词到后处理全链路实践
  • Unity音频优化实战:基于Audio Mixer构建专业级BGM系统
  • 抖音批量下载工具完整指南:如何轻松保存无水印视频与音乐
  • 2026年武汉正规知名的钢管租赁公司甄选指南:如何择优避开租赁陷阱? - geo交流
  • 文本分块:RAG系统的“黄金切割术“——从暴力拆解到智能感知的演进之路
  • 【读书笔记】《如何快速了解一个行业》
  • BlenderKit插件:三步解决3D建模资产获取难题,零门槛提升创作效率
  • GULP:宇宙的演化---第2章光速不变原理
  • 2026年气密性检测设备工厂有哪些?推荐广州岳信仪器 - 汇聚至此
  • OpenAI收购Python工具,开发者慌了?
  • TCP四次挥手详解:从状态机到异常排查与编程实践
  • MistyR空间转录组分析:量化细胞间空间依赖性的统计框架与实战
  • Arduino入门全攻略:从零搭建环境到实战项目开发
  • CnOpenData 管理层讨论与分析基本信息表
  • Java构建棋牌H5:高效开发实战指南
  • 小熊猫Dev-C++:如何用5分钟搭建高效的C++学习环境
  • OpenClaw替代方案:生物信息学AI工具链迁移与成本优化实战
  • 终极指南:为qBittorrent配置搜索插件,一站式搞定资源查找
  • 虚拟电厂“牌照“及准入门槛
  • 2026.08.03