Summing 与 Aggregating 引擎
上一章的 Replacing/Collapsing 解决"更新"。这一章的 Summing/Aggregating 解决另一个问题:明细数据太大,每次查询都从头聚合太慢。
思路是把聚合提前到写入时做,让 merge 顺手把同键的行"加起来"。
1. 为什么需要预聚合
events 表一天 1 亿行。运营看板要看「每天每个国家每种事件的 PV 和平均时长」。
-- 每次刷新看板都要扫 30 亿行
SELECT
toDate(event_time) AS d, country, event_type,
count() AS pv, avg(duration_ms) AS avg_ms
FROM events
WHERE event_time >= '2024-06-01'
GROUP BY d, country, event_type;但结果只有 30 天 × 200 国 × 4 类型 = 24000 行。输入 30 亿行,输出 2.4 万行——压缩比 12 万倍。这种情况就该预聚合。
明细表 events(30 亿行,180 GB)
│
│ 预聚合(写入时或 merge 时完成)
▼
汇总表 events_daily(2.4 万行,2 MB)
│
│ 查询直接读汇总表
▼
0.003 秒返回2. SummingMergeTree
最简单的预聚合引擎:merge 时,同一个 ORDER BY 键的行,数值列自动相加。
2.1 建表
CREATE TABLE events_daily
(
day Date,
country LowCardinality(String),
event_type LowCardinality(String),
pv UInt64,
total_ms UInt64,
uv_approx UInt64 -- 注意:这个字段有问题,后面讲
)
ENGINE = SummingMergeTree() -- 不指定列 = 所有非键数值列都求和
PARTITION BY toYYYYMM(day)
ORDER BY (day, country, event_type); -- 聚合的分组键也可以指定只对某几列求和:
ENGINE = SummingMergeTree((pv, total_ms)) -- 只有这两列相加,其他列取第一行的值2.2 写入与合并
-- 第一批数据
INSERT INTO events_daily VALUES
('2024-06-15', 'CN', 'view', 1000, 250000, 0),
('2024-06-15', 'CN', 'click', 200, 40000, 0);
-- 第二批(同样的键)
INSERT INTO events_daily VALUES
('2024-06-15', 'CN', 'view', 500, 130000, 0),
('2024-06-15', 'US', 'view', 800, 190000, 0);
SELECT * FROM events_daily ORDER BY day, country, event_type;merge 前:
┌────────day─┬─country─┬─event_type─┬───pv─┬─total_ms─┐
│ 2024-06-15 │ CN │ click │ 200 │ 40000 │
│ 2024-06-15 │ CN │ view │ 1000 │ 250000 │
│ 2024-06-15 │ CN │ view │ 500 │ 130000 │ ← 还没合并
│ 2024-06-15 │ US │ view │ 800 │ 190000 │
└────────────┴─────────┴────────────┴──────┴──────────┘OPTIMIZE TABLE events_daily FINAL;
SELECT * FROM events_daily ORDER BY day, country, event_type;┌────────day─┬─country─┬─event_type─┬───pv─┬─total_ms─┐
│ 2024-06-15 │ CN │ click │ 200 │ 40000 │
│ 2024-06-15 │ CN │ view │ 1500 │ 380000 │ ← 相加了
│ 2024-06-15 │ US │ view │ 800 │ 190000 │
└────────────┴─────────┴────────────┴──────┴──────────┘2.3 查询必须再 GROUP BY 一次
和 ReplacingMergeTree 一样,merge 是最终一致的。正确查法是在汇总表上再套一层聚合:
SELECT
day, country, event_type,
sum(pv) AS pv,
sum(total_ms) / sum(pv) AS avg_ms -- 加权平均
FROM events_daily
WHERE day >= '2024-06-01'
GROUP BY day, country, event_type
ORDER BY day, pv DESC;这个 GROUP BY 只在 2.4 万行上做,成本可以忽略,但保证了结果永远正确(无论 merge 有没有发生)。
- 存
pv(可加)和total_ms(可加),查询时算total_ms / pv - 不要直接存
avg_ms——两个平均值不能相加,merge 会把它们错误累加
同理:存 success_cnt 和 total_cnt,不存 success_rate。
2.4 SummingMergeTree 的两个坑
上面表里的 uv_approx 是个陷阱。UV 是去重计数,不可加:
2024-06-15 CN view uv=100 (part A,来自上午的数据)
2024-06-15 CN view uv=80 (part B,来自下午的数据)
merge 后 uv=180 ← 错!同一批用户上下午都访问了,真实 UV 可能只有 120UV 必须用下一节的 AggregatingMergeTree + uniqState。
CREATE TABLE t (k String, v UInt64, note String)
ENGINE = SummingMergeTree ORDER BY k;
INSERT INTO t VALUES ('a', 1, 'first'), ('a', 2, 'second');
OPTIMIZE TABLE t FINAL;
-- 结果:('a', 3, 'first') —— note 保留了排序后第一行的值,另一行丢了如果你需要保留某个字符串列的最新值,SummingMergeTree 做不到,得用 AggregatingMergeTree 加 argMaxState。
3. AggregatingMergeTree
SummingMergeTree 只会求和。AggregatingMergeTree 支持任意聚合函数,包括 uniq、quantile、argMax、max/min。
要理解它,得先理解 ClickHouse 的核心抽象:聚合状态。
3.1 聚合状态:-State 与 -Merge
普通聚合函数返回最终结果:
SELECT uniq(user_id) FROM events; -- 返回一个数字,比如 48213加 -State 后缀,返回中间状态(未完成的聚合数据结构):
SELECT uniqState(user_id) FROM events;
-- 返回二进制的 HyperLogLog 结构,无法直接阅读加 -Merge 后缀,把多个状态合并并得出最终结果:
SELECT uniqMerge(state) FROM (SELECT uniqState(user_id) AS state FROM events);
-- 48213聚合状态的生命周期
明细数据 中间状态 最终结果
┌──────────────┐
part A 的 user_id ──▶│ HLL 位图 A │─┐
└──────────────┘ │ uniqMerge
part B 的 user_id ──▶│ HLL 位图 B │─┼────────────▶ 48213
└──────────────┘ │
part C 的 user_id ──▶│ HLL 位图 C │─┘
└──────────────┘
关键:位图之间可以合并(求并集),所以可以分段计算再汇总
而最终结果(48213 这个数字)之间无法合并这就是 UV 能被预聚合的原理:存的不是数字,是可合并的数据结构。
3.2 AggregateFunction 类型
聚合状态在表里的存储类型是 AggregateFunction(函数名, 参数类型...):
CREATE TABLE events_agg
(
day Date,
country LowCardinality(String),
event_type LowCardinality(String),
pv SimpleAggregateFunction(sum, UInt64),
uv AggregateFunction(uniq, UInt64),
total_ms SimpleAggregateFunction(sum, UInt64),
max_ms SimpleAggregateFunction(max, UInt32),
p95_ms AggregateFunction(quantile(0.95), UInt32),
top_page AggregateFunction(argMax, String, DateTime)
)
ENGINE = AggregatingMergeTree()
PARTITION BY toYYYYMM(day)
ORDER BY (day, country, event_type);3.3 SimpleAggregateFunction vs AggregateFunction
| SimpleAggregateFunction | AggregateFunction | |
|---|---|---|
| 适用函数 | 状态就是结果本身的函数:sum min max any anyLast groupBitOr | 任意聚合函数 |
| 存储 | 就是原始类型(如 UInt64) | 二进制状态 blob |
| 写入 | 直接写值 | 必须写 xxxState(...) |
| 读取 | 直接读 | 必须用 xxxMerge(...) |
| 开销 | 无额外开销 | 有序列化开销 |
能用 SimpleAggregateFunction 就用它。sum、max、min 这些"状态即结果"的函数没必要走复杂路径。
3.4 写入 AggregatingMergeTree
必须用 -State 后缀的函数:
INSERT INTO events_agg
SELECT
toDate(event_time) AS day,
country,
event_type,
count() AS pv, -- Simple,直接写值
uniqState(user_id) AS uv, -- 需要 State
sum(duration_ms) AS total_ms, -- Simple
max(duration_ms) AS max_ms, -- Simple
quantileState(0.95)(duration_ms) AS p95_ms, -- 需要 State
argMaxState(page, event_time) AS top_page -- 需要 State
FROM events
WHERE event_time >= '2024-06-01' AND event_time < '2024-07-01'
GROUP BY day, country, event_type;3.5 查询 AggregatingMergeTree
必须用 -Merge 后缀:
SELECT
day,
country,
sum(pv) AS pv,
uniqMerge(uv) AS uv, -- 正确的跨行去重计数
sum(total_ms) / sum(pv) AS avg_ms,
max(max_ms) AS max_ms,
quantileMerge(0.95)(p95_ms) AS p95_ms,
argMaxMerge(top_page) AS latest_page
FROM events_agg
WHERE day >= '2024-06-01'
GROUP BY day, country
ORDER BY day, pv DESC
LIMIT 10;┌────────day─┬─country─┬─────pv─┬────uv─┬──avg_ms─┬─max_ms─┬─p95_ms─┬─latest_page─┐
│ 2024-06-01 │ CN │ 682341 │ 41203 │ 2521.44 │ 5049 │ 4788 │ /p/318 │
│ 2024-06-01 │ US │ 679812 │ 40988 │ 2519.87 │ 5049 │ 4791 │ /p/104 │
│ 2024-06-01 │ JP │ 681003 │ 41156 │ 2523.10 │ 5048 │ 4790 │ /p/442 │
└────────────┴─────────┴────────┴───────┴─────────┴────────┴────────┴─────────────┘注意 uniqMerge(uv) 的威力:它把 day, country, event_type 三个维度的 HLL 位图合并,得到 day, country 两个维度的正确 UV。这是 SummingMergeTree 永远做不到的。
SELECT uv FROM events_agg LIMIT 1;直接查 AggregateFunction 列会返回二进制状态的字符串表示,不是数字。在 AggregatingMergeTree 上永远不要写 SELECT *。
如果需要看状态的实际值,用 finalizeAggregation(uv):
SELECT day, country, finalizeAggregation(uv) AS uv_this_row FROM events_agg LIMIT 3;注意这只是单行状态的值,不是跨行合并的结果。
4. -State 的另一个用途:跨表传递
聚合状态可以序列化后在表之间传递,这让多级预聚合成为可能:
events(明细,1 亿行/天)
│ uniqState(user_id)
▼
events_hourly(小时级,24 × 200 × 4 = 19200 行/天)
│ uniqMergeState(uv) ← Merge + State 组合,合并后仍是状态
▼
events_daily(天级,800 行/天)
│ uniqMerge(uv)
▼
最终 UV 数字中间层用 -MergeState 组合子:
INSERT INTO events_daily
SELECT
toDate(hour) AS day,
country,
sum(pv) AS pv,
uniqMergeState(uv) AS uv -- 合并小时状态,输出天级状态
FROM events_hourly
WHERE hour >= today() - 1 AND hour < today()
GROUP BY day, country;关键性质:无论怎么分层合并,最终 uniqMerge 出来的 UV 都是正确的(HLL 的并集运算满足结合律)。这就是可加性抽象的价值。
5. Summing vs Aggregating 怎么选
| 场景 | 选择 |
|---|---|
| 只需要 sum / count | SummingMergeTree(更简单,无序列化开销) |
| 需要 UV、分位数、argMax | AggregatingMergeTree |
| 混合需求 | AggregatingMergeTree + SimpleAggregateFunction(sum, ...) |
实际上,AggregatingMergeTree 配合 SimpleAggregateFunction 可以完全替代 SummingMergeTree,且更灵活。新表建议直接用 AggregatingMergeTree。
6. 存储收益实测
SELECT
table,
sum(rows) AS rows,
formatReadableSize(sum(bytes_on_disk)) AS size
FROM system.parts
WHERE database = 'demo' AND active AND table IN ('events', 'events_agg')
GROUP BY table;┌─table──────┬──────rows─┬─size──────┐
│ events │ 100000000 │ 1.82 GiB │
│ events_agg │ 24000 │ 3.14 MiB │
└────────────┴───────────┴───────────┘查询对比:
-- 从明细算
SELECT country, count(), uniq(user_id) FROM events
WHERE event_time >= '2024-06-01' GROUP BY country;
-- Elapsed: 1.842 sec. Processed 100.00 million rows
-- 从预聚合算
SELECT country, sum(pv), uniqMerge(uv) FROM events_agg
WHERE day >= '2024-06-01' GROUP BY country;
-- Elapsed: 0.008 sec. Processed 24.00 thousand rows230 倍加速,存储只多 0.17%。
上面的 INSERT INTO events_agg SELECT ... FROM events 需要定时执行,而且要处理重复执行、迟到数据等问题。
更好的方案是让 ClickHouse 在数据写入 events 时自动同步写 events_agg——这就是下一章的物化视图。
7. 完整可运行示例
-- 1. 预聚合表
CREATE TABLE demo.events_agg
(
day Date,
country LowCardinality(String),
event_type LowCardinality(String),
pv SimpleAggregateFunction(sum, UInt64),
total_ms SimpleAggregateFunction(sum, UInt64),
max_ms SimpleAggregateFunction(max, UInt32),
uv AggregateFunction(uniq, UInt64),
p95_ms AggregateFunction(quantile(0.95), UInt32)
)
ENGINE = AggregatingMergeTree()
PARTITION BY toYYYYMM(day)
ORDER BY (day, country, event_type);
-- 2. 从明细回填
INSERT INTO demo.events_agg
SELECT
toDate(event_time), country, event_type,
count(), sum(duration_ms), max(duration_ms),
uniqState(user_id), quantileState(0.95)(duration_ms)
FROM demo.events
GROUP BY toDate(event_time), country, event_type;
-- 3. 按天查
SELECT day, sum(pv) AS pv, uniqMerge(uv) AS uv,
round(sum(total_ms) / sum(pv), 1) AS avg_ms,
round(quantileMerge(0.95)(p95_ms), 1) AS p95
FROM demo.events_agg
GROUP BY day ORDER BY day LIMIT 5;┌────────day─┬────pv─┬────uv─┬──avg_ms─┬────p95─┐
│ 2024-06-01 │ 33452 │ 22104 │ 2524.3 │ 4787.5 │
│ 2024-06-02 │ 33280 │ 22019 │ 2521.8 │ 4790.1 │
│ 2024-06-03 │ 33391 │ 22077 │ 2525.6 │ 4788.9 │
│ 2024-06-04 │ 33188 │ 21955 │ 2519.4 │ 4791.2 │
│ 2024-06-05 │ 33507 │ 22138 │ 2527.1 │ 4786.3 │
└────────────┴───────┴───────┴─────────┴────────┘- 建
events_dailySummingMergeTree 表,插入两批同键数据,观察 OPTIMIZE 前后的差异。故意存一列avg_ms,看看 merge 后它变成了什么荒谬的值。 - 建
events_aggAggregatingMergeTree,用uniqState存 UV。分别按(day, country, event_type)和只按(day)聚合,验证uniqMerge在降维时给出的 UV 小于各维度 UV 之和(因为去重了)。 - 对比
SELECT uv FROM events_agg LIMIT 1(乱码)和SELECT finalizeAggregation(uv) FROM events_agg LIMIT 1(数字),理解 AggregateFunction 的存储形态。 - 用明细表和预聚合表跑同一个业务查询,对比 Elapsed 和 Processed rows,算出加速比。
小结
- 当「输入行数 / 输出行数」比值很大时(超过 1000 倍),就该考虑预聚合
- SummingMergeTree 在 merge 时把同 ORDER BY 键的数值列相加,简单但只支持求和
- 预聚合表只存可加的量(count、sum),不存比率(avg、rate),查询时再算
- UV 不可加,SummingMergeTree 处理不了
- AggregatingMergeTree 通过
AggregateFunction类型存储聚合中间状态,支持任意聚合函数 -State生成状态,-Merge合并状态并出结果,-MergeState用于多级聚合SimpleAggregateFunction适用于状态即结果的函数(sum/max/min),无序列化开销,优先使用- 无论 merge 是否发生,查询时都要再
GROUP BY一次保证正确 - 下一章用物化视图让预聚合表自动维护 →