欢迎光临

Spark 动态资源分配深度解析:从原理到生产调优的完整实战指南

在大规模数据处理场景中,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 个,以此类推,上限为
    1
    spark.dynamicAllocation.maxExecutors

  • Executor 释放:当某个 Executor 空闲时间超过
    1
    spark.dynamicAllocation.executorIdleTimeout

    (默认 60s)后,将其释放回资源管理器。释放前会先将其上面的 Shuffle 数据迁移到其他 Executor。

1.2 Shuffle 文件保留机制

这是动态分配最精妙的设计。在静态模式下,每个 Executor 的 Shuffle map 输出直接存在本地磁盘,reduce 阶段通过

1
BlockTransferService

远程拉取。但动态分配下,如果 reduce 阶段需要的 map 输出所在的 Executor 已被释放,这些数据将无法访问。

Spark 引入了 External Shuffle Service 机制来解决这一问题:

  • 每个节点运行一个独立的 Shuffle Service 进程(YARN 上的
    1
    YarnShuffleService

  • 该进程与 Executor 生命周期解耦,即使 Executor 被释放,Shuffle 文件仍可通过该服务访问
  • Reduce 任务拉取数据时,直接与 Shuffle Service 通信,而非与原始 Executor 通信

这意味着启用动态分配的前提条件是必须配置 External Shuffle Service,否则 Spark 只能缓存被移除 Executor 的 Shuffle 数据,严重消耗内存。

二、核心配置参数详解

要正确使用动态资源分配,需要理解每个参数的含义和影响。下面按照功能分组详细说明。

2.1 基础开关参数

参数 默认值 说明
1
spark.dynamicAllocation.enabled
false 是否启用动态分配,必须显式设为 true
1
spark.shuffle.service.enabled
false 是否启用 External Shuffle Service,动态分配下必须为 true
1
spark.dynamicAllocation.minExecutors
0 最少保留的 Executor 数量,即使全部空闲也不释放
1
spark.dynamicAllocation.maxExecutors
Infinity 最多申请的 Executor 数量,防止失控
1
spark.dynamicAllocation.initialExecutors
minExecutors 初始 Executor 数量,若未设置则取 minExecutors

2.2 超时与调度参数

参数 默认值 说明
1
spark.dynamicAllocation.executorIdleTimeout
60s Executor 空闲多久后被释放
1
spark.dynamicAllocation.cachedExecutorIdleTimeout
Infinity 缓存了 RDD 数据的 Executor 空闲多久后释放
1
spark.dynamicAllocation.schedulerBacklogTimeout
1s 有 pending task 多久后开始申请新 Executor
1
spark.dynamicAllocation.sustainedSchedulerBacklogTimeout
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%

这个案例说明:动态分配可能略微增加单作业延迟,但整体集群效率大幅提升,在多租户场景下尤其值得。

八、最佳实践总结

最后,将动态资源分配的最佳实践总结如下:

  1. 始终启用 External Shuffle Service(YARN 模式)或 Shuffle Tracking(K8s 模式),这是动态分配正常工作的前提。
  2. 设置合理的 min/max 边界。minExecutors 保证最低处理能力,maxExecutors 防止资源失控。
  3. initialExecutors 接近平均负载,避免启动阶段的资源瓶颈。
  4. executorIdleTimeout 设置 120-300 秒,避免多 Stage 作业中 Executor 反复创建销毁。
  5. 与 AQE 协同使用,在分区合并和 Executor 调整两个层面实现精细化的资源适配。
  6. 持续监控,通过 Spark UI 和日志分析调优效果,定期调整参数。
  7. 多租户场景下设置队列限制,配合 YARN 队列容量或 K8s ResourceQuota 防止单作业过度占用资源。

Spark 动态资源分配是一项成熟且经过大规模验证的特性,合理配置后可以在几乎不牺牲单作业性能的前提下,将集群资源利用率提升一倍以上。掌握其原理和调优方法,是每个 Spark 工程师的必修课。

【本站文章皆为原创,未经允许不得转载】:汤不热吧 » Spark 动态资源分配深度解析:从原理到生产调优的完整实战指南
分享到: 更多 (0)