在大规模数据处理场景中,Spark 集群的资源利用率一直是工程师们关注的核心问题。传统的静态资源分配模式下,每个 Executor 从作业启动到结束都占用固定资源,即使处于空闲状态也无法释放,导致大量资源浪费。Spark 动态资源分配(Dynamic Allocation)机制应运而生,它能够根据工作负载实时增减 Executor 数量,在保证作业性能的前提下显著提升集群资源利用率。本文将从底层原理、核心配置、生产调优策略三个维度,全面解析这一关键特性。
一、动态资源分配的核心原理
Spark 动态资源分配的核心思想很简单:忙碌时申请更多 Executor,空闲时及时释放。但在分布式系统中实现这一目标并不容易,需要解决两个关键问题:何时申请/释放 Executor,以及如何在释放 Executor 后保留其计算结果供后续任务使用。
1.1 Executor 申请与释放策略
Spark 通过
1 | ExecutorAllocationManager |
组件管理动态分配逻辑。其工作机制分为两个方向:
- Executor 申请:当存在待调度任务(pending task)且已有 Executor 数量不足以在合理时间内完成时,向资源管理器(YARN/K8s)请求新的 Executor。申请速率采用指数增长策略——首轮申请 1 个,若仍不能满足则下次申请 2 个,再下次 4 个,以此类推,上限为
1spark.dynamicAllocation.maxExecutors
。
- Executor 释放:当某个 Executor 空闲时间超过
1spark.dynamicAllocation.executorIdleTimeout
(默认 60s)后,将其释放回资源管理器。释放前会先将其上面的 Shuffle 数据迁移到其他 Executor。
1.2 Shuffle 文件保留机制
这是动态分配最精妙的设计。在静态模式下,每个 Executor 的 Shuffle map 输出直接存在本地磁盘,reduce 阶段通过
1 | BlockTransferService |
远程拉取。但动态分配下,如果 reduce 阶段需要的 map 输出所在的 Executor 已被释放,这些数据将无法访问。
Spark 引入了 External Shuffle Service 机制来解决这一问题:
- 每个节点运行一个独立的 Shuffle Service 进程(YARN 上的
1YarnShuffleService
)
- 该进程与 Executor 生命周期解耦,即使 Executor 被释放,Shuffle 文件仍可通过该服务访问
- Reduce 任务拉取数据时,直接与 Shuffle Service 通信,而非与原始 Executor 通信
这意味着启用动态分配的前提条件是必须配置 External Shuffle Service,否则 Spark 只能缓存被移除 Executor 的 Shuffle 数据,严重消耗内存。
二、核心配置参数详解
要正确使用动态资源分配,需要理解每个参数的含义和影响。下面按照功能分组详细说明。
2.1 基础开关参数
| 参数 | 默认值 | 说明 | ||
|---|---|---|---|---|
|
false | 是否启用动态分配,必须显式设为 true | ||
|
false | 是否启用 External Shuffle Service,动态分配下必须为 true | ||
|
0 | 最少保留的 Executor 数量,即使全部空闲也不释放 | ||
|
Infinity | 最多申请的 Executor 数量,防止失控 | ||
|
minExecutors | 初始 Executor 数量,若未设置则取 minExecutors |
2.2 超时与调度参数
| 参数 | 默认值 | 说明 | ||
|---|---|---|---|---|
|
60s | Executor 空闲多久后被释放 | ||
|
Infinity | 缓存了 RDD 数据的 Executor 空闲多久后释放 | ||
|
1s | 有 pending task 多久后开始申请新 Executor | ||
|
schedulerBacklogTimeout | 持续有 pending task 时,两次请求间的间隔 |
关键理解:
1 | schedulerBacklogTimeout |
决定首次申请的触发时间,
1 | sustainedSchedulerBacklogTimeout |
决定后续每次申请的间隔。如果后者也设为 1s,则每秒都可能申请新 Executor,增长速度非常快。在生产中通常将
1 | sustainedSchedulerBacklogTimeout |
设为 5-10s,避免 Executor 数量暴涨。
2.3 一个典型的生产配置
1
2
3
4
5
6
7
8
9 spark.dynamicAllocation.enabled true
spark.shuffle.service.enabled true
spark.dynamicAllocation.minExecutors 2
spark.dynamicAllocation.maxExecutors 50
spark.dynamicAllocation.initialExecutors 5
spark.dynamicAllocation.executorIdleTimeout 120s
spark.dynamicAllocation.cachedExecutorIdleTimeout 300s
spark.dynamicAllocation.schedulerBacklogTimeout 5s
spark.dynamicAllocation.sustainedSchedulerBacklogTimeout 10s
这个配置的含义是:启动时分配 5 个 Executor,最小保留 2 个,最大不超过 50 个;Executor 空闲 120 秒后释放,缓存了 RDD 数据的 Executor 空闲 300 秒后释放;有 pending task 5 秒后开始申请新 Executor,之后每 10 秒评估一次。
三、YARN 模式下的完整部署
在 YARN 上部署动态分配需要额外配置 Shuffle Service,这是最常见的踩坑点。以下是完整的操作步骤。
3.1 配置 YARN Auxiliary Service
在每台 NodeManager 的
1 | yarn-site.xml |
中添加:
1
2
3
4
5
6
7
8
9
10
11
12
13
14 <configuration>
<property>
<name>yarn.nodemanager.aux-services</name>
<value>spark_shuffle,mapreduce_shuffle</value>
</property>
<property>
<name>yarn.nodemanager.aux-services.spark_shuffle.class</name>
<value>org.apache.spark.network.yarn.YarnShuffleService</value>
</property>
<property>
<name>spark.shuffle.service.port</name>
<value>7337</value>
</property>
</configuration>
然后将 Spark 的
1 | spark-[version]-yarn-shuffle.jar |
复制到每台 NodeManager 的 classpath 下:
1
2
3
4
5
6
7
8
9 # 方法一:放到 Hadoop 的 lib 目录
cp $SPARK_HOME/yarn/spark-3.5.0-yarn-shuffle.jar \
$HADOOP_HOME/share/hadoop/yarn/lib/
# 方法二:在 yarn-site.xml 中指定路径
<property>
<name>yarn.nodemanager.aux-services.spark_shuffle.classpath</name>
<value>/opt/spark/yarn/spark-3.5.0-yarn-shuffle.jar</value>
</property>
修改后必须重启所有 NodeManager,否则 Shuffle Service 不会生效。
3.2 验证 Shuffle Service
在 NodeManager 日志中确认:
1
2
3
4
5 grep "ShuffleService" $YARN_HOME/logs/nodemanager*.log
# 应看到类似输出:
# INFO service.AbstractService: Service org.apache.spark.network.yarn.YarnShuffleService started
# INFO shuffle.YarnShuffleService: YarnShuffleService started on port 7337
如果日志中没有出现 Shuffle Service 启动信息,最常见的原因是 jar 包路径不对或 NodeManager 未重启。
3.3 Spark Submit 提交示例
1
2
3
4
5
6
7
8
9
10
11
12
13
14 spark-submit \
--master yarn \
--deploy-mode cluster \
--conf spark.dynamicAllocation.enabled=true \
--conf spark.shuffle.service.enabled=true \
--conf spark.dynamicAllocation.minExecutors=2 \
--conf spark.dynamicAllocation.maxExecutors=100 \
--conf spark.dynamicAllocation.initialExecutors=10 \
--conf spark.dynamicAllocation.executorIdleTimeout=120s \
--conf spark.dynamicAllocation.schedulerBacklogTimeout=5s \
--conf spark.dynamicAllocation.sustainedSchedulerBacklogTimeout=10s \
--executor-memory 4g \
--executor-cores 2 \
your-app.jar
四、Kubernetes 模式下的动态分配
Spark 3.x 开始原生支持 Kubernetes 上的动态分配,但机制与 YARN 有所不同。K8s 模式不使用 NodeManager 上的 Shuffle Service,而是使用 Shuffle Tracking 机制。
4.1 Shuffle Tracking 原理
在 K8s 模式下,Driver 会跟踪每个 Executor 上的 Shuffle 数据状态。只有当某个 Executor 上的所有 Shuffle 数据都不再被需要时(即所有依赖该数据的 reduce 任务已完成),才会释放该 Executor。这避免了 External Shuffle Service 的依赖,但也意味着 Executor 可能保留更长时间。
4.2 K8s 模式配置
1
2
3
4
5
6
7
8
9
10
11 spark-submit \
--master k8s://https://k8s-api-server:6443 \
--deploy-mode cluster \
--conf spark.dynamicAllocation.enabled=true \
--conf spark.dynamicAllocation.shuffleTracking.enabled=true \
--conf spark.dynamicAllocation.shuffleTracking.timeout=300s \
--conf spark.dynamicAllocation.minExecutors=1 \
--conf spark.dynamicAllocation.maxExecutors=50 \
--conf spark.dynamicAllocation.executorIdleTimeout=120s \
--conf spark.kubernetes.container.image=spark:3.5.0 \
your-app.jar
关键参数
1 | spark.dynamicAllocation.shuffleTracking.timeout |
控制 Shuffle Tracking 的最长等待时间。如果超过这个时间仍有未完成的 Shuffle 依赖,Executor 仍会被强制释放,此时后续 Shuffle 读取会失败并触发重算。
五、生产环境调优实战
理解原理和配置只是第一步,真正的挑战在于生产环境中如何根据实际负载调优参数,避免常见的陷阱。
5.1 参数调优的黄金法则
法则一:
1 | initialExecutors |
应接近平均负载。如果平时需要 20 个 Executor,初始值设为 5 则意味着作业启动时有一段时间资源不足,大量任务等待调度。建议设置为你经验的 P50 值。
法则二:
1 | executorIdleTimeout |
不要设太小。60 秒的默认值在多 Stage 作业中会导致频繁的 Executor 申请和释放,产生大量开销。特别是 Stage 间隔时间超过 idle timeout 时,每个 Stage 开始都要重新申请 Executor。生产环境建议 120-300 秒。
法则三:
1 | sustainedSchedulerBacklogTimeout |
要大于
1 | schedulerBacklogTimeout |
。首次快速响应(5s),后续逐步扩展(10-15s),避免 Executor 爆炸式增长。
5.2 常见陷阱与解决方案
陷阱一:Executor 反复创建和销毁
症状:Spark UI 上 Executor 列表频繁变化,集群负载图呈锯齿状。这通常是因为
1 | executorIdleTimeout |
设置过小,Stage 间的空隙导致 Executor 被释放又重新申请。
解决方案:增大
1 | executorIdleTimeout |
,或设置合理的
1 | minExecutors |
兜底。
陷阱二:Shuffle 数据丢失导致 FetchFailedException
症状:日志中大量
1 | FetchFailedException |
,任务反复重试。原因是 External Shuffle Service 未正确配置或已崩溃。
解决方案:确认 NodeManager 上 Shuffle Service 正常运行,检查端口 7337 是否可达:
1 | nc -zv node-manager-host 7337 |
陷阱三:多租户场景下的资源饥饿
症状:动态分配导致大作业抢占过多资源,小作业长时间等待。YARN 队列容量上限可能被快速触及。
解决方案:合理设置
1 | maxExecutors |
,结合 YARN 队列的
1 | maximum-capacity |
和
1 | maximum-am-resource-percent |
进行限制。
5.3 与 Spark SQL 的协同优化
当动态分配与 Spark SQL 的 AQE(自适应查询执行)结合时,可以发挥更大的优化效果:
1
2
3
4
5
6
7
8
9
10 # AQE 配置
spark.sql.adaptive.enabled=true
spark.sql.adaptive.shuffle.targetPostShuffleInputSize=64m
spark.sql.adaptive.coalescePartitions.enabled=true
# 动态分配配置
spark.dynamicAllocation.enabled=true
spark.dynamicAllocation.minExecutors=5
spark.dynamicAllocation.maxExecutors=200
spark.dynamicAllocation.executorIdleTimeout=180s
工作流程:AQE 在运行时根据 Shuffle 数据量动态合并分区,减少小任务数量;动态分配则根据实际任务并行度调整 Executor 数量。两者协同,使得资源分配与计算需求精确匹配,避免了空跑 Executor 和排队等 Executor 两个极端。
六、监控与诊断
生产环境中,持续监控动态分配的效果至关重要。以下是关键的监控指标和诊断方法。
6.1 关键监控指标
| 指标 | 来源 | 理想值 |
|---|---|---|
| Executor 数量波动 | Spark UI Executors 页面 | 平滑增长,而非锯齿形 |
| 已完成 Stage 的 Executor 利用率 | Spark UI Stage 详情 | 平均大于 70% |
| Executor 申请频率 | Driver 日志(ExecutorAllocationManager) | 每个 Stage 增长不超过 2-3 次 |
| Shuffle Fetch 失败率 | Spark UI Task 详情 | 0% |
| 集群资源利用率 | YARN RM UI / K8s Dashboard | P50 大于 60% |
6.2 日志诊断技巧
通过 Driver 日志中
1 | ExecutorAllocationManager |
的输出可以追踪动态分配的决策过程:
1
2
3
4
5
6
7
8 # 查看 Executor 申请记录
grep "Requesting new executors" $SPARK_LOG_DIR/driver.log
# 查看 Executor 释放记录
grep "Removing executor" $SPARK_LOG_DIR/driver.log
# 查看空闲超时触发
grep "idle for" $SPARK_LOG_DIR/driver.log
如果发现日志中频繁出现
1 | Requesting new executors |
,说明
1 | sustainedSchedulerBacklogTimeout |
需要调大;如果频繁出现
1 | Removing executor |
,说明
1 | executorIdleTimeout |
需要调大。
6.3 通过 Spark Event Log 回溯分析
Spark 的事件日志(Event Log)记录了所有 Executor 的申请和释放事件,可以用于事后分析:
1
2
3
4
5
6
7
8 import json
with open("/path/to/spark-event-log", "rb") as f:
for line in f:
event = json.loads(line)
if event["Event"] in ["SparkListenerExecutorAdded", "SparkListenerExecutorRemoved"]:
print(f"{event['Event']}: executor {event.get('Executor ID', 'N/A')} "
f"at {event.get('Timestamp', 'N/A')}")
七、实战案例:从静态到动态的资源优化
以下是一个真实的数据管道优化案例。原始配置是静态分配 50 个 Executor,处理包含 ETL、聚合、写入三个阶段的作业:
1
2
3
4
5
6 # 原始静态配置
spark-submit --master yarn \
--num-executors 50 \
--executor-memory 8g \
--executor-cores 4 \
etl-pipeline.jar
分析 Spark UI 发现:ETL 阶段需要 50 个 Executor,但聚合阶段只需要 20 个,写入阶段仅需 10 个。大量 Executor 在聚合和写入阶段处于空闲状态。
切换到动态分配:
1
2
3
4
5
6
7
8
9
10
11
12
13
14 # 优化后的动态分配配置
spark-submit --master yarn \
--conf spark.dynamicAllocation.enabled=true \
--conf spark.shuffle.service.enabled=true \
--conf spark.dynamicAllocation.minExecutors=5 \
--conf spark.dynamicAllocation.maxExecutors=60 \
--conf spark.dynamicAllocation.initialExecutors=20 \
--conf spark.dynamicAllocation.executorIdleTimeout=180s \
--conf spark.dynamicAllocation.cachedExecutorIdleTimeout=600s \
--conf spark.dynamicAllocation.schedulerBacklogTimeout=5s \
--conf spark.dynamicAllocation.sustainedSchedulerBacklogTimeout=15s \
--executor-memory 8g \
--executor-cores 4 \
etl-pipeline.jar
优化结果:
- 作业总运行时间从 42 分钟增加到 44 分钟(增加约 5%,因 Executor 扩缩容开销)
- 集群平均资源利用率从 35% 提升到 72%
- 相同时间段可并发运行的作业数从 3 个提升到 7 个
- 月度 YARN 资源消耗下降 48%
这个案例说明:动态分配可能略微增加单作业延迟,但整体集群效率大幅提升,在多租户场景下尤其值得。
八、最佳实践总结
最后,将动态资源分配的最佳实践总结如下:
- 始终启用 External Shuffle Service(YARN 模式)或 Shuffle Tracking(K8s 模式),这是动态分配正常工作的前提。
- 设置合理的 min/max 边界。minExecutors 保证最低处理能力,maxExecutors 防止资源失控。
- initialExecutors 接近平均负载,避免启动阶段的资源瓶颈。
- executorIdleTimeout 设置 120-300 秒,避免多 Stage 作业中 Executor 反复创建销毁。
- 与 AQE 协同使用,在分区合并和 Executor 调整两个层面实现精细化的资源适配。
- 持续监控,通过 Spark UI 和日志分析调优效果,定期调整参数。
- 多租户场景下设置队列限制,配合 YARN 队列容量或 K8s ResourceQuota 防止单作业过度占用资源。
Spark 动态资源分配是一项成熟且经过大规模验证的特性,合理配置后可以在几乎不牺牲单作业性能的前提下,将集群资源利用率提升一倍以上。掌握其原理和调优方法,是每个 Spark 工程师的必修课。
汤不热吧