主题
第 6 章 MergeTree 家族进阶
学习目标:理解 MergeTree 的 5 个变种为什么存在、它们各自在 Merge 时额外做什么"加料动作"、对应什么真实业务;能画出"每个引擎 Merge 前 / Merge 后"的数据变化图;知道
FINAL的代价与替代方案;面试官问「怎么实现亿级日活的实时 PV/UV」能立刻抛出AggregatingMergeTree + MaterializedView。
6.0 一张总图先立起来
MergeTree
(基础:分区 + 排序 + 稀疏索引 + 合并)
│
┌────────────┬────────────┬───┴─────────┬──────────────────┬──────────────┐
▼ ▼ ▼ ▼ ▼ ▼
Replacing Summing Aggregating Collapsing VersionedCollapsing Graphite
保留主键 主键求和 聚合状态合并 +1/-1 行折叠 带版本号的折叠 时序 rollup
最新版本 (数值列) (配合 MV 用) (适合 CDC) (乱序写入 OK) (Graphite 专用)
👆 共同点:所有家族成员在"Merge 过程中"多做一步事
普通 MergeTree : 直接归并
Replacing : 归并 + 同主键保留最新
Summing : 归并 + 同主键数值列累加
Aggregating : 归并 + 同主键聚合状态合并
Collapsing : 归并 + 同主键按 sign 折叠
Graphite : 归并 + 按时间粗化 rollup核心心智:它们都是 MergeTree —— 存储格式一样、Part 结构一样、稀疏索引一样、分区/排序键一样;唯一区别是 Merge 时对"同主键的多行"做什么处理。
📌 "同主键" = ORDER BY 的完整列相等(不是你写的 PRIMARY KEY,注意!)。
6.1 ReplacingMergeTree —— 主键去重(最终一致)
6.1.1 它在 Merge 时做什么
普通 MergeTree:同主键多行 → 全部保留。 ReplacingMergeTree:同主键多行 → 合并到一起时只保留一条(根据 ver 列决定保留谁)。
INSERT 顺序:
(user_id=1, name='Alice', ver=1)
(user_id=1, name='Alicia', ver=2)
(user_id=1, name='A.Smith', ver=3)
写入立刻:磁盘上真的有 3 行(3 个 Part,每个 1 行)
后台 Merge 发生后:只剩 (user_id=1, name='A.Smith', ver=3)
(ver 最大的那一行胜出)生活类比:快递柜每个格子只放最新一件。你往同一个格子塞了 3 次快递,柜子系统会保留最后一次放的那件,旧的会被"合并扫一扫"扔掉。
6.1.2 语法
sql
CREATE TABLE learn_ck.user_profile
(
user_id UInt64,
name String,
email String,
city String,
updated_at DateTime
)
ENGINE = ReplacingMergeTree(updated_at) -- ver 列
PARTITION BY tuple()
ORDER BY user_id; -- 主键就是 user_idReplacingMergeTree(ver):ver参数可选,可以不写;不写的话同主键保留最后插入的一条(取决于内部顺序,不可靠)。- 强烈建议写
ver:保留ver最大的那条;并列时保留最后一条。 - 可选
is_deleted列:ClickHouse 23.x 起支持ReplacingMergeTree(ver, is_deleted),is_deleted=1的行合并时会被删除 —— 实现"软删除"。
6.1.3 最终一致 ≠ 实时去重
重要! 去重发生在 Merge 时,而 Merge 是异步后台任务:
sql
INSERT INTO user_profile VALUES (1, 'A', 'a@x', 'BJ', '2024-01-01');
INSERT INTO user_profile VALUES (1, 'B', 'b@x', 'SH', '2024-01-02');
SELECT * FROM user_profile WHERE user_id = 1;
-- 可能看到 2 行!!因为 Merge 还没跑解决办法:
SELECT ... FINAL:查询时实时去重,但代价很大(读时归并)。OPTIMIZE TABLE ... PARTITION ... FINAL:强制合并一次(运维操作,不能当查询手段)。argMax()聚合:最推荐 —— 直接在 SELECT 里用聚合函数实现"每组取最新"。
sql
-- ① FINAL:可靠但慢
SELECT * FROM user_profile FINAL WHERE user_id = 1;
-- ② argMax:自行实现最终一致,通常比 FINAL 快
SELECT user_id,
argMax(name, updated_at) AS name,
argMax(email, updated_at) AS email,
argMax(city, updated_at) AS city,
max(updated_at) AS updated_at
FROM user_profile
WHERE user_id = 1
GROUP BY user_id;生产铁律:业务实时查询用
argMax,离线批处理用FINAL。千万不要在高频查询里用FINAL。
6.1.4 FINAL 的真实代价
普通 SELECT:
所有 Part 并行读 → 聚合计算 → 返回
SELECT ... FINAL:
所有 Part 读完,按 ORDER BY 再做一次归并排序,
对每组同主键的多行做去重(取 ver 最大),
然后再做业务聚合。
→ 读 IO 一样,但额外加 CPU 排序 + 去重 → 通常慢 2-10 倍6.1.5 典型场景:用户最新画像表
业务:每次用户修改资料,写入一条新行。查询"用户当前的昵称"。
sql
-- 建表
CREATE TABLE learn_ck.user_latest
(
user_id UInt64,
nickname String,
city String,
level UInt8,
version UInt64 DEFAULT toUnixTimestamp(now())
)
ENGINE = ReplacingMergeTree(version)
ORDER BY user_id;
-- 模拟 3 次修改
INSERT INTO user_latest(user_id, nickname, city, level, version) VALUES
(1001, 'Alice', 'BJ', 1, 100),
(1001, 'Alicia', 'SH', 2, 200),
(1001, 'Lily', 'SZ', 3, 300);
-- 查询(推荐用 argMax)
SELECT user_id,
argMax(nickname, version) AS nickname,
argMax(city, version) AS city,
argMax(level, version) AS level
FROM user_latest
WHERE user_id = 1001
GROUP BY user_id;
-- ┌─ user_id ─┬─ nickname ─┬─ city ─┬─ level ─┐
-- │ 1001 │ Lily │ SZ │ 3 │
-- └───────────┴────────────┴────────┴─────────┘6.1.6 为什么不能当"精确去重"
- 如果你刚写入,Merge 还没发生 → 不去重;
- 不同分区的同主键数据永远不合并 → 永远不去重;
FINAL只保证查询那一刻的去重,不是持久化去重。
结论:ReplacingMergeTree 是"最终一致"的去重,适合"读多写多、能容忍短暂重复"的场景(最新画像、维度表、缓慢变化维)。不适合"主键必须唯一否则会算错钱"的金融场景。
6.2 SummingMergeTree —— 数值列自动求和
6.2.1 它在 Merge 时做什么
同主键多行 → 在 Merge 时把所有数值类型的列(非主键)累加,只保留一条结果。
INSERT 顺序:
(date=2024-01-01, event='click', user=1001, count=3)
(date=2024-01-01, event='click', user=1001, count=5)
(date=2024-01-01, event='click', user=1001, count=2)
Merge 后:
(date=2024-01-01, event='click', user=1001, count=10) ← 累加了生活类比:记账本同主题自动合并。你给"2024-01-01 吃饭"这个条目写了 3 次(8元、12元、5元),合并后只保留一条"2024-01-01 吃饭 25元"。
6.2.2 语法
sql
CREATE TABLE learn_ck.daily_stats
(
event_date Date,
event_name LowCardinality(String),
user_id UInt64,
pv UInt32, -- 自动累加
uv_est UInt32, -- 自动累加
revenue Decimal(18, 2) -- 自动累加
)
ENGINE = SummingMergeTree((pv, uv_est, revenue)) -- 可指定要求和的列(可选)
PARTITION BY toYYYYMM(event_date)
ORDER BY (event_date, event_name, user_id);- 参数可选,不写则"所有非主键的数值列"都被 Sum。
- 非数值列(如 String):Merge 时保留同组的第一条(其实是未定义行为,建议同组 String 值必须一致)。
6.2.3 不能代替 GROUP BY
sql
SELECT event_date, event_name, sum(pv), sum(revenue)
FROM daily_stats
WHERE event_date = '2024-01-01'
GROUP BY event_date, event_name;还是要写 GROUP BY + sum()! 因为:
- Merge 是异步的,查询瞬间可能还有未合并的 Part;
- 不同分区 / 不同 Part 之间的同主键数据不会自动合并。
SummingMergeTree 的作用是"减少 GROUP BY 要扫的行数",而不是"替代 GROUP BY"。
6.2.4 典型场景:每日指标预聚合
sql
-- 从 events_mt 每秒把原始事件聚合到 daily_stats
-- (用物化视图串起来,见第 10 章;这里先手动 INSERT)
INSERT INTO daily_stats
SELECT event_date, event_name, user_id, count(), uniq(user_id), sum(revenue)
FROM events_mt
WHERE event_date = today() - 1
GROUP BY event_date, event_name, user_id;- 原始表每天 10 亿行;
- 预聚合表每天可能只有 1000 万行(按 user × event_name 聚合);
- 查"每日 PV 总和"时扫 1000 万行 vs 10 亿行,差 100 倍。
6.2.5 坑
- 主键设计错位:如果 ORDER BY 里没有"你要的 GROUP BY 列",Sum 不会发生。
- 非数值列变"随机值":有 String 列在排序键之外时,Merge 保留的是"不确定的第一条"。
- 溢出:
UInt32的 sum 可能溢出,大流量场景用UInt64更安全。
6.3 AggregatingMergeTree —— 聚合状态的"积木"
这是家族里最强大也最绕的一个。 一旦理解,就能用它 + 物化视图做"亿级实时大屏"。
6.3.1 先理解一个概念:聚合状态 AggregateFunction(T)
普通聚合函数的输入是"行",输出是"一个值":
sum(revenue) 行→标量
uniq(user_id) 行→标量
avg(price) 行→标量但 ClickHouse 有一个中间状态的概念:你可以拿到 "sum 已经算了多少、还没收尾" 的状态本身。
sumState(revenue) 行→状态(二进制 blob,里面有 running_sum 累加器)
uniqState(user_id) 行→状态(HyperLogLog 二进制)
avgState(price) 行→状态(里面是 {sum, count})两个状态可以合并成一个状态:
sumMerge(state_a, state_b) = sumState(a+b)
uniqMerge(state_a, state_b) = HyperLogLog 合并
avgMerge(state_a, state_b) = {sum_a+sum_b, count_a+count_b}状态可以收尾为最终结果:
sumMerge(state) → 一个数
uniqMerge(state) → 一个估算基数
avgMerge(state) → sum/count 的最终 avg📌 这就是 "分布式聚合" 的基础:一大堆状态在各节点算出来 → 传回中心 →
merge合并 → 得到最终结果。
6.3.2 AggregatingMergeTree 做什么
让聚合状态在 Merge 时自动合并:
INSERT 顺序:
(date=2024-01-01, uv_state=HLL(users={1,2,3}))
(date=2024-01-01, uv_state=HLL(users={2,3,4}))
Merge 后:
(date=2024-01-01, uv_state=HLL(users={1,2,3,4})) ← 合并了 HyperLogLog6.3.3 完整示例:实时 PV/UV
Step 1:建表
sql
-- 明细事件表
CREATE TABLE learn_ck.pv_events
(
event_date Date,
ts DateTime,
user_id UInt64,
url String
) ENGINE = MergeTree
ORDER BY (event_date, ts);
-- 聚合表:一行 = 一天一个页面的 PV/UV 状态
CREATE TABLE learn_ck.pv_agg
(
event_date Date,
url String,
pv_state AggregateFunction(sum, UInt64),
uv_state AggregateFunction(uniq, UInt64)
)
ENGINE = AggregatingMergeTree
PARTITION BY toYYYYMM(event_date)
ORDER BY (event_date, url);
-- 物化视图:明细表 INSERT 时自动写到聚合表
CREATE MATERIALIZED VIEW learn_ck.pv_mv
TO learn_ck.pv_agg
AS
SELECT event_date,
url,
sumState(toUInt64(1)) AS pv_state,
uniqState(user_id) AS uv_state
FROM learn_ck.pv_events
GROUP BY event_date, url;Step 2:写入明细
sql
INSERT INTO pv_events VALUES
('2024-01-01','2024-01-01 10:00:00',1001,'/home'),
('2024-01-01','2024-01-01 10:01:00',1001,'/home'),
('2024-01-01','2024-01-01 10:02:00',1002,'/home'),
('2024-01-01','2024-01-01 10:03:00',1003,'/product/1');MV 触发,把 4 行的状态写到 pv_agg 里(存的是二进制 HLL 等状态)。
Step 3:查询最终 PV/UV
sql
SELECT event_date,
url,
sumMerge(pv_state) AS pv,
uniqMerge(uv_state) AS uv
FROM pv_agg
WHERE event_date = '2024-01-01'
GROUP BY event_date, url
ORDER BY pv DESC;
-- ┌─ event_date ─┬─ url ─┬─ pv ─┬─ uv ─┐
-- │ 2024-01-01 │ /home │ 3 │ 2 │
-- │ 2024-01-01 │ /product/1 │ 1 │ 1 │
-- └──────────────┴──────────────┴──────┴──────┘6.3.4 状态函数家族(常用)
| 聚合 | State 函数 | Merge 函数 | 说明 |
|---|---|---|---|
sum(x) | sumState(x) | sumMerge(s) | 精确求和 |
count() | countState() | countMerge(s) | 精确计数 |
avg(x) | avgState(x) | avgMerge(s) | 精确平均 |
uniq(x) | uniqState(x) | uniqMerge(s) | 基数估算(HLL) |
uniqExact(x) | uniqExactState(x) | uniqExactMerge(s) | 精确去重基数(内存代价大) |
uniqCombined(x) | uniqCombinedState(x) | uniqCombinedMerge(s) | 精度/内存折中 |
quantile(x) | quantileState(x) | quantileMerge(s) | 分位数估算 |
argMax(x,y) | argMaxState(x,y) | argMaxMerge(s) | max(y) 对应的 x |
topK(K)(x) | topKState(K)(x) | topKMerge(K)(s) | Top-K 估算 |
6.3.5 典型杀手级场景:实时大屏
Kafka → events_raw (MergeTree)
│
│ MaterializedView (每 INSERT 自动触发)
▼
agg_by_minute (AggregatingMergeTree)
agg_by_hour (AggregatingMergeTree)
agg_by_day (AggregatingMergeTree)
│
▼
实时大屏查 agg_by_*_Merge
→ 亿级明细被压到千级 → 毫秒响应6.4 CollapsingMergeTree / VersionedCollapsingMergeTree —— 行折叠
6.4.1 "折叠" 是什么
业务上有些事件是成对撤销的:订单状态变化、CDC 同步、增减粉丝。 思路:每次"状态变更"写两条 —— 一条 -1(表示"之前那条不算了")+ 一条 +1(表示"新的状态")。Merge 时一对 +1/-1 相互抵消(折叠)。
INSERT 顺序:
(order_id=10, status='pending', Sign=+1) ← 1. 订单创建
(order_id=10, status='pending', Sign=-1) ← 2. 撤销旧状态
(order_id=10, status='confirmed', Sign=+1) ← 3. 新状态
Merge 后:
行 1 和行 2 正负抵消 → 消失
只留:(order_id=10, status='confirmed', Sign=+1)6.4.2 语法
sql
CREATE TABLE learn_ck.order_state
(
order_id UInt64,
user_id UInt64,
status LowCardinality(String),
amount Decimal(18, 2),
Sign Int8 -- +1 / -1
)
ENGINE = CollapsingMergeTree(Sign)
PARTITION BY tuple()
ORDER BY (order_id, user_id);- 同主键相邻的 +1/-1 对会在 Merge 时被消掉。
- 查询时用
sum(Sign)作为行数:
sql
-- 查 confirmed 状态的订单金额
SELECT order_id, sum(amount * Sign) AS amount
FROM order_state
WHERE status = 'confirmed'
GROUP BY order_id
HAVING sum(Sign) > 0; -- 过滤"已折叠掉"的行6.4.3 坑:同主键 +1/-1 顺序必须对
sql
-- 正确顺序(-1 先于 +1):
INSERT INTO order_state VALUES (10, 1, 'pending', 99, +1);
INSERT INTO order_state VALUES (10, 1, 'pending', 99, -1); -- 撤销
INSERT INTO order_state VALUES (10, 1, 'confirmed', 99, +1);
-- 错误顺序(同组 +1 多于 -1):
INSERT INTO order_state VALUES (10, 1, 'pending', 99, +1);
INSERT INTO order_state VALUES (10, 1, 'pending', 99, +1); -- 不是折叠,是多算一次CollapsingMergeTree 要求客户端自己保证顺序。乱序写入 → 折叠错位。
6.4.4 VersionedCollapsingMergeTree —— 乱序写入也能折对
加一个 Version 列,Merge 时按 (主键, Version) 排序,再按 Sign 折叠:
sql
CREATE TABLE learn_ck.order_state_v
(
order_id UInt64,
user_id UInt64,
status String,
amount Decimal(18, 2),
Sign Int8,
Version UInt32
)
ENGINE = VersionedCollapsingMergeTree(Sign, Version)
ORDER BY (order_id, user_id);- 适合多客户端并发写入、乱序到达的 CDC 场景(Debezium → Kafka → CH)。
- 缺点:Version 需要业务层生成(通常是 DB 的 binlog 位点)。
6.4.5 典型场景:MySQL → ClickHouse 增量同步
MySQL binlog:
UPDATE order SET status='paid' WHERE id=10
↓ Debezium 解析
↓ before=(10,'pending',99) after=(10,'paid',99)
↓
发到 Kafka 两条消息
↓
消费进 CH:
INSERT (10,'pending',99, -1, 123) ← 旧 before + Sign=-1
INSERT (10,'paid', 99, +1, 124) ← 新 after + Sign=+1
↓
Merge 后只剩 (10,'paid',99,+1,124) ✅6.5 GraphiteMergeTree —— 时序数据专用 rollup
最特化的一个引擎,专门为 Graphite(一个老牌时序系统)设计,在普通业务里用得不多。
它在 Merge 时做的事:按配置的"老化策略"把高频点位粗化成低频点位:
5 分钟内的点 → 保留 10 秒粒度
1 天内的点 → 合并成 1 分钟粒度
1 周前的点 → 合并成 1 小时粒度
1 年前的点 → 合并成 1 天粒度sql
CREATE TABLE learn_ck.metrics_graphite
(
Path String,
Value Float64,
Time DateTime,
Date Date,
Timestamp UInt32
)
ENGINE = GraphiteMergeTree('graphite_rollup') -- 引用 config 里的 rollup 配置
PARTITION BY toYYYYMM(Date)
ORDER BY (Path, Time);graphite_rollup 由服务端 config.xml 指定。业务上基本只有做监控自建的人会用。
6.6 FINAL 的代价 + 替代方案总结
| 场景 | 推荐 | 备选 | ❌ 反推荐 |
|---|---|---|---|
| ReplacingMT 读最新版本 | argMax(col, ver) GROUP BY pk | OPTIMIZE PARTITION FINAL(离线) | SELECT * FINAL(高频) |
| SummingMT 读聚合值 | sum(col) GROUP BY pk | FINAL 仅调试用 | SELECT col FINAL |
| AggregatingMT 读最终值 | xxMerge(state) GROUP BY pk | FINAL | xxMerge(xxState(col)) 等奇怪写法 |
| CollapsingMT 读有效行 | sum(col*Sign) GROUP BY pk HAVING sum(Sign)>0 | FINAL | WHERE Sign=1 会漏折叠 |
| 精确去重 | 迁 Replicated + 业务主键去重 + 用 uniqExact | FINAL DEDUPLICATE(离线) | 高频 SELECT FINAL |
再啰嗦一次:
SELECT ... FINAL= 查询时强制按主键归并 + 执行变种规则 → 比普通 SELECT 慢 2-10 倍;OPTIMIZE TABLE ... FINAL= 强制后台合并成 1 个 Part → 巨耗 I/O,仅离线 / 冷分区用。
6.7 真实案例汇总
6.7.1 场景 A:用户最新画像(Replacing)
sql
CREATE TABLE learn_ck.user_profile_latest
(
user_id UInt64,
nickname String,
city String,
level UInt8,
version UInt64 DEFAULT toUnixTimestamp64Milli(now64(3))
)
ENGINE = ReplacingMergeTree(version)
ORDER BY user_id;
-- 业务查询:每次用 argMax
SELECT user_id,
argMax(nickname, version) AS nickname,
argMax(level, version) AS level
FROM user_profile_latest
WHERE user_id IN (1,2,3)
GROUP BY user_id;6.7.2 场景 B:实时大屏 PV/UV(Aggregating + MV)
sql
-- 1) 事件明细表
CREATE TABLE learn_ck.events_raw
( event_date Date, ts DateTime, user_id UInt64, url String
) ENGINE = MergeTree ORDER BY (event_date, ts);
-- 2) 聚合表
CREATE TABLE learn_ck.events_agg
( event_date Date, url String,
pv_state AggregateFunction(sum, UInt64),
uv_state AggregateFunction(uniq, UInt64)
) ENGINE = AggregatingMergeTree ORDER BY (event_date, url);
-- 3) 物化视图串起来
CREATE MATERIALIZED VIEW learn_ck.events_mv TO learn_ck.events_agg AS
SELECT event_date, url,
sumState(toUInt64(1)) AS pv_state,
uniqState(user_id) AS uv_state
FROM learn_ck.events_raw
GROUP BY event_date, url;
-- 4) 实时查大屏
SELECT event_date, url,
sumMerge(pv_state) AS pv,
uniqMerge(uv_state) AS uv
FROM events_agg
WHERE event_date = today()
GROUP BY event_date, url
ORDER BY pv DESC LIMIT 10;6.7.3 场景 C:订单状态 CDC(VersionedCollapsing)
sql
CREATE TABLE learn_ck.order_cdc
( order_id UInt64, status LowCardinality(String),
amount Decimal(18,2),
Sign Int8, Version UInt64
)
ENGINE = VersionedCollapsingMergeTree(Sign, Version)
ORDER BY order_id;
-- 查询
SELECT order_id, argMax(status, Version) AS status, sum(amount * Sign) AS amount
FROM order_cdc
GROUP BY order_id
HAVING sum(Sign) > 0;6.8 📌 与 MySQL / PG 的对比小框
| 场景 | MySQL | PostgreSQL | ClickHouse |
|---|---|---|---|
| "主键存在则更新" | INSERT ... ON DUPLICATE KEY UPDATE | INSERT ... ON CONFLICT ... DO UPDATE | ReplacingMergeTree + argMax(最终一致) |
| 实时累加计数 | UPDATE t SET c=c+1 | UPDATE t SET c=c+1 | SummingMergeTree + INSERT |
| 实时 UV 估算 | 业务层配 Redis HLL | pg_hll 扩展 | AggregatingMergeTree + uniqState |
| CDC 同步 | binlog + 目标端 UPDATE | 逻辑复制 pub/sub | CollapsingMergeTree / VersionedCollapsingMergeTree |
| 事务 | ACID 完整 | ACID 完整 | 无跨行事务 —— 用 Sign 折叠实现逻辑事务 |
关键反差:
- MySQL/PG 的 UPSERT 是同步的,插入瞬间就生效;
- CH 的"UPSERT"都是异步合并语义,查询时要配合
FINAL或argMax才能看到最终结果。
为什么 CH 不做同步 UPSERT? 因为它是列存 + 面向亿级扫描的 OLAP 引擎 —— 同步 UPSERT 必须先"原地改某行",这会破坏列存压缩块,IO 代价巨大。异步 Merge 用"批量归并"摊销 IO,更符合 OLAP 本质。
6.9 本章小结
┌───────────────────────────────────────────────────────────┐
│ 家族成员共性:都是 MergeTree,只是 Merge 时多做一步事 │
├───────────────────────────────────────────────────────────┤
│ · Replacing → 同主键保留 ver 最大的一行 │
│ · Summing → 同主键数值列累加 │
│ · Aggregating → 同主键聚合状态合并 │
│ · Collapsing → +1/-1 折叠(需客户端保证顺序) │
│ · VersionedColl→ +1/-1 + Version 折叠(允许乱序写) │
│ · Graphite → 时序点位粗化 │
├───────────────────────────────────────────────────────────┤
│ 通用查询套路(避开 FINAL): │
│ Replacing → argMax(col, ver) GROUP BY pk │
│ Summing → sum(col) GROUP BY pk │
│ Aggregating→ xxMerge(state) GROUP BY pk │
│ Collapsing → sum(col*Sign) GROUP BY pk HAVING sum(Sign)>0│
├───────────────────────────────────────────────────────────┤
│ 生产铁律: │
│ ① FINAL 代价极大,高频查询禁用 │
│ ② OPTIMIZE FINAL 仅用于离线压缩 / 冷分区 │
│ ③ 所有家族成员都该叠上 Replicated 前缀 │
└───────────────────────────────────────────────────────────┘6.10 面试高频题
Q1:ReplacingMergeTree 能做精确去重吗?为什么?
考察点:最终一致 vs 精确去重的本质区别。
标准答案:
- 不能做精确去重。ReplacingMergeTree 的去重是最终一致的:
- 去重发生在 Merge 时,而 Merge 是异步后台任务;
- 同主键但在不同分区的数据永远不会合并;
- 刚 INSERT 的数据,查询瞬间可能还没 Merge → 看到重复行。
- 强行去重的手段:
SELECT ... FINAL:查询时临时归并,但慢 2-10 倍;OPTIMIZE PARTITION ... FINAL:强制合并(运维手段);argMax(col, ver) GROUP BY pk:最推荐,SQL 层表达清晰且通常比 FINAL 快。
- 真正要"精确唯一主键",请用 MySQL/PG 这类 OLTP,或者在业务写入前去重后再写入。
加分项:能补"23.x 起支持 ReplacingMergeTree(ver, is_deleted) 实现软删除";能讲清 "ver 列并列时取最后一条" 是实现定义而非 SQL 标准。
易错点:以为 ORDER BY user_id 就会自动去重,忽视了 Merge 的异步性。
Q2:SELECT ... FINAL 和 OPTIMIZE TABLE ... FINAL 的区别?代价分别是什么?
考察点:两个最容易混的 FINAL。
标准答案:
SELECT ... FINAL:查询时按主键临时归并 + 执行家族规则(去重/求和/聚合);不修改磁盘。- 代价:比普通 SELECT 慢 2-10 倍;读的数据量一样,但多了归并排序 + 规则计算;
- 频率:每次 SELECT 都重算,高频查询不能用;
OPTIMIZE TABLE ... FINAL:后台立刻执行一次全表 Merge,把所有 Part 合并成一个大 Part;修改磁盘。- 代价:I/O 巨爆炸,亿级表可能跑几小时;占满 Merge 线程池,阻塞其他 DDL;
- 频率:只在需要时执行一次;不能当成查询手段。
加分项:能讲替代方案(argMax / xxMerge / sum(x*Sign)),能讲 OPTIMIZE ... PARTITION ... FINAL 可以只压某个冷分区来降低代价。
易错点:把它俩当成"开关" —— 其实 OPTIMIZE 是运维动作、SELECT FINAL 是查询动作。
Q3:AggregatingMergeTree + 物化视图是怎么实现实时大屏的?
考察点:CH 最强大的"预聚合 + 状态 merge"模式。
标准答案(按数据流走一遍):
- 明细表
events_raw用普通MergeTree接收原始事件。 - 聚合表
events_agg用AggregatingMergeTree,列类型是AggregateFunction(sum, UInt64)/AggregateFunction(uniq, UInt64),存聚合的中间状态(二进制)。 - 物化视图
events_mv把events_raw的 INSERT 转成聚合状态写入events_agg:sqlSELECT event_date, url, sumState(toUInt64(1)) AS pv_state, uniqState(user_id) AS uv_state FROM events_raw GROUP BY event_date, url - 后台 Merge 会把同
(event_date, url)的状态合并(sumMerge / uniqMerge HyperLogLog 合并)。 - 大屏查询时:
SELECT sumMerge(pv_state), uniqMerge(uv_state) FROM events_agg GROUP BY ...→ 毫秒响应。
本质:把"亿级明细扫 + 聚合"压缩为"千万级聚合状态扫 + 状态合并"。
加分项:
- 能讲
uniqState存的是 HyperLogLog(定长状态),哪怕 user_id 有上亿种也不会爆; - 能讲"链式 MV":一个 MV 的结果作为另一个 MV 的输入;
- 能讲
POPULATE的坑(第 10 章)。
易错点:把 AggregateFunction(sum, UInt64) 当成"就是个 UInt64" —— 它是二进制状态,必须用 sumState / sumMerge 访问。
Q4:CollapsingMergeTree 和 VersionedCollapsingMergeTree 的区别?各自适合什么场景?
考察点:CDC 同步场景的引擎选择。
标准答案:
CollapsingMergeTree(Sign):- 必须客户端保证同主键 +1/-1 的写入顺序(-1 要先于后续的 +1);
- 适合"单客户端、顺序严格"的场景;
- 乱序写入会导致折叠失败(出现"净 Sign > 1"的脏行)。
VersionedCollapsingMergeTree(Sign, Version):- 多了一个
Version列,Merge 时按(主键, Version)排序后再按 Sign 折叠; - 允许乱序写入:只要最终 Version 大的在后,就能折对;
- 适合"多消费者并发写"、Kafka → Debezium → CH 的 CDC 场景;
- 代价:Version 需要业务层正确生成(通常是 binlog 位点、Kafka offset、MySQL 的
(binlog_file, pos)哈希等)。
- 多了一个
加分项:能讲"查询时都要用 sum(col * Sign) GROUP BY pk HAVING sum(Sign) > 0"的通用套路;能举 MySQL → Debezium → Kafka → CH 的完整链路。
易错点:在并发写入场景用 CollapsingMergeTree,导致数据错位。
Q5:能否用 SummingMergeTree 完全替代 GROUP BY sum()?
考察点:对 Summing 引擎"减少扫描但不替代聚合"的理解。
标准答案:
不能。 原因:
- Merge 是异步的,查询瞬间可能还有未合并的 Part,SummingMergeTree 表内仍然可能有"同主键多行"。
- 不同分区的数据永远不合并,跨分区查询必然要 GROUP BY。
- Summing 的作用是"减少 GROUP BY 要扫的行数"(从明细亿级压到聚合表千万级),而不是"免去 GROUP BY"。
正确用法:在 SummingMergeTree 上依然要写 GROUP BY pk,只是扫的是聚合后的行,非常快。
加分项:能补"非数值列 + 非主键列"的合并行为是不确定的(保留同组任意一条,通常是第一条);建议要么全部进主键、要么进数值列。
易错点:看到 SummingMergeTree 就写 SELECT col FROM t 不加 GROUP BY,以为结果已经是累加完的。
Q6:什么是"聚合状态" AggregateFunction?它解决了什么问题?
考察点:聚合的分布式/增量合并能力。
标准答案:
- 普通聚合函数(
sum,uniq,quantile…)是"一次性算出最终结果";但在分布式 / 增量场景里,我们希望能把"半成品状态"保留下来,再和新状态合并。 - ClickHouse 为几乎所有聚合函数都提供了 3 个变种:
xxState(col):把行聚合成中间状态(二进制结构:累加器 / HLL / 样本…);xxMerge(state)/xxMerge(state1, state2, …):把多个中间状态合并,得到最终标量;xxMergeState(state):把多个中间状态合并成新的中间状态(用于链式 MV)。
- 用途:
AggregatingMergeTree在 Merge 时自动合并同主键的聚合状态;- 分布式查询时,各分片返回状态,协调节点做最终 Merge;
- 物化视图可以存状态,业务查询时再按需 Merge。
- 存储的是定长二进制状态,比存原始明细省几个数量级。
加分项:能讲 uniq 是 HyperLogLog、uniqCombined 内部用了 HLL+LL、quantile 是 TDigest(不同算法会影响精度);能讲 argMaxState 这种特殊聚合状态。
易错点:以为 AggregateFunction(sum, UInt64) 就是普通 UInt64 字段 —— 其实它是 blob,不能直接当数用。
Q7:为什么 ClickHouse 没有真正的 UPSERT?它用什么实现这个语义?
考察点:OLAP vs OLTP 的本质差异。
标准答案:
- 真正的
UPSERT需要"定位到某一行,原地改"。在列存表里,一行的不同列分散在不同的.bin文件里;原地改意味着"解压列块 → 找到行 → 改 → 重压缩 → 写回"——单行改动的 I/O 代价巨大,完全违背 OLAP 的"批量扫描" 本性。 - ClickHouse 用 "写入+异步合并" 实现逻辑 UPSERT:
ReplacingMergeTree实现"主键存在则替换",取 ver 最大一条;SummingMergeTree实现"主键存在则累加数值";CollapsingMergeTree(+Version)实现"先撤销旧状态再写新状态";- 都是"写多行 + 合并时收敛"而非"原地改"。
- 副作用:强一致场景要么用
argMax/sum(x*Sign)在查询时计算、要么接受"最终一致"的语义。 - 还有
ALTER TABLE ... UPDATE/DELETE的 Mutation 机制,但本质是"重写整个 Part",不是原地修改;适合批量后台操作,不能高频触发(第 12 章详讲)。
加分项:能补"MySQL/PG 是行存 + B+Tree / Heap,天然适合单行改;OLAP 必须放弃这条才能拿到列存 + 向量化带来的百倍扫描速度"。
易错点:以为"CH 只是没实现 UPSERT",其实是故意不做。
Q8:生产环境怎么用 MergeTree 家族做一套"实时数仓"?给个架构。
考察点:综合应用题,体现对整个家族的理解。
标准答案(参考):
Kafka
│
▼ Engine=Kafka
events_kafka ──(MV)──▶ events_raw (MergeTree,ODS 明细层)
│
│ MV1
▼
events_cdc (VersionedCollapsingMergeTree,
接 CDC 同步,带 Version)
│
│ MV2
▼
events_agg_daily (AggregatingMergeTree,
按 day × url × event 预聚合状态)
│
│ MV3 (链式)
▼
events_agg_hourly (AggregatingMergeTree,
按小时状态二次聚合)
│
▼
📺 实时大屏 / BI (查 events_agg_*)
维度/最新画像:user_profile_latest (ReplacingMergeTree)
配合 字典表(Dictionary) 做大表 JOIN
所有表统一 Replicated* 前缀 + Keeper 协调要点:
- ODS 层用
MergeTree存全量明细; - CDC 同步用
VersionedCollapsingMergeTree; - 预聚合层
AggregatingMergeTree+ 物化视图; - 用户维度用
ReplacingMergeTree+argMax取最新; - 生产环境全部换成
Replicated*前缀; - 用分区(日/月)+ TTL 做冷热分层。
加分项:能讲"MV3 链式"时要用 xxMergeState 而非 xxState 避免重复聚合;能讲 Projection 与 MV 的取舍。
易错点:把所有表都用 MergeTree 堆原始明细,查询时跑 sum() / uniq() 全表扫 —— 没有体现家族引擎的真正价值。
📌 下一章预告:第 7 章讲"写入与读取最佳实践" —— 批量写入、
async_insert、max_threads、FORMAT全家桶,解决"怎么把数据高速灌进 MergeTree 而不炸 Part"。
🎬 可视化演示
演示加载缓慢或样式异常?点此在新标签页打开 ↗
💻 示例代码
python
"""
第 6 章 MergeTree 家族进阶 · 5 种引擎行为演示
依次展示:
1. ReplacingMergeTree 同主键去重 + argMax vs FINAL 对比
2. SummingMergeTree 数值列自动累加(合并前 vs 合并后)
3. AggregatingMergeTree 聚合状态合并 + uniqMerge 实时 PV/UV
4. CollapsingMergeTree +1/-1 折叠
5. VersionedCollapsing 乱序 + Version 折叠
前置:
先跑 init.sql 建表,再跑 seed.py 灌数
运行:
python family_play.py
"""
from __future__ import annotations
import sys
import time
import textwrap
import clickhouse_connect
HOST = "127.0.0.1"
PORT = 8123
USER = "default"
PASSWORD = ""
DB = "learn_ck"
def banner(s: str) -> None:
print("\n" + "=" * 78)
print(" " + s)
print("=" * 78)
def show(client, sql: str, label: str = "") -> None:
sql = textwrap.dedent(sql).strip()
if label:
print(f"\n-- {label} --")
print(">>> " + sql.replace("\n", "\n "))
try:
res = client.query(sql)
rows, cols = res.result_rows, res.column_names
if not rows:
print(" (0 rows)"); return
widths = [max(len(str(c)), max((len(str(r[i])) for r in rows), default=0))
for i, c in enumerate(cols)]
line = "+".join("-"*(w+2) for w in widths)
print(" " + line)
print(" |" + "|".join(f" {c:<{w}} " for c, w in zip(cols, widths)) + "|")
print(" " + line)
for r in rows[:20]:
print(" |" + "|".join(f" {str(v):<{w}} " for v, w in zip(r, widths)) + "|")
if len(rows) > 20:
print(f" ... ({len(rows)} rows total, showing first 20)")
print(" " + line)
except Exception as e:
print(f" !! error: {e}")
def demo_replacing(client):
banner("① ReplacingMergeTree —— 用户最新画像")
show(client, f"SELECT count() FROM {DB}.replacing_user", "Merge 前的物理行数")
show(client,
f"""SELECT user_id, count() AS dup_rows
FROM {DB}.replacing_user GROUP BY user_id
ORDER BY dup_rows DESC LIMIT 5""",
"每个 user_id 的重复行数(说明没去重)")
show(client,
f"""SELECT user_id,
argMax(nickname, version) AS nickname,
argMax(city, version) AS city,
max(version) AS latest_ver
FROM {DB}.replacing_user WHERE user_id <= 3
GROUP BY user_id ORDER BY user_id""",
"【推荐】用 argMax 取最新 —— 不阻塞、不归并")
show(client,
f"""SELECT user_id, nickname, city, version
FROM {DB}.replacing_user FINAL
WHERE user_id <= 3 ORDER BY user_id""",
"【慎用】SELECT … FINAL 查询时即时去重(慢 2-10 倍)")
print("\n→ 手动触发合并:OPTIMIZE TABLE learn_ck.replacing_user FINAL")
try:
client.command(f"OPTIMIZE TABLE {DB}.replacing_user FINAL")
print(" OK")
except Exception as e:
print(f" !! {e}")
show(client,
f"""SELECT count() FROM {DB}.replacing_user""",
"OPTIMIZE FINAL 后的行数(等于唯一 user_id 数)")
def demo_summing(client):
banner("② SummingMergeTree —— 数值自动累加")
show(client,
f"SELECT count() FROM {DB}.summing_daily",
"Merge 前的物理行数")
show(client,
f"""SELECT event_date, event_name, count() AS raw_rows, sum(pv) AS pv_sum
FROM {DB}.summing_daily
WHERE event_name = 'click'
GROUP BY event_date, event_name
ORDER BY event_date LIMIT 5""",
"正确查询仍要 GROUP BY + sum(),因为 Merge 是异步")
try:
client.command(f"OPTIMIZE TABLE {DB}.summing_daily FINAL")
except Exception:
pass
show(client,
f"SELECT count() FROM {DB}.summing_daily",
"OPTIMIZE FINAL 后的行数(相同 (date,event,user) 被累加成一条)")
def demo_aggregating(client):
banner("③ AggregatingMergeTree —— 实时 PV/UV 大屏")
show(client,
f"SELECT count() FROM {DB}.agg_events_raw",
"明细表行数")
show(client,
f"SELECT count() FROM {DB}.agg_events_state",
"状态表行数(被物化视图写入)")
show(client,
f"""SELECT event_date, url,
sumMerge(pv_state) AS pv,
uniqMerge(uv_state) AS uv
FROM {DB}.agg_events_state
WHERE event_date = '2024-01-15'
GROUP BY event_date, url
ORDER BY pv DESC LIMIT 5""",
"大屏查询:只扫千级状态行,毫秒返回")
# 对照:直接在原始表上算
show(client,
f"""SELECT event_date, url, count() AS pv, uniq(user_id) AS uv
FROM {DB}.agg_events_raw
WHERE event_date = '2024-01-15'
GROUP BY event_date, url
ORDER BY pv DESC LIMIT 5""",
"原始明细直接算(对照用,生产中亿级明细会慢很多)")
def demo_collapsing(client):
banner("④ CollapsingMergeTree —— +1/-1 折叠")
show(client,
f"SELECT count() AS total_rows, sum(Sign) AS live_rows FROM {DB}.collapsing_order",
"合并前:物理行数 vs 折叠后有效行数")
show(client,
f"""SELECT order_id,
argMax(status, Sign) AS latest_status,
sum(amount * Sign) AS live_amount
FROM {DB}.collapsing_order
GROUP BY order_id HAVING sum(Sign) > 0
LIMIT 5""",
"【推荐】sum(col*Sign) + HAVING sum(Sign)>0 —— 通用折叠查询套路")
try:
client.command(f"OPTIMIZE TABLE {DB}.collapsing_order FINAL")
except Exception:
pass
show(client,
f"SELECT count() FROM {DB}.collapsing_order",
"OPTIMIZE FINAL 后,物理折叠剩余行数")
def demo_versioned(client):
banner("⑤ VersionedCollapsingMergeTree —— 允许乱序写入")
show(client,
f"SELECT count() FROM {DB}.versioned_order",
"乱序写入后的物理行数")
show(client,
f"""SELECT order_id,
argMax(status, Version) AS latest_status,
max(Version) AS ver
FROM {DB}.versioned_order GROUP BY order_id
HAVING sum(Sign) > 0 LIMIT 5""",
"带 Version 的查询(即使乱序,按 Version 取最新)")
try:
client.command(f"OPTIMIZE TABLE {DB}.versioned_order FINAL")
except Exception:
pass
show(client,
f"SELECT count() FROM {DB}.versioned_order",
"OPTIMIZE FINAL 后,折叠剩余行数(等于独立 order_id 数)")
def main() -> int:
client = clickhouse_connect.get_client(
host=HOST, port=PORT, username=USER, password=PASSWORD, database=DB)
print(f"Connected. Database = {DB}")
demo_replacing(client)
time.sleep(1)
demo_summing(client)
time.sleep(1)
demo_aggregating(client)
time.sleep(1)
demo_collapsing(client)
time.sleep(1)
demo_versioned(client)
banner("🎉 全部演示完成。要点回顾:")
print(" · Replacing → argMax(col, ver) / SELECT ... FINAL(慎用)")
print(" · Summing → GROUP BY + sum() (Merge 只是减少扫描行数)")
print(" · Aggregating → xxMerge(state) 取最终值,毫秒级大屏")
print(" · Collapsing → sum(col*Sign) + HAVING sum(Sign) > 0")
print(" · VersionedColl→ 允许乱序,按 (pk, Version) 折叠")
return 0
if __name__ == "__main__":
sys.exit(main())