Spark 调度机制
概述
一个 Spark 应用从提交到跑完,涉及两层调度:作业间调度(FIFO/Fair)与作业内调度(Stage/Task)。本文讲清调度器选择、Stage 划分算法、Task 的数据本地性调度、推测执行与动态资源分配。
一、作业间调度:FIFO 与 Fair
1.1 FIFO 调度
| 特点 | 说明 |
|---|---|
| 先进先出 | 队列中第一个 Job 独占资源 |
| 简单 | 默认调度器 |
| 问题 | 大作业阻塞后续小作业 |
1.2 Fair 调度
| 特点 | 说明 |
|---|---|
| 公平分享 | 多个 Job 轮流占用资源 |
| 短作业友好 | 小作业快速获得资源启动 |
| 权重可配 | spark.scheduler.allocation.file 配置池权重 |
1.3 配置
bash
# 提交时指定
spark-submit --conf spark.scheduler.mode=FAIR ...
# 或
spark-submit --conf spark.scheduler.mode=FIFO ...| 场景 | 调度器 |
|---|---|
| 单作业批处理 | FIFO |
| 多用户共享集群 | Fair |
二、Stage 划分算法
2.1 核心逻辑
DAGScheduler 从 RDD 反向遍历,遇到宽依赖就切出一个新 Stage:
算法:
1. 从触发 Action 的 RDD 开始
2. 反向遍历父 RDD 依赖
3. 窄依赖:合并到当前 Stage
4. 宽依赖:当前 Stage 结束,父 RDD 开新 Stage
5. 直到数据源2.2 划分结果
rdd1 = textFile
rdd2 = rdd1.map(...) # 窄
rdd3 = rdd2.filter(...) # 窄
rdd4 = rdd3.groupByKey() # 宽 ← Stage 边界
rdd5 = rdd4.mapValues(...) # 窄
rdd5.count() # Action
Stage 0: rdd1 → rdd2 → rdd3(管道式)
Stage 1: rdd4 → rdd5(Shuffle 后)2.3 Stage 类型
| 类型 | 说明 |
|---|---|
| ShuffleMapStage | 产生 Shuffle 输出,供下游读取 |
| ResultStage | 最终输出结果的 Stage |
三、Task 调度
3.1 任务生成
Stage 内的 Task 数 = 分区数
每个分区生成一个 Task,提交给 TaskScheduler3.2 TaskScheduler 调度流程
1. TaskSet 提交到调度队列
2. 按本地性级别调度:
PROCESS_LOCAL → NODE_LOCAL → RACK_LOCAL → ANY
3. 每个级别等待 spark.locality.wait 时间
4. 超时降级到下一级别| 本地性级别 | 含义 |
|---|---|
| PROCESS_LOCAL | Task 所在 Executor 已有数据(最优) |
| NODE_LOCAL | 同节点其他 Executor 有数据 |
| RACK_LOCAL | 同机架 |
| ANY | 任意节点(网络传输) |
3.3 数据本地性调优
| 参数 | 说明 |
|---|---|
spark.locality.wait | 各级别等待时间,默认 3s |
spark.locality.wait.process | 单独配置 PROCESS_LOCAL 等待 |
本地性差(大量 ANY)时,可等待更久或优化数据分区与 Task 分布。
四、推测执行
4.1 原理
当某个 Task 运行时间远长于同 Stage 其他任务时,Spark 在另一 Executor 推测启动一个备份任务,谁先完成算谁:
| 参数 | 默认 | 说明 |
|---|---|---|
spark.speculation | false | 是否开启 |
spark.speculation.multiplier | 1.5 | 落后任务的倍率阈值 |
spark.speculation.quantile | 0.75 | 触发推测的任务完成比例 |
4.2 适用场景
| 场景 | 建议 |
|---|---|
| 节点异构、慢节点明显 | 开启推测执行 |
| 数据倾斜(个别任务本就慢) | 谨慎开启,会浪费资源 |
| 均匀集群 | 可关闭减少开销 |
注意:推测执行不是倾斜的解药,倾斜要治数据,否则备份任务一样慢。
五、动态资源分配
5.1 原理
根据负载动态增减 Executor,空闲时释放,繁忙时申请:
| 参数 | 默认 | 说明 |
|---|---|---|
spark.dynamicAllocation.enabled | false | 是否开启 |
spark.dynamicAllocation.initialExecutors | - | 初始 Executor 数 |
spark.dynamicAllocation.minExecutors | 0 | 最小数 |
spark.dynamicAllocation.maxExecutors | - | 最大数 |
spark.dynamicAllocation.executorIdleTimeout | 60s | 空闲释放时间 |
5.2 开启条件
- 集群运行 Shuffle 服务(external shuffle service)或 Executor 回收安全。
- 适合多租户共享集群、任务负载波动场景。
5.3 对比固定资源
| 方式 | 优点 | 缺点 |
|---|---|---|
| 固定 Executor | 资源稳定、可预测 | 浪费或不足 |
| 动态分配 | 按需伸缩、省资源 | 调度延迟、波动 |
六、调度完整链路
spark-submit → ClusterManager 分配 Driver
→ SparkContext 启动
→ 申请 Executor
→ Action 触发 Job
→ DAGScheduler 切 Stage
→ TaskScheduler 按本地性派发 Task
→ Executor 执行,Shuffle 处重分区
→ 结果回传,释放资源常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 小作业被大作业阻塞 | 切 Fair 调度器 |
| Task 本地性全是 ANY | 检查分区与数据位置,调整 locality.wait |
| 某 Task 特别慢 | 数据倾斜或慢节点,治理数据或开推测 |
| Executor 波动频繁 | 动态分配参数过小,调大 min/max 与空闲时间 |
| 提交后一直等待资源 | 队列资源不足,检查集群容量 |