0%

Spark data skew

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,是倾斜的高发区:distinctgroupByKeyreduceByKeyaggregateByKeyjoincogrouprepartition 等。

它们”为什么会倾斜”本质是同一件事——按 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 两条底层规律

  1. map 端 combine 只减”分区内重复”,不减”跨分区热点”reduceByKey / aggregateByKey 把每个 map 分区内的同 key 先合并成一份,Shuffle 数据量下降;但 1000 个 map 各发一份,reduce 端仍要处理这 1000 份——热点 key 的倾斜只是被”缓解”,不会被”消除”。
  2. 聚合是收敛的,join 是发散的:groupBy 类把 M 条收成 1 条(越聚越小);join 则把左表 M 条 × 右表 N 条配对(越配越大)。同一个热点 key,join 的倾斜烈度通常远高于聚合,所以 join 需要单独的治理手段(Broadcast Join / 热键隔离 / AQE skew join,见第 3 节)。

参考:Spark性能优化指南——高级篇

3. 治理手段(四类)

3.1 加盐打散(Salting)

给热点 key 拼接随机前缀/后缀(如 key + "_" + random(N)),把一个热点打散成 N 份落到不同 Task 并行处理。

  • 适用:可拆分的聚合(计数、求和)。
  • Tradeoff:多一次 Shuffle 阶段、逻辑更复杂;不适用于需要精确去重(distinct count)的场景——加盐后同一原始 key 被拆散,distinct 会算错。

3.2 两阶段聚合(局部聚合 + 全局聚合)

也称 Map 端预聚合 / 两遍聚合,专治 groupBy 类倾斜:

  1. 第一阶段(局部):给 key 加随机前缀,按 前缀+key 做局部聚合(reduceByKey / aggregateByKey),先把热点大幅缩小。
  2. 第二阶段(全局):去掉前缀,按原始 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 推理调优,是面试里最强的差异化叙事。