Spark 性能调优总结
概述
Spark 性能调优是一个系统性工程:先给足资源、再配好并行度,然后优化序列化与 GC,最后治理数据倾斜等特殊问题。本文按维度给出可直接落地的调优清单与排查路径。
一、资源分配
1.1 核心参数
| 参数 | 默认 | 说明 |
|---|---|---|
spark.executor.memory | 1g | Executor 堆内存 |
spark.executor.cores | 1 | Executor 核数 |
spark.executor.instances | - | Executor 数量 |
spark.driver.memory | 1g | Driver 内存 |
spark.executor.memoryOverhead | 10% | 堆外开销 |
1.2 分配原则
1. 内存:8-16g 常见,过大 GC 难、过小效率低
2. 核数:2-4 核/Executor,核数与内存匹配
3. 数量:总资源 ÷ 单 Executor 规格
4. Driver 内存:collect/广播大时加大| 场景 | 建议 |
|---|---|
| 计算密集 | 核多内存适中 |
| 缓存/迭代 | 内存大核少 |
| 数据倾斜 | 均匀分区 + 足够内存 |
二、并行度
2.1 原则
并行度(分区数)≈ Executor 核数总量 × 2~3
分区太多:调度开销大
分区太少:资源闲置2.2 调整方式
| 算子 | 场景 |
|---|---|
coalesce | 缩小分区(无 Shuffle) |
repartition | 增大分区(有 Shuffle) |
spark.sql.shuffle.partitions | SQL Shuffle 分区数 |
| AQE 合并 | 自动调节 |
2.3 判断标准
看 Spark UI:
Task 平均耗时短(< 数秒)→ 分区多
大量空 Task → 分区过多
任务等待资源 → 并行度超资源三、序列化
3.1 Java vs Kryo
| 对比 | Java | Kryo |
|---|---|---|
| 速度 | 慢 | 快 10 倍 |
| 大小 | 大 | 小(省 30-50%) |
| 注册 | 无需 | 推荐注册类 |
scala
sparkConf.set("spark.serializer",
"org.apache.spark.serializer.KryoSerializer")
sparkConf.registerKryoClasses(Array(classOf[MyClass]))3.2 注意
| 注意 | 说明 |
|---|---|
| 注册类 | 注册避免全限定名开销 |
| 反序列化失败 | 类不可变/无默认构造问题 |
| Shuffle 压缩 | 与压缩配合效果更佳 |
四、GC 调优
4.1 GC 问题信号
| 信号 | 含义 |
|---|---|
| GC 时间占比高 | 对象分配过多/内存小 |
| 频繁 Full GC | 老年代压力大 |
| 偶发 OOM | 内存不足或倾斜 |
4.2 优化手段
| 手段 | 说明 |
|---|---|
| 调大 Executor 内存 | 减少 GC 频率 |
| 用 Kryo | 减少对象与内存 |
| 堆外内存 | 减少堆内压力 |
| 避免大对象 | 减小广播与缓存粒度 |
| JVM 参数 | -XX:+UseG1GC 等 |
评估:
UI 看 GC Time
GC 时间 < 总执行 10% 为佳五、数据倾斜处理
5.1 识别
| 症状 | 排查 |
|---|---|
| 个别 Task 极慢 | UI 看 Task 耗时分布 |
| 某 Executor OOM | 单 key 数据过大 |
| Shuffle 阶段卡 | 倾斜分区 |
5.2 处理手段(按优先级)
| 手段 | 说明 |
|---|---|
| AQE 倾斜优化 | 自动拆分区(Join) |
| 加盐 | 热点 key 加随机前缀 |
| 两阶段聚合 | 局部聚合 + 全局聚合 |
| 广播小表 | 免大 Shuffle |
| 过滤极端 key | 单独处理后 union |
| 提高并行度 | 分散压力 |
5.3 加盐示例
scala
// 第一阶段:加盐聚合
rdd.map { case (key, v) =>
val salt = Random.nextInt(10)
((key, salt), v)
}.reduceByKey(_ + _)
// 第二阶段:去盐聚合
.map { case ((key, _), sum) => (key, sum) }
.reduceByKey(_ + _)六、Shuffle 优化
| 手段 | 说明 |
|---|---|
| 预聚合 | reduceByKey 而非 groupByKey |
| 压缩 | spark.shuffle.compress=true |
| Kryo | 更快更小 |
| 分区数合理 | 并行度匹配 |
| 分桶表 | BucketedJoin 免 Shuffle |
| 广播 join | 小表广播 |
七、动态资源
7.1 配置
| 参数 | 默认 | 说明 |
|---|---|---|
spark.dynamicAllocation.enabled | false | 开关 |
spark.dynamicAllocation.minExecutors | 0 | 最小 |
spark.dynamicAllocation.maxExecutors | - | 最大 |
spark.dynamicAllocation.executorIdleTimeout | 60s | 空闲回收 |
7.2 前提
YARN 下必须开启 External Shuffle Service
否则 Executor 回收导致 Shuffle 数据丢失| 场景 | 建议 |
|---|---|
| 共享集群 | 开启,省资源 |
| 独占 SLA | 固定资源稳定 |
八、SQL 场景专项
| 优化点 | 配置 |
|---|---|
| AQE | spark.sql.adaptive.enabled=true |
| 广播阈值 | spark.sql.autoBroadcastJoinThreshold |
| 动态分区裁剪 | 默认开启 |
| 分区合并 | AQE 自动 |
| 向量化读 | Parquet 默认开启 |
| 缓存 | 复用中间结果 |
九、调优总流程
调优顺序(先宏观后微观):
1. 资源:内存/核数/数量是否匹配数据规模
2. 并行度:Task 是否打满资源
3. Shuffle:是否可避免/减少
4. 序列化 + GC:开销是否过大
5. 倾斜:是否有慢 Task
6. 动态资源:负载波动是否利用
排查入口:
Spark UI → 慢 Stage → 慢 Task → 定位原因常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 资源打不满 | 分区数 < 核数,增大分区 |
| GC 严重 | 加内存、用 Kryo、堆外 |
| 个别 Task 慢 | 倾斜,加盐/两阶段/AQE |
| 广播变量 OOM | 缩小数据或提高阈值 |
| Shuffle 文件爆炸 | 分区过多,AQE 合并 |
| 动态分配丢数据 | 未开 External Shuffle Service |