大表 GROUP BY 易倾斜主因是分组字段存在高频值(如"0"、空值),导致 shuffle 时大量数据集中到同一 task;需通过加盐、两阶段聚合或过滤热点 key 主动分散负载,而非仅调高并行度。

为什么大表 GROUP BY 容易倾斜
GROUP BY 本身不倾斜,但当某个分组字段(比如 user_id)存在极少数高频值(如默认值 "0"、空字符串、测试账号),而这些值对应的数据量远超其他 key(例如百万级 vs 十几条),shuffle 阶段就会把所有相同 key 的数据拉到同一个 task 处理——这就是倾斜的根源。Spark SQL 不会自动识别“热门 key”,它只按 hash 值分发,所以 GROUP BY user_id 在遇到 user_id = "0" 占 40% 数据时,必然卡在单个 task 上。
用两阶段聚合(局部 + 全局)绕过单点瓶颈
核心思路是:先在 map 端做一次轻量聚合,打散热点;再 shuffle 后做最终合并。Spark SQL 本身不直接支持“加盐”,但可通过子查询模拟:
- 第一阶段:对原表加随机前缀(比如
concat(user_id, '_', cast(rand() * 10 as int))),然后GROUP BY这个新字段,得到中间结果 - 第二阶段:去掉前缀,按原始
user_id再次GROUP BY并聚合中间结果的统计值(如sum(cnt))
注意:rand() 要控制盐值范围(如 * 10 生成 0–9),避免分区爆炸;盐值太大反而增加 shuffle 数据量。
提前过滤或重映射明确的倾斜 key
如果已知几个导致倾斜的 key(比如 user_id IN ("0", "", "test")),不要等它们进主流程。直接拆开处理:
- 用
WHERE user_id NOT IN ("0", "", "test")跑主逻辑 - 单独对这几个 key 做聚合(甚至可转成广播变量查维表),最后用
UNION ALL合并结果 - 比“加盐”更干净,但前提是能枚举出倾斜 key,且业务允许忽略或特殊处理它们
若倾斜 key 是 null,别用 WHERE user_id IS NOT NULL 简单过滤——先确认 null 是否业务有效值;否则应统一替换为一个带业务含义的占位符(如 "unknown"),再参与分组。
调整 shuffle 并行度只是辅助手段
spark.sql.shuffle.partitions 默认 200,对倾斜无实质缓解。它只是把“一个超重 task”变成“多个较重 task”,但热点 key 仍全落在某几个 partition 里。只有配合加盐或过滤才有意义:
- 加盐后,需同步调高该参数(如设为 400–800),确保盐值足够分散
- 单纯调高参数反而可能让小 key 的 task 更多,增加调度开销
- 线上建议先观察 stage 中各 task 输入数据量(Spark UI 的 “Input Size / Records” 列),再决定是否调整
真正关键的不是并行度数字,而是让倾斜 key 的数据流主动“变细”或“绕开” shuffle —— 过滤、加盐、两阶段聚合,三者选其一,比盲目调参管用得多。

















