Spark 数据倾斜的成因与治理
一句话定位:数据倾斜是分布式计算里最典型、最高频的性能杀手——少量 key 承载了不成比例的数据量,拖出长尾 Task,让整个 Stage 的吞吐受限于最慢的那一个。
面试里它既是”会不会调优”的试金石,也是把大数据经验迁移到 GPU 推理调优的最佳桥梁(长尾 Task ≈ 静态 Batching 木桶效应)。
1. 成因:为什么会发生倾斜
分布式计算按 key 分区(默认 HashPartitioner:同一 key 必落同一 Task)。当数据本身的 key 分布不均(存在热点 key / 大 key),或分区键 / 哈希函数选得不合适时,少数 Task 会分到远超均值的数据量,成为长尾(straggler),拖慢整个 Stage——因为 Stage 的完成取决于最慢的 Task。
关键结论:倾斜只发生在 Shuffle 过程中。窄依赖(map / filter / mapPartitions)不跨节点重分布、不按 key 汇聚,因此不存在按 key 汇聚的倾斜;只有 Shuffle(按 key 重分区)才会把热点 key 集中到一个 reduce Task 上。
2. 触发 Shuffle 的算子(何时会踩坑)
以下算子会产生 Shuffle,是倾斜的高发区:distinct、groupByKey、reduceByKey、aggregateByKey、join、cogroup、repartition 等。
它们”为什么会倾斜”本质是同一件事——按 key 哈希把数据打进同一个 reduce Task;区别只在于每个算子对 key 的处理方式不同,导致倾斜的烈度不同。
2.1 算子功能与倾斜原理对照
| 算子 | 功能 | 倾斜原因 |
|---|---|---|
groupByKey |
按 key 分组,把同 key 的 value 收集成 Iterable,不聚合 |
无 map 端预聚合,热点 key 的全部 value 原样 shuffle 到同一 Task,Shuffle 量最大、最易 OOM |
reduceByKey |
按 key 用可结合函数聚合(sum / count / max) | 自带 map 端 combine,Shuffle 量小;但热点 key 跨分区汇聚到同一 reduce,长尾仍在 |
aggregateByKey |
reduceByKey 的通用版:分区内 seqOp + 跨分区 combOp |
同 reduceByKey,有 map 端预聚合,热点 key 仍会汇聚 |
distinct |
去重(等价 map(_ -> null).reduceByKey(...)) |
按值分组去重,热点值(如空串 / 默认值)汇聚到同一 Task;不可加盐,否则破坏去重语义 |
join |
两表按 key 关联配对 | 热点 key 两侧记录汇聚到同一 Task,且配对结果是 M×N 乘积放大,比聚合倾斜更危险 |
cogroup |
多个 RDD 按 key 联合分组,各取 value 的 Iterable |
多个 groupByKey 叠加、无预聚合,多数据源的热点 key 叠加 |
repartition |
重新分区以调整并行度 / 分布 | RDD 的 repartition(n) 均匀打散、不按业务 key,本身不引入倾斜;DataFrame 的 repartition(n, col) / repartitionByRange(col) 按列分布,遇热点列会重新引入倾斜 |
2.2 两条底层规律
- map 端 combine 只减”分区内重复”,不减”跨分区热点”:
reduceByKey/aggregateByKey把每个 map 分区内的同 key 先合并成一份,Shuffle 数据量下降;但 1000 个 map 各发一份,reduce 端仍要处理这 1000 份——热点 key 的倾斜只是被”缓解”,不会被”消除”。 - 聚合是收敛的,join 是发散的:groupBy 类把 M 条收成 1 条(越聚越小);join 则把左表 M 条 × 右表 N 条配对(越配越大)。同一个热点 key,join 的倾斜烈度通常远高于聚合,所以 join 需要单独的治理手段(Broadcast Join / 热键隔离 / AQE skew join,见第 3 节)。
3. 治理手段(四类)
3.1 加盐打散(Salting)
给热点 key 拼接随机前缀/后缀(如 key + "_" + random(N)),把一个热点打散成 N 份落到不同 Task 并行处理。
- 适用:可拆分的聚合(计数、求和)。
- Tradeoff:多一次 Shuffle 阶段、逻辑更复杂;不适用于需要精确去重(distinct count)的场景——加盐后同一原始 key 被拆散,distinct 会算错。
3.2 两阶段聚合(局部聚合 + 全局聚合)
也称 Map 端预聚合 / 两遍聚合,专治 groupBy 类倾斜:
- 第一阶段(局部):给 key 加随机前缀,按
前缀+key做局部聚合(reduceByKey/aggregateByKey),先把热点大幅缩小。 - 第二阶段(全局):去掉前缀,按原始 key 做全局聚合。
- 适用:可累加的聚合(sum / count / 可两阶段 avg)。
- 不适用:不可累加(distinct、median、全局 topN)。
- 关键配合:优先用
reduceByKey/aggregateByKey(自带 map 端 combine,shuffle 数据量小),避免groupByKey(无 map 端聚合,shuffle 量最大,最易倾斜)。
3.3 Broadcast Join 消除 Shuffle
小表 join 大表时,把小表广播到所有 Executor,在 map 端完成 join,完全避免大表的 Shuffle——倾斜自然消失。
- 适用条件:小表能放进 Executor 内存(默认阈值
spark.sql.autoBroadcastJoinThreshold = 10MB,可上调)。 - Tradeoff:小表过大广播会撑爆内存/网络;大表 join 大表不适用(考虑 Bucket Join / Sort Merge Join)。
- 若倾斜源于 join 的倾斜 key,且维表(小表)不大,Broadcast Join 是直接”绕开”倾斜的利器。
3.4 倾斜 key 单独处理(隔离热键)
工程上最常用、最可控:先识别热点 key,再把它与非倾斜 key 分开处理。
- 非倾斜 key 走正常 join / 聚合;
- 倾斜 key 单独用 Broadcast Join,或拆成多个子 key 分别 join 后再
union。 - 例:大表 join 维表但该大表某 key 极热 → 把热 key 从大表过滤出来单独和维表广播 join,其余走正常路径。
3.5 兜底:AQE 自动倾斜处理
Spark 3.0+ 的 AQE(spark.sql.adaptive.enabled=true)默认开启 spark.sql.adaptive.skewJoin.enabled:
运行时统计各分区大小,若某分区 > 中位数 N 倍且超过阈值,自动拆分为多个 Task 并行。
- 主要解决 join 倾斜;对非 join 的 groupBy 倾斜帮助有限(AQE 的 coalesce partitions 仅缓解小分区过多)。
- 面试要点:AQE 是”白嫖”的兜底,但要知道原理、知道它的边界,不能只依赖它。
4. 面试落点:提前抽样发现,而非事后救火
事后救火(作业跑慢/失败、OOM、长尾 Task 几小时跑不完才去猜)是被动且昂贵的。要做到提前发现:
- 抽样定位:全量跑之前,对 join key / group key 做
sample()+countByValue()(或groupBy(key).count()),看 key 分布——若某 key 的 count 远高于均值即为倾斜 key。 - Spark UI 看板:Stage 的 Task 耗时分布、Shuffle Read 大小分布;单个 Task 的 Shuffle Read 特别大、duration 特别长 → 倾斜。
- 建基线:对作业建立 key 分布基线,新数据接入时校验是否出现新热点。
数字直觉:若某 key 占全量 30%、并行度 100,则该 key 对应的单个 reduce 独吞约 30% 数据,其余 99 个 Task 分摊 70%——这就是典型长尾。
可迁移话术(回指 1.1.1):Spark 长尾 Task ≈ 静态 Batching 木桶效应;提前抽样发现倾斜 ≈ 推理前先做输入分布 Profiling / 建立基线,而非 GPU 利用率已呈锯齿才救火。把这套”提前量化基线 → 定位 → 验证”的方法论直接迁移到 GPU 推理调优,是面试里最强的差异化叙事。