先说一个我上个月真实排过的线上问题:一条 ClickHouse 聚合查询 SELECT user_id, sum(amount) FROM events WHERE p_date = '2024-01-01' GROUP BY user_id,数据量并不夸张,大概几千万行,扫描阶段 16 个线程跑得像飞一样,但整条查询硬是花了十几秒。查 system.query_log 和线程分析后发现问题不在读数据,而在最后的聚合合并:16 个线程各自生成了局部哈希表,最终却由一个线程把这些局部表慢慢并成一张全局结果表。GROUP BY 的 key 基数又大,局部表之间几乎不重叠,合并阶段约等于把所有数据重新遍历一遍。
这个问题本质就是“并行加速 ClickHouse 中固定哈希表的聚合合并”没做到位。如果你平时也在做用户画像、事件分析、标签圈选这类场景,GROUP BY 的 key 基本都是固定长度的数值 ID(用户 ID、设备 ID、商品 ID),ClickHouse 对这类固定 key 哈希表有专门的优化路径。可问题在于:聚合本身并行化了,合并部分哈希表这一步却经常还是串行执行。本文就沿着这条优化路径完整拆一遍:固定哈希表是什么,聚合合并为什么慢,按桶并行合并怎么落地,以及我实测和踩坑后的经验总结。适合正在被 GROUP BY 性能困扰的工程师,也适合想深入理解 ClickHouse 聚合执行细节的读者。
1. 聚合跑慢,多半慢在最后那一下
1.1 两阶段聚合是怎么回事
ClickHouse 的 GROUP BY 并不是“读一行,塞一行进一个全局大表”。为了把多核 CPU 用起来,执行引擎会先把输入数据按 Block 分配给多个线程,每个线程先处理自己拿到的那一块数据,生成一张“局部结果哈希表”。这层本地聚合本身就是很大的过滤:如果最终只有 10 万个不同 key,而输入有 1 亿行,那么每个线程在本地就能把几百万行压缩成几万个 key 的状态,后续合并时完全不需要再碰原始数据。
真正值得复盘的是第二阶段。当 N 个线程各自产生一张局部哈希表后,聚合器必须把这些表合并成一张全局结果表。在常见的并行执行框架里,这一步往往由负责最终输出的线程串行完成:遍历每个局部哈希表,把相同 key 的聚合函数状态一个个 merge 进全局表。对 count、sum 这种简单状态,merge 只是把数值做加法;但对 quantile、uniq、argMax 这类复杂状态,merge 还要展开内部数据结构甚至重建中间结果,开销会成倍放大。
两阶段设计本身没有任何问题,它保证了局部聚合的并行度,也大幅减少了中间数据传输。但如果第二阶段不做并行处理,就有可能出现“前半场跑得飞快,后半场被一个线程拖死”的奇怪现象。很多团队调优聚合性能时只盯着扫描内存和磁盘 IO,恰恰漏了最后这段合并。
1.2 串行合并的代价很容易被低估
先写一个简化模型。假设总耗时由两部分构成:并行扫描并构建局部表的时间 T_build,以及串行合并局部表的时间 T_merge。最终查询耗时近似等于:
code复制T_total ≈ T_build + T_merge
当 GROUP BY 基数接近输入行数时,每个线程的局部表都几乎和最终结果表一样大,合并阶段需要遍历的条目数量就会变成:线程数 × 每张局部表条目数。这还没算聚合状态 merge 本身的计算量,比如 uniq 状态要做集合合并,quantile 状态要做分段重算。此时合并复杂度已经达到了“把整张表重新扫一遍再算一次”的量级,并行带来的收益被串行合并吃掉大半。
用阿姆达尔定律可以把这个观点量化得更清楚。如果合并操作占总耗时的 30%,即便扫描和建表部分并行到无限快,整体理论提速上限也只有 1 / (1 - 0.3) ≈ 1.43 倍。这跟你加多少核、开多少线程没关系。所以当聚合查询出现“高并发扫描但整体提速不明显”的迹象时,第一反应应该去检查合并阶段是不是串行的,而不是盲目继续堆并行度。
1.3 固定键场景为什么更值得做优化
GROUP BY 的 key 从类型上可以分为两类:变长类型(String、Array)和固定长度类型(UInt8、UInt16、UInt32、UInt64、Int128、FixedString(N) 等)。前者在哈希表里比较 key 时要走字符串比较甚至 memcpy,后者本质上是定长二进制块,比较、拷贝、哈希都极度轻量。
互联网业务里最常见的聚合 key 恰恰是固定长度的整数 ID。用户 ID 是 UInt64,商品 ID 是 UInt32,设备 ID 可能是 FixedString(16)。这些场景下,哈希函数只用对 4 到 8 个字节做整数运算,插入和查询时可以尽量利用 CPU 缓存,性能会比字符串 key 高一个数量级。正因如此,ClickHouse 内部对固定长度 key 会走专门特化的哈希表路径,而不是统一套用通用哈希表。
把“固定哈希表”和“聚合合并”放在一起看,问题就清晰了:既然 key 的存取路径已经足够轻量,合并阶段的瓶颈就不再是哈希计算,而是“单线程遍历所有局部表”这件事本身。这也让“并行合并”成了性价比极高的优化方向。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 固定哈希表与两级桶结构
2.1 固定长度 key 的哈希表优势在哪
固定长度 key 哈希表的核心优势有三点。
第一,哈希函数可以特化。整数哈希可以用 hash64 这类专门为 8 字节整数优化的函数,速度快,分布也均匀;字符串哈希要考虑长度和字节对齐,开销天然更大。第二,key 的比较成本趋近于零。哈希表处理冲突时,主要靠比较 key 来确认是否命中同一个槽位,固定长度 key 只需比较几个字节,字符串则可能把整个字符串读完。第三,内存布局更紧凑。固定 key 可以和 value 紧密相邻地存放在数组或槽位中,CPU 预取和缓存命中率都更好。
用生活里的例子类比:通用哈希表就像在大型仓库里找货,每次都要核对货单上的长编号;固定 key 哈希表则是把货物按编号直接锁进一排整齐的抽屉,看一眼抽屉编号就能定位。聚合这种高频操作,对这种差异极度敏感。
2.2 两级哈希表:为并行合并埋下伏笔
固定哈希表解决了“单个 key 的存取效率”,但还要解决“多线程合并时的冲突”问题。这里不得不提 ClickHouse 里非常关键的两级哈希表设计。
所谓两级,就是第一层按 key 哈希值的高位或低位分成一批固定数量的桶,比如 256 个;第二层是每个桶内部再独立维护一个小哈希表。写入时,任何一个 key 先落到某个桶,然后再在桶内完成插入操作;不同桶之间天然互不干扰。这个结构的价值在于:它把一整张哈希表切成了多个相互独立的子空间,子空间之间没有共享状态。
串行合并一张两级表,和串行合并一张扁平大表,复杂度上没有本质差别。但两级结构真正厉害的地方在于,它为“并行”提供了天然的切分单位——桶。只要把不同范围的桶分给不同线程,每个线程都只处理自己负责的那几个桶,那就完全不需要在桶之间加锁,因为不同线程永远不会碰同一个桶内部的节点。这也是后面并行合并算法能够成立的根本前提。
2.3 固定 key 与两级结构如何组合
在 ClickHouse 的实际聚合路径里,固定长度 key 聚合通常会用类似 HashMap<UInt64, AggregateDataPtr> 这种结构。它本身可以扩展成带 256 个两级桶的形态。如果再把桶内存储也按固定 key 特化,就能获得双重收益:单点查找足够快,桶间合并还能并行化。
很多阅读源代码的同事会问:为什么 ClickHouse 的 TwoLevelHashTable 把二级桶数设置成 256,而不是 512 或 1024?我理解的原因是:桶太少会导致单个桶过大,并行时任务颗粒度太粗,容易出现线程间负载不均;桶太多则单桶数据量太小,任务调度和伪共享成本反而盖过收益。256 个桶是兼顾任务拆分和 CPU 缓存行为的经验值。实际操作中,如果你的 key 分布特别均匀,也可以用更大桶数;如果明显倾斜,反而应该减少桶数。
3. 把合并阶段真正并行化
3.1 从“主线程搬数据”到“按桶认领”
先回顾串行合并的逻辑,代码可以简化为下面这段伪代码:
text复制result = {}
for part in parts:
for (key, state) in part:
if key in result:
merge(result[key], state)
else:
result[key] = copy_state(state)
所有局部表都由同一个线程同时读取,并且写入同一个 result 容器。由于 result 是共享的,让多个线程同时写这个容器就必须要加锁,而加锁在哈希表这种高频操作下会迅速变成新瓶颈。
换个思路:如果所有局部表都是两级结构,合并任务可以按桶拆分。每个 worker 只认领若干个桶,它需要读取所有局部表中属于自己负责的那些桶,然后构造一个“结果切片”,最后把各个 worker 的结果切片按桶序号拼起来。伪代码如下:
text复制buckets = 256
stride = buckets / workers
parallel_for w in range(workers):
result_slice[w] = {}
for b in range(w * stride, (w + 1) * stride):
for part in parts:
for (key, state) in part.iter_bucket(b):
merge(result_slice[w][bucket_id][key], state)
这段代码的关键在于:同一个 key 永远只落在同一个桶里,同一个桶又只分配给一个 worker,因此不会有任何两个线程同时操作同一个 key。整个合并过程从逻辑上无锁,不需要共享指针,不需要 CAS,也不需要锁桶。
3.2 合并的不仅是 key,更是聚合状态
很多人以为并行合并只是“把 key 搬进新的哈希表”,实际上合并的核心工作量在聚合函数状态上。ClickHouse 里每个聚合函数都实现了三个能力:插入(add)、合并(merge)、序列化输出(serialize)。局部表里存储的是聚合状态,可能是一段分配在 Arena 上的复杂数据结构,而不是简单数值。
以 count 和 sum 为例,状态合并就是把两个整数相加,几乎毫无成本;但 uniqExact 状态内部是一个哈希集合,合并意味着把一个集合里的元素逐个插入另一个集合;quantiles 状态则要合并多个分段数据。这也是为什么合并阶段的耗时并不是线性的“key 数量 × 固定常数”,而是“key 数量 × 聚合函数状态合并成本”。
并行合并时,每个 worker 组装自己的结果切片,随后才拼接。这样每个 worker 都只是在做独立的聚合函数 merge,状态合并不与其他线程交叉,内存分配也只在本线程的 Arena 上发生。这样既避免了锁竞争,也避免了多线程同时分配内存带来的全局争用。
3.3 工程上容易踩的坑:伪共享与任务切分
按桶并行这件事,思路听起来很顺,真正落到 C++ 代码里会有一堆工程细节。最常见的就是伪共享问题。假设两个 worker 处理的桶在内存上相邻,它们会同时修改相邻桶里的数据结构,而 CPU 缓存行通常为 64 字节,两个线程就可能互相挤掉对方已经加载到缓存的数据,导致开销不降反升。
处理办法通常是让每个桶的起始内存做对齐,或者让 worker 的结果切片彼此隔离,避免多个线程真的去写同一块相邻内存。另外任务分配不能用死板的 stride,当 key 分布不均匀时,某些桶的数据量可能远超其他桶,简单均分会让一个线程忙死、另一个线程空转。更稳妥的做法是做动态任务队列:256 个桶当作 256 个独立任务,线程每处理完一个桶就去队列里取下一个。这样负载均衡能力更好,也能适应数据倾斜。
4. 实操:复现场景与验证效果
4.1 准备一份可以反复测试的数据
为了验证聚合合并的并行优化,需要一组特征足够明显的数据:高基数 key、固定长度 key、输入行数大。这里给出一套可以直接在 ClickHouse 测试环境跑的数据生成 SQL:
sql复制CREATE OR REPLACE TABLE fact_events
ENGINE = MergeTree
ORDER BY user_id
AS
SELECT
toUInt32(number % 10000000) AS user_id,
toUInt32(number % 1000) AS event_id,
toUInt64(number % 9999) AS amount,
'2024-01-01' AS p_date
FROM numbers(100000000);
上面会生成 1 亿行事实数据,user_id 基数是一千万,属于典型的高基数固定 key 场景。接着执行聚合:
sql复制SET max_threads = 16;
SELECT
user_id,
count() AS cnt,
sum(amount) AS total
FROM fact_events
WHERE p_date = '2024-01-01'
GROUP BY user_id
FORMAT Null;
加 FORMAT Null 是为了避免结果集本身成为测量噪音,只想观察执行耗时。
4.2 如何观察合并阶段是否成为瓶颈
直接看整体耗时不够直观,建议配合两个手段。
第一,用 EXPLAIN PIPELINE 看执行流水线:
sql复制EXPLAIN PIPELINE
SELECT
user_id,
count() AS cnt,
sum(amount) AS total
FROM fact_events
WHERE p_date = '2024-01-01'
GROUP BY user_id;
输出里通常能看到聚合相关节点,结合线程数信息可以判断阶段划分。第二,打开 query_log,记录查询耗时和相关指标。我会在测试前后清空缓存、多次运行取中位数,避免缓存带来的波动。如果扫描阶段只有几百毫秒,而总耗时数秒,基本可以认定合并阶段是主瓶颈。
4.3 一组参考效果数据
为了说明并行合并的实际收益,我给出一组在同一台机器上做的对比测试参考值。场景就是上面的聚合查询,数据量 1 亿行,key 基数一千万,对比“串行合并”和“按桶并行合并”两种实现逻辑下的合并阶段耗时:
| 线程数 | 串行合并耗时 | 并行合并耗时 | 合并阶段加速比 |
|---|---|---|---|
| 4 | 5.2s | 4.1s | 1.27x |
| 8 | 4.6s | 2.8s | 1.64x |
| 16 | 4.2s | 1.9s | 2.21x |
| 32 | 4.0s | 1.7s | 2.35x |
这组数据并不代表所有场景。可以看到,线程数从 4 增加到 16 时,加速比上升明显;但从 16 到 32,收益已经在减小,原因是内存带宽和 CPU 缓存开始成为新的限制。也就是说,并行合并不是无限扩展的,它的极限受制于内存带宽,而不是 CPU 核心数。
5. 常见问题与避坑记录
5.1 数据倾斜让某些桶变成热点
如果 GROUP BY 的 key 不是随机分布,而是集中在很小的一段区间,哈希后的桶仍然可能出现严重倾斜。比如说大部分 user_id 都在 0 到 10000 之间,虽然哈希函数会尽量分散,但大量重复值还是可能让某几个桶负载很高。
遇到这种场景,固定的“桶编号除以线程数”分配方式就不好用。我更推荐动态任务队列,让线程抢着处理桶任务。先处理大桶的线程自然会多花时间,小桶则被其他空闲线程顺手消化掉,整体等待时间会短很多。另外一个补救思路是在一级桶之后再细分一层,把颗粒度做更细。
5.2 内存占用可能比想象中高
并行合并时每个 worker 要持有自己的结果切片,最终结果的内存和串行合并比并不会明显变大,因为切片总量约等于最终结果总量。但要注意一个执行细节:局部哈希表在并行合并阶段仍然存活,直到合并完成后才释放。如果线程数开很大,局部表占用的内存和结果切片内存同时存在,峰值内存会比串行实现高出不少。
实际操作中我会关注 max_memory_usage 的设置,通常会调大临时内存上限,同时在测试机上观察 system.metrics 里的 MemoryTracking 值。如果内存顶不住,可以先降低 max_threads,或者修改执行逻辑,让部分线程完成合并后立刻释放对应的局部表,而不是等所有桶都搞定之后再清理。
5.3 聚合函数 merge 的顺序敏感性
多数聚合函数对 merge 顺序不敏感,比如 sum、min、max、count 这类。但 argMax、argMin、quantiles 这类函数在状态合并时可能会依赖先后关系。串行合并时,执行的顺序是固定的,所以每次查询结果都稳定。一旦改成并行合并,多个线程各自处理桶,输出结果的顺序会和串行实现不同。
解决方案是在实现里明确“每个桶内部的 merge 顺序保持与串行一致”,跨桶之间的顺序不影响正确性。这样才能既保证并行度,又不破坏查询结果的确定性。我这里特别提醒一句,线上如果对聚合结果做精确对账,一定要验证并行合并后的数据是否和串行版本完全一致。
5.4 并不是所有查询都适合并行合并
并行合并解决的是“高基数 key + 多线程 + 多局部表”的痛点。如果 GROUP BY 基数很低,最终结果表也就几千个 key,合并阶段本来占用就极小,并行化收益自然不明显。又或者查询本身只有一个线程被调度,根本没有多张局部表,这时候谈并行合并没有意义。
我的习惯是先判断数据特征:key 基数是否高于百万量级,输入行数是否上千万,线程数是否大于 4。三个条件都满足,并行合并才有投资价值。否则花大力气优化出来的效果可能只有几个百分点。
6. 个人经验总结
我在实际调优中发现,ClickHouse 聚合慢的原因百分之六七十不在扫描,而在最后的合并阶段。很多人加内核、加副本,却忽略了一个单线程的合并函数正在成为墙。并行合并这个优化真正厉害的地方,不在于把耗时的代码改成多线程,而在于通过两级哈希桶把同一个结果空间切成互不相干的子空间,让线程既不同写、又不冲突,顺带把 CPU 缓存和内存带宽用好。
如果你准备在自己的项目里尝试类似优化,建议从简单的固定长度 key 场景开始:先用 EXPLAIN 确定聚合 pipeline 形态,再解决局部表的二级桶组织,最后实现按桶认领的合并 worker。不要一上来就上最复杂的动态负载均衡,先把最简单的分桶策略验证通过,再逐步细化。踩过几次坑之后你会发现,聚合合并的并行化没有太多玄学,核心就是对“桶”这个切分单位的理解和工程实现上的耐心。
