Spark核心概念深度解析:Job、Stage、Task三者关系

Spark核心概念深度解析:Job、Stage、Task三者关系

    • 一、三者关系全景图
      • 1.1 一句话定义
    • 二、Job(作业)深度解析
      • 2.1 什么是Job?
      • 2.2 Job的生命周期
      • 2.3 Job的监控
    • 三、Stage(阶段)深度解析
      • 3.1 什么是Stage?
      • 3.2 Stage的两种类型
      • 3.3 Stage的并行度
    • 四、Task(任务)深度解析
      • 4.1 什么是Task?
      • 4.2 Task的两种类型
      • 4.3 Task的数量计算
    • 五、Job、Stage、Task的关系公式
      • 5.1 数学关系
      • 5.2 实例推演
    • 六、为什么要划分Stage?
      • 6.1 Stage划分的核心原因
      • 6.2 流水线执行的威力
      • 6.3 没有Stage划分的后果
    • 七、三者在Spark UI中的体现
      • 7.1 Spark UI界面结构
      • 7.2 如何通过UI定位问题
    • 八、面试高频问题
      • Q1:Job、Stage、Task分别是什么?有什么关系?
      • Q2:Stage是如何划分的?
      • Q3:Task的数量如何确定?
      • Q4:为什么要划分Stage?
      • Q5:一个Stage中的Task是并行执行的吗?
    • 九、总结
      • 9.1 三层次执行单元
      • 9.2 核心公式
      • 9.3 记忆口诀

🌺The Begin🌺点点关注,收藏不迷路🌺

关键词:Spark作业、Stage划分、Task执行、DAG调度、并行度、执行单元

在Spark中,JobStageTask是三个层次分明的执行单元概念。理解它们的定义、关系以及划分逻辑,是掌握Spark分布式执行原理的关键,也是面试中的五星级考点。

今天,我们将深入剖析这三个核心概念,以及它们如何协同工作完成分布式计算任务。


一、三者关系全景图

Stage内部

Job内部

Application

Shuffle

Shuffle

Job 1
Action触发

Job 2
Action触发

Job 3
Action触发

Stage 1
窄依赖阶段

Stage 2
窄依赖阶段

Stage 3
窄依赖阶段

Task 1
分区1处理

Task 2
分区2处理

Task 3
分区3处理

1.1 一句话定义

概念 定义 触发方式
Job 一次Action操作触发的完整作业 由Action算子触发
Stage Job中按照Shuffle划分的任务组 由宽依赖切分
Task Stage中的最小执行单元 每个分区对应一个Task

二、Job(作业)深度解析

2.1 什么是Job?

Job是Spark中最大的执行单元,由一次Action操作触发。每个Action(如count()collect()saveAsTextFile()等)都会生成一个Job。

// 一个Application中可以包含多个Job
val rdd = sc.parallelize(Seq(1, 2, 3, 4, 5), 3)
// 第一个Action → 第一个Job
val count = rdd.count()  // Job 1
// 第二个Action → 第二个Job
val collected = rdd.map(_ * 2).collect()  // Job 2
// 第三个Action → 第三个Job
rdd.filter(_ > 2).saveAsTextFile("/output")  // Job 3

2.2 Job的生命周期

Job生命周期

Action触发

DAGScheduler构建DAG

划分为多个Stage

提交Task执行

返回结果/写入存储

2.3 Job的监控

// 在Spark UI中查看Job
// http://driver:4040/jobs/
// 或通过代码获取Job信息
val jobGroup = "my-job-group"
sc.setJobGroup(jobGroup, "My Important Job")
val result = rdd.count()
sc.clearJobGroup()

三、Stage(阶段)深度解析

3.1 什么是Stage?

Stage是Job中按照宽依赖(Shuffle)划分的任务组。每个Stage包含一组可以流水线执行的窄依赖转换。

Stage划分结果

Stage 0
textFile → flatMap → map

Stage 1
reduceByKey

Stage 2
map → filter → collect

Job划分Stage示例

窄依赖

窄依赖

宽依赖
reduceByKey

窄依赖

窄依赖

Action

textFile

flatMap

map

Stage边界

map

filter

collect

3.2 Stage的两种类型

// 1. ShuffleMapStage
// 最终输出是Shuffle文件,供下游Stage使用
val mapStageRDD = rdd.map(x => (x % 10, x)).reduceByKey(_ + _)
// 这里的reduceByKey会创建ShuffleMapStage
// 2. ResultStage
// 最终输出结果(返回Driver或写入存储)
val result = mapStageRDD.collect()
// collect()触发的最后一个Stage是ResultStage

Stage类型对比

Stage类型 作用 输出 位置
ShuffleMapStage 为Shuffle准备数据 Shuffle文件 中间Stage
ResultStage 生成最终结果 结果数据 最后一个Stage

3.3 Stage的并行度

// Stage的并行度由分区数决定
val rdd = sc.parallelize(Seq(1, 2, 3, 4, 5), 10)  // 10个分区
val mapped = rdd.map(_ * 2)  // 窄依赖,仍然10个分区
// 这个Stage会有10个Task并行执行
val result = mapped.collect()

四、Task(任务)深度解析

4.1 什么是Task?

Task是Spark中最小的执行单元,负责处理一个RDD分区的数据。每个Stage由多个Task组成,Task数量等于该Stage最后一个RDD的分区数。

Stage内部

Stage的RDD

分区1

Task 1

分区2

Task 2

分区3

Task 3

分区n

Task n

4.2 Task的两种类型

// 1. ShuffleMapTask
// 执行ShuffleMapStage中的任务,输出Shuffle数据
// 对应ShuffleMapStage中的Task
// 2. ResultTask
// 执行ResultStage中的任务,输出最终结果
// 对应ResultStage中的Task

Task类型对比

Task类型 所属Stage 输出 执行位置
ShuffleMapTask ShuffleMapStage Shuffle文件 Executor
ResultTask ResultStage 最终结果 Executor

4.3 Task的数量计算

// Task数 = 所有Stage的Task数之和
// 每个Stage的Task数 = 该Stage最后RDD的分区数
val rdd = sc.parallelize(Seq(1, 2, 3, 4, 5), 5)  // 5个分区
val mapped = rdd.map(_ * 2)  // 窄依赖,仍5个分区
val reduced = mapped.map(x => (x % 3, x)).reduceByKey(_ + _)  // 宽依赖
// 假设reduceByKey后分区数为默认的200(或通过参数指定)
// Stage 0: 5个ShuffleMapTask
// Stage 1: 200个ShuffleMapTask或ResultTask(取决于后续操作)
// 总Task数 = 5 + 200 = 205

五、Job、Stage、Task的关系公式

5.1 数学关系

// 一个Application → 多个Job
Application =(Job)
// 一个Job → 多个Stage
Job =(Stage)
// 一个Stage → 多个Task
Stage =(Task) = RDD分区数
// 总Task数 = ∑(每个Stage的分区数)
TotalTasks =(Stage.partitions)

5.2 实例推演

// 完整示例
val rdd = sc.textFile("hdfs://data/input.txt", 10)  // 10个分区
val words = rdd.flatMap(_.split(" "))               // 窄依赖
val pairs = words.map(word => (word, 1))            // 窄依赖
val counts = pairs.reduceByKey(_ + _, 20)            // 宽依赖,20个分区
val result = counts.filter(_._2 > 10)                // 窄依赖
result.saveAsTextFile("hdfs://data/output")          // Action
// Job划分:
// Stage 0: textFile + flatMap + map (10个分区 → 10个ShuffleMapTask)
// Stage 1: reduceByKey (20个分区 → 20个ShuffleMapTask)
// Stage 2: filter + saveAsTextFile (20个分区 → 20个ResultTask)
// 总Task数 = 10 + 20 + 20 = 50个Task

执行流程图

Job

Stage2

Stage1

Stage0

Shuffle

Shuffle

Shuffle

Task 2-1
分区1

Task 0-1
分区1

S1

Task 0-2
分区2

Task 0-10
分区10

Task 1-1
分区1

Task 1-2
分区2

Task 1-20
分区20

Task 2-2
分区2

Task 2-20
分区20


六、为什么要划分Stage?

6.1 Stage划分的核心原因

root(Stage划分的价值)

流水线执行

窄依赖在同一Stage

多个算子串行执行

减少中间数据落盘

Shuffle管理

宽依赖作为Stage边界

数据落盘与拉取

容错恢复点

并行度控制

不同Stage不同并行度

资源动态调整

任务调度优化

依赖关系清晰

任务依赖管理

容错恢复

以Stage为单位重试

检查点优化

6.2 流水线执行的威力

// 同一个Stage中的窄依赖可以流水线执行
val stage0RDD = rdd
  .flatMap(_.split(" "))     // 窄依赖
  .map(word => (word, 1))    // 窄依赖
  .filter(_._1.length > 3)   // 窄依赖
// 实际执行:一个Task中依次执行flatMap → map → filter
// 没有中间数据落盘,全部在内存管道中完成

流水线执行示意图

单个Task执行

flatMap

map

filter

读取分区数据

中间结果

中间结果

最终输出

6.3 没有Stage划分的后果

// 如果不划分Stage
// 1. 无法区分哪些算子可以流水线执行
// 2. 每个算子都独立执行,产生大量中间文件
// 3. 性能急剧下降
// 4. 容错恢复困难

七、三者在Spark UI中的体现

7.1 Spark UI界面结构

# Spark Web UI (http://driver:4040)
# Jobs 标签页
- 列出所有Job
- 每个Job的状态、持续时间、Stage数量
# Stages 标签页
- 列出所有Stage
- 每个Stage的Task数量、执行时间、Shuffle读写量
# Executors 标签页
- 每个Executor执行的Task数量
- 内存使用情况
# Storage 标签页
- 缓存RDD的信息

7.2 如何通过UI定位问题

现象 查看位置 可能问题
某个Job特别慢 Stages页面 查看Stage耗时
部分Task极慢 Stage详情 数据倾斜
Shuffle读写量大 Stage详情 Shuffle需要优化
Executor内存溢出 Executors页面 分区过大或数据倾斜

八、面试高频问题

Q1:Job、Stage、Task分别是什么?有什么关系?

  • Job:一次Action触发的完整作业
  • Stage:Job中按宽依赖划分的任务组
  • Task:Stage中处理单个分区的执行单元
  • 关系:1个Job包含多个Stage,1个Stage包含多个Task(Task数=分区数)

Q2:Stage是如何划分的?

:从后往前遍历RDD依赖链:

  • 遇到窄依赖,继续向前
  • 遇到宽依赖(Shuffle),创建新的Stage
  • 每个宽依赖都是Stage的边界

Q3:Task的数量如何确定?

:Task数 = 每个Stage的Task数之和

  • 每个Stage的Task数 = 该Stage最后RDD的分区数
  • 分区数可以通过repartition()coalesce()等算子调整

Q4:为什么要划分Stage?

  1. 流水线执行:同一Stage内窄依赖可流水线执行
  2. Shuffle管理:宽依赖作为边界,统一管理Shuffle
  3. 容错恢复:以Stage为单位进行任务重试
  4. 资源优化:不同Stage可设置不同并行度

Q5:一个Stage中的Task是并行执行的吗?

:是的,同一个Stage中的Task是并行执行的,每个Task处理一个分区。并行度受集群资源和Task数量共同影响。


九、总结

9.1 三层次执行单元

Task (任务)

Stage (阶段)

Job (作业)

Stage 1

Stage 2

Stage 3

Application (应用)

Job 1

Job 2

Job 3

Task 1

Task 2

Task 3

Task N

任务执行单元

Job

Stage

Task

9.2 核心公式

Job = Σ(Stage) = Σ(Σ(Task))
Task数 = Σ(每个Stage的RDD分区数)

9.3 记忆口诀

Action触发一个Job,Shuffle切分成Stage,分区个数定Task

理解这三者的关系,你就能深入理解Spark的分布式执行模型,为性能调优和问题排查打下坚实基础!


思考题:在Spark Structured Streaming中,还会沿用Job、Stage、Task这套模型吗?实时流处理的任务划分有什么不同?欢迎在评论区讨论!

在这里插入图片描述

🌺The End🌺点点关注,收藏不迷路🌺
© 版权声明

相关文章