大基数group by易卡在shuffle read阶段,因key分布不均导致数据倾斜、中间膨胀及内存压力;应通过salting+pre-aggregate拆解热key,并避免count(distinct)与grouping sets共用。

超大数据集 + 大基数 GROUP BY 在 Spark SQL 中极易触发 shuffle 瓶颈、内存溢出和长尾 task,根本原因不是数据量大,而是 key 分布不可控、中间数据膨胀严重。直接加资源或调 spark.sql.shuffle.partitions 通常收效甚微,甚至让问题更隐蔽。
为什么大基数 GROUP BY 容易卡在 Shuffle Read 阶段
Spark 的 GROUP BY 必须把相同 key 的所有 record 拉到同一个 Executor 做聚合。当 key 基数高达千万甚至上亿(比如用户 ID、设备 ID、URL 哈希),就出现两个典型问题:
- Shuffle Write 数据量远超原始输入(例如输入 500GB,Write 达 2TB+),网络和磁盘 IO 成瓶颈
- 每个 key 对应的数据量极不均衡:99% 的 key 只有几条记录,但 Top 100 热 key 占据 30% 总数据量 → 数据倾斜
- Executor 内存压力陡增:大量小对象进入 HashAggregate,GC time 占比飙升,频繁 Full GC
Spark UI 上最明显的信号是:某个 Stage 的 Shuffle Read Size / Records 异常高,且 Task 耗时分布极度右偏(多数 20s,个别 40min+)。
用 salting + pre-aggregate 拆解热 key
对已知的热 key(如 top 100 用户、高频 URL),不能让它原样进 shuffle,必须提前“打散”。核心思路是:给热 key 加随机后缀,使其在 shuffle 阶段被分散到多个 partition,再二次聚合。
- 先用
approx_count_distinct()或采样统计识别热 key(避免全表扫描) - 对原始数据做两路处理:
– 非热 key 走正常GROUP BY
– 热 key 加rand(100)后缀,GROUP BY key, salt,再GROUP BY key汇总 - 最后
UNION ALL两路结果;注意列名、类型、顺序必须严格一致
示例伪代码:
WITH hot_keys AS (
SELECT user_id FROM logs
GROUP BY user_id
HAVING count(*) > 1000000
),
salted AS (
SELECT
CASE WHEN h.user_id IS NOT NULL THEN CONCAT(user_id, '_', CAST(rand(100) AS INT)) ELSE user_id END AS key,
COUNT(*) AS cnt
FROM logs l
LEFT JOIN hot_keys h ON l.user_id = h.user_id
GROUP BY key
)
SELECT
CASE WHEN key RLIKE '_[0-9]+' THEN SPLIT(key, '_')[0] ELSE key END AS user_id,
SUM(cnt) AS total_cnt
FROM salted
GROUP BY CASE WHEN key RLIKE '_[0-9]+' THEN SPLIT(key, '_')[0] ELSE key END
避免 COUNT(DISTINCT) 和 GROUPING SETS 同时出现
这是 EXPAND 算子的双重触发点:一旦查询里既有 COUNT(DISTINCT x) 又有 GROUPING SETS 或 CUBE,Spark 会强制生成 EXPAND 节点,导致行数指数级膨胀(100 亿输入 → 3000 亿中间行)。这不是配置能绕过的机制。
- 优先用近似去重:
approx_count_distinct(x, 0.01)(误差率 1%,性能提升 5–10 倍) - 若必须精确,把
COUNT(DISTINCT)拆成子查询 +collect_set+size(),但仅适用于中等基数( - 绝对不要在同一个 GROUP BY 中混用
GROUPING SETS和多个COUNT(DISTINCT);改用显式多路GROUP BY + UNION ALL
错误写法:GROUP BY a, b GROUPING SETS ((a), (b), ()) + COUNT(DISTINCT c) → 必炸
正确替代:SELECT a, NULL AS b, COUNT(DISTINCT c) FROM t GROUP BY a UNION ALL SELECT NULL AS a, b, COUNT(DISTINCT c) FROM t GROUP BY b
shuffle 分区数不是越大越好,要匹配 key 基数
spark.sql.shuffle.partitions 默认 200,对大基数场景完全不够——它决定的是 shuffle 后每个 reduce task 处理多少 key,而不是多少数据量。如果 key 基数是 5000 万,200 个分区意味着平均每个 task 要处理 25 万个 key,哈希表内存占用爆炸。
- 合理值 ≈ key 基数 ÷ 10 万(目标:每个 partition 平均承载 10 万左右 distinct key)
- 例如预估 key 基数 2000 万 → 设为 200;若 2 亿 → 设为 2000;超过 5000,建议配合 salting 使用
- 同时调大
spark.executor.memory和spark.memory.fraction,确保执行内存足够容纳哈希聚合结构
真正容易被忽略的是:key 基数必须是**预估**而非硬算。用 approx_count_distinct 或采样(TABLESAMPLE(0.1))快速探查,别一上来就全表 count distinct。











