spark sql的count(distinct)默认采用两阶段聚合优化,但当distinct字段存在大量重复值(如null、空字符串)时仍会倾斜;可通过过滤、group by替代、加盐等方式解决。

Spark SQL 中 COUNT(DISTINCT) 本身已默认做了两阶段聚合优化,一般不会直接引发严重倾斜;但当 distinct 字段存在大量重复值(尤其是 null、空字符串、默认值)时,仍会集中在少数 task 上处理——这是实际生产中最常踩的坑。
为什么 Spark 的 COUNT DISTINCT 不像 Hive 那样容易倾斜
Spark 在物理执行计划里对 COUNT(DISTINCT) 自动展开为 Expand + HashAggregate 两阶段:先按 distinct 字段分组打散(shuffle 到多个 partition),再全局合并计数。这和 Hive 默认用单个 reducer 处理完全不同。
但这个优化有个前提:distinct 字段的 key 分布必须能被 hash 均匀切分。一旦出现成千上万条记录共享同一个值(比如 user_id IS NULL 占 30%),这些 null 全被 hash 到同一个 partition,第二阶段就又卡住了。
验证方式很简单:在 Spark UI 的最后一个 HashAggregate stage 里看各 task 的 Shuffle Read Size —— 如果某 task 读了 5GB,其他都不到 10MB,说明还是倾斜了。
过滤 null 或默认值后再 COUNT DISTINCT
适用于业务允许忽略空值或已知倾斜值明确(如 ''、'unknown'、-1)的场景。这是最轻量、见效最快的手段。
- 直接在
WHERE子句中排除:SELECT COUNT(DISTINCT user_id) FROM log WHERE user_id IS NOT NULL AND user_id != '' - 如果必须保留 null 的语义(比如统计“有标识用户数”+“无标识用户数”),可拆成两部分:
SELECT COUNT(DISTINCT user_id) + (CASE WHEN COUNT(*) - COUNT(user_id) > 0 THEN 1 ELSE 0 END),避免把 null 当作一个 key 参与 shuffle - 注意:不能写成
WHERE user_id IS NOT NULL OR user_id != ''—— 这会导致逻辑短路失效,null 仍会进入后续计算
用 GROUP BY + COUNT 替代 COUNT DISTINCT(显式两阶段)
当自动优化不可靠(比如字段类型隐式转换导致 hash 不一致)、或需要复用中间结果时,手动控制更稳妥。
核心思路是把去重动作提前到 map 端完成一次局部去重,再 shuffle 聚合:
SELECT COUNT(*) FROM ( SELECT user_id FROM log GROUP BY user_id ) t
这个写法比 COUNT(DISTINCT) 多一次 shuffle,但好处是:
- 可以加
DISTRIBUTE BY控制分发逻辑(比如DISTRIBUTE BY hash(user_id)配合自定义 salt) - 方便在子查询里先过滤、加盐、或对倾斜 key 单独处理
- 执行计划清晰,便于定位哪个 stage 出问题
注意:GROUP BY user_id 本身也可能倾斜,所以它只是“可控的起点”,不是万能解药。
对倾斜 key 加随机前缀再聚合(Salted Aggregation)
当明确知道某些 key 极度集中(比如 user_id = 'guest' 有 200 万条),且无法过滤时,必须打散。
关键不是“加盐”,而是加盐后要能还原回原始 key:
SELECT COUNT(*) FROM (
SELECT
CASE WHEN user_id = 'guest' THEN concat('guest_', cast(rand() * 10 AS INT)) ELSE user_id END AS salted_id
FROM log
GROUP BY salted_id
) t
这样原本所有 'guest' 都被打散到 10 个不同 salted_id 下,各自聚合后,最终 count 就是真实去重数。但要注意:
- 随机数范围(如
rand() * 10)要足够覆盖倾斜程度,否则仍可能二次倾斜 不能在最外层再 - 如果倾斜 key 有多个(比如
'guest','test',null),需统一规则处理,否则逻辑不一致
GROUP BY 原始 key —— 那会把打散的结果又收回来,白忙活真正难的不是写对这几句 SQL,而是怎么低成本发现这些 key —— 得先 SELECT user_id, COUNT(*) FROM log GROUP BY user_id ORDER BY COUNT(*) DESC LIMIT 10,再决定是否加盐。










