<p>直接执行SELECT key, COUNT() AS cnt FROM tbl GROUP BY key ORDER BY cnt DESC LIMIT 10,观察首行cnt是否远超其余(如相差10倍以上),即可定位倾斜key;已知倾斜key后,用CASE WHEN对齐加盐打散(如CONCAT(key, '_', FLOOR(RAND() 100)))实现两阶段聚合。</p>

怎么查出是哪个key导致GROUP BY倾斜
直接在SQL里跑 SELECT key, COUNT(*) AS cnt FROM tbl GROUP BY key ORDER BY cnt DESC LIMIT 10,看前几行的cnt是不是远超其他值。比如其他key都是几百条,但有个device_id突然有250万条,基本就是它了。别依赖直觉,Spark UI里看task耗时差异大、某个task卡在99%、executor报java.lang.OutOfMemoryError,都只是现象,根源还得落到具体key上。
已知倾斜key,怎么加盐做两阶段GROUP BY
核心是把大key打散,先局部聚合再合并。假设已知倾斜key是'abc123',用CASE WHEN加随机前缀是最轻量的做法:
SELECT key, SUM(cnt) AS total_cnt
FROM (
SELECT
CASE WHEN key = 'abc123' THEN CONCAT(key, '_', FLOOR(RAND() * 100)) ELSE key END AS key,
COUNT(*) AS cnt
FROM tbl
GROUP BY
CASE WHEN key = 'abc123' THEN CONCAT(key, '_', FLOOR(RAND() * 100)) ELSE key END
) t
GROUP BY key要点:
- RAND()必须在GROUP BY和SELECT中保持一致,否则会漏统计
- 随机数范围(比如* 100)要足够大,确保原key被打散到至少5–10个子key,但也不能过大导致小key也被误拆
- 如果有多个倾斜key,CASE WHEN可以嵌套,或改用WHEN key IN ('k1', 'k2') THEN ...
不知道哪些key倾斜,但想自动应对
开启Spark SQL自适应执行(AQE)是最省事的兜底方式,尤其适合临时作业或无法预知数据分布的场景:
- 设置
spark.sql.adaptive.enabled=true - 必须同时开
spark.sql.adaptive.skewedJoin.enabled=true(注意:名字带Join,但它对GROUP BY也生效) -
spark.sql.adaptive.skewedPartitionMaxSplits=10可调,默认5,最大10;值太小打散不够,太大增加调度开销
这个方案不改SQL,只调参数,但要求Spark版本 ≥ 3.2,且作业不能禁用AQE(比如显式设了spark.sql.adaptive.enabled=false)。
为什么不能只靠提高spark.sql.shuffle.partitions
把spark.sql.shuffle.partitions从默认200调到1000,确实能让每个task处理的数据量变少,但治标不治本:
- 如果倾斜key本身占全量90%,那1000个task里,仍有约900个task空转,100个task抢着处理那90%数据——资源浪费+延迟没本质改善
- 分区数不是越高越好,超过集群core总数后,task频繁切换上下文,反而拖慢整体速度
- 对极端倾斜(比如单key百万级),单纯加分区可能仍OOM,因为单个task内存压力没真正释放
真正难啃的倾斜,从来不是靠“分得更细”,而是让大key“不再是一个key”。加盐、预过滤、AQE打散,才是动了数据结构本身。

















