考察点
数据倾斜是 Spark 生产问题的头号选手,这题区分「背过方案清单」和「真处理过倾斜」的人。面试官想听:怎么从 UI/日志确认是倾斜而不是其他慢;不同算子(聚合、join)倾斜的治理手段差异;加盐方案的完整推导(为什么另一侧要膨胀 N 倍);以及什么情况下倾斜根本不该在 Spark 层面治。追问常往「倾斜和 OOM 的关系」「AQE 能替代手动处理吗」「热点 key 占比多少才算倾斜」方向走。
参考答案
先定位:怎么确认是倾斜
症状很典型:一个 stage 里 99% 的 task 几分钟跑完,剩下两三个 task 跑几十分钟甚至 OOM。确认靠 Spark UI 的 Stage tab,看 task 耗时和 Shuffle Read Size 的四分位数分布——max 是中位数的 5 倍以上基本就是倾斜。再进一步,开 spark.sql.adaptive 相关的 metrics 或在代码里对 key 做 count 排序,找出热点 key 是谁(往往是空值、默认值「-1」、或者大促当天的爆款商品 ID)。
倾斜的本质:key 分布极度不均,而 shuffle 按 key 分区,热点 key 的全量数据压到单个 partition、单个 task、单个 executor。CPU、内存、网络全卡在这一个点上,集群其他资源闲着。
手段一:加盐打散(两阶段聚合)
聚合类算子(reduceByKey、countByKey、groupBy 后 sum/count)的经典方案。分两步:
// 第一阶段:给 key 加随机前缀,把热点 key 拆成 N 份
val salted = rdd.map { case (key, v) => ((key, Random.nextInt(N)), v) }
val partial = salted.reduceByKey(_ + _) // 每个 (key, salt) 子集先聚合
// 第二阶段:去掉盐,对聚合后的中间结果再聚合一次
val result = partial.map { case ((key, _), sum) => (key, sum) }
.reduceByKey(_ + _)
为什么有效:第一阶段 shuffle 把热点 key 的数据均匀分到 N 个分区,单个 task 压力降为 1/N;第一阶段的 combine 已经把数据量大幅削减,第二阶段聚合的数据量很小。N 的选择:热点 key 的单分区数据量 / 期望单 task 处理量,常见取 10~100,别拍脑袋取 1000——盐太多会让所有 key 的聚合效率下降。
适用范围:聚合满足结合律的算子。collect_list 这种不行。
手段二:join 倾斜的三种治法
小表 join 大表,直接广播。小表不 shuffle,大表按 key 分布天然均匀,倾斜问题不存在。这是最简单的解法,也是为什么生产上小表阈值调得比较大。
大表 join 大表,热点 key 加盐 + 另一侧膨胀:
// 热点侧(事实表):key 加随机盐 0..N-1
val saltedFact = fact.map(r => ((r.key, Random.nextInt(N)), r))
// 维度侧:每行复制 N 份,分别配上所有可能的盐
val explodedDim = dim.flatMap(r => (0 until N).map(i => ((r.key, i), r)))
saltedFact.join(explodedDim).map { case ((k, _), (a, b)) => (k, (a, b)) }
推导逻辑:热点 key 的事实数据被随机分到 N 个桶,维度数据必须和每个桶都匹配得上,所以膨胀 N 倍——用维度侧 N 倍的存储换事实侧热点 key 1/N 的单点压力。只对确认的热点 key 做这个处理(用热点 key 列表过滤后分别 join 再 union),全表膨胀 N 倍代价太大。
过滤热点 key。如果热点是脏数据(空 key、爬虫流量、测试账号),最正确的做法是过滤或单独处理,而不是技术硬扛。先在业务上问一句「这些 key 真的需要 join 吗」。
手段三:让 AQE 兜底
Spark 3.x 开 AQE 后,SMJ 的倾斜分区自动拆分(spark.sql.adaptive.skewJoin.enabled 默认 true),聚合倾斜也有 spark.sql.adaptive.skewedPartition 相关优化。AQE 能处理中度倾斜,但极端倾斜(单 key 占数据 50%)或 outer join 场景它管不了,还是得手动。
治标还是治本
倾斜治理有个经常被忽略的判断:倾斜到底是数据特性还是上游问题。
- 数据特性(天然的长尾分布、爆款商品):按上面的手段处理。
- 上游问题(join key 生成逻辑有 bug、默认值填充错误、join 条件写错导致笛卡尔积):治本在修 bug。我见过「倾斜」查到最后是 join 条件里有个
OR把不该关联的都关联上了,数据量膨胀三个数量级。
还有一个结构性解法:热点 key 如果可预测(大促爆款、头部主播),预处理阶段把热点拆出去单独算,非热点正常跑,最后 union。把热点从分布式问题降级为单机问题——单 key 再大,一台机器内存也装得下的时候,单机处理比分布式调度快得多。
可能的追问
- 加盐打散会不会影响结果正确性? 聚合满足结合律/交换律时两阶段结果和单阶段一致;但 collect 类(要保留全组数据)不能加盐。join 加盐 + 膨胀的结果等价性依赖于「每份事实数据恰好匹配一份膨胀后的维度数据」,数学上是笛卡尔分解。
- 倾斜和 OOM 是什么关系? 倾斜是 OOM 最常见的诱因之一:热点 task 内存需求远超平均,execution/storage 都不够时就 OOM。但也有纯倾斜不 OOM 的情况——数据能 spill,只是慢。反过来 OOM 不一定是倾斜,user memory 泄漏也会 OOM。
- 怎么预防倾斜而不是事后治? ETL 侧:key 设计时避免单值热点(用复合 key、分桶键);监控侧:对 key 分布做例行统计,热点 key 清单维护成元数据;架构侧:极端热点单独走预聚合或外部缓存(Redis 计数),不进 Spark。