根本原因不是函数本身慢,而是它触发了执行引擎中代价极高的数据膨胀或多次 shuffle;尤其在 sparksql、hive 中,count(distinct) 会因 expand 节点导致数十倍数据膨胀,或因热点 key 引发 shuffle 倾斜,且缺乏自动优化机制。

根本原因不是函数本身慢,而是它触发了执行引擎中代价极高的数据膨胀或多次 shuffle —— 尤其在 SparkSQL、Hive 等分布式引擎里,COUNT(DISTINCT) 往往会把单条记录复制成 N 条(N = distinct 表达式个数),中间数据量暴增数十倍。
SparkSQL 中 Expand 节点导致 30 倍数据膨胀
当 SQL 包含多个 COUNT(DISTINCT)(比如 30 个维度 UV 统计),Spark 会生成 Expand 节点,将每行原始输入“展开”为多行,每行只保留一个 distinct 字段的值,其余置为 NULL。这看似语义清晰,实际后果严重:
- 原始 15 亿行 × 30 个 distinct → 中间数据达 450 亿行
- 所有后续 shuffle、sort、aggregate 都基于膨胀后的数据,内存和磁盘压力陡增
-
Expand后必须接Aggregate,而该阶段无法复用 hash 结构,只能全量排序或构建巨型哈希表 - SparkUI 中可明显看到
Expand节点的outputRows是输入的整数倍,且下游 stage 持续失败重试
Hive/Spark 中单个 COUNT(DISTINCT) 的 shuffle 瓶颈
即使只有一个 COUNT(DISTINCT user_id),在超大数据量下仍可能慢得离谱,关键不在去重逻辑,而在 shuffle 分布不均:
- 默认使用
user_id做 hash 分区,若存在热点 ID(如测试账号、机器人用户),大量数据被路由到同一 reducer - Hive 3+ 虽支持
hive.optimize.countdistinct=true自动改写为两阶段聚合,但前提是统计信息准确且user_id基数预估合理;否则仍走单阶段,OOM 风险极高 - SparkSQL 不自动启用类似优化,需手动加 hint:
/*+ REPARTITION(1000) */强制打散,但 repartition 本身又引入额外 shuffle - 对比
GROUP BY user_id+COUNT(*),前者能天然利用 map-side combine 减少网络传输,后者不能
替代方案:COLLECT_SET + SIZE 或近似算法
对实时性要求不高的场景,绕开精确去重是最快落地的解法:
- 用
SIZE(COLLECT_SET(user_id))替代COUNT(DISTINCT user_id):避免 expand,但内存占用仍随 distinct 基数线性增长,适合基数 - Spark 3.0+ 支持
APPROX_COUNT_DISTINCT(user_id, 0.01):误差率 1%,底层用 HyperLogLog,内存恒定,速度提升 5–10 倍 - 提前物化:每天跑一次
INSERT OVERWRITE TABLE uv_daily SELECT dt, COUNT(DISTINCT user_id) ... GROUP BY dt,查询直接读物化表 - 绝对要避免:在子查询里嵌套多个
COUNT(DISTINCT),再 join —— 这会让膨胀叠加,复杂度从 O(N) 变成 O(N²)
真正卡住性能的,往往不是你写的那行 COUNT(DISTINCT),而是它背后隐式触发的执行计划变形。查慢查询,第一件事不是改 SQL,而是看 SparkUI 或 EXPLAIN 输出里有没有 Expand、有没有单个 reducer 处理 90% 数据的 shuffle 阶段 —— 这些信号比函数名更真实。











