Skip to content

第 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_id
  • ReplacingMergeTree(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 还没跑

解决办法

  1. SELECT ... FINAL查询时实时去重,但代价很大(读时归并)。
  2. OPTIMIZE TABLE ... PARTITION ... FINAL:强制合并一次(运维操作,不能当查询手段)。
  3. 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 坑

  1. 主键设计错位:如果 ORDER BY 里没有"你要的 GROUP BY 列",Sum 不会发生。
  2. 非数值列变"随机值":有 String 列在排序键之外时,Merge 保留的是"不确定的第一条"。
  3. 溢出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}))   ← 合并了 HyperLogLog

6.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 pkOPTIMIZE PARTITION FINAL(离线)SELECT * FINAL(高频)
SummingMT 读聚合值sum(col) GROUP BY pkFINAL 仅调试用SELECT col FINAL
AggregatingMT 读最终值xxMerge(state) GROUP BY pkFINALxxMerge(xxState(col)) 等奇怪写法
CollapsingMT 读有效行sum(col*Sign) GROUP BY pk HAVING sum(Sign)>0FINALWHERE Sign=1 会漏折叠
精确去重迁 Replicated + 业务主键去重 + 用 uniqExactFINAL 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 的对比小框

场景MySQLPostgreSQLClickHouse
"主键存在则更新"INSERT ... ON DUPLICATE KEY UPDATEINSERT ... ON CONFLICT ... DO UPDATEReplacingMergeTree + argMax(最终一致)
实时累加计数UPDATE t SET c=c+1UPDATE t SET c=c+1SummingMergeTree + INSERT
实时 UV 估算业务层配 Redis HLLpg_hll 扩展AggregatingMergeTree + uniqState
CDC 同步binlog + 目标端 UPDATE逻辑复制 pub/subCollapsingMergeTree / VersionedCollapsingMergeTree
事务ACID 完整ACID 完整无跨行事务 —— 用 Sign 折叠实现逻辑事务

关键反差

  • MySQL/PG 的 UPSERT 是同步的,插入瞬间就生效;
  • CH 的"UPSERT"都是异步合并语义,查询时要配合 FINALargMax 才能看到最终结果。

为什么 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 精确去重的本质区别。

标准答案

  1. 不能做精确去重。ReplacingMergeTree 的去重是最终一致的
    • 去重发生在 Merge 时,而 Merge 是异步后台任务
    • 同主键但在不同分区的数据永远不会合并
    • 刚 INSERT 的数据,查询瞬间可能还没 Merge → 看到重复行。
  2. 强行去重的手段:
    • SELECT ... FINAL:查询时临时归并,但慢 2-10 倍;
    • OPTIMIZE PARTITION ... FINAL:强制合并(运维手段);
    • argMax(col, ver) GROUP BY pk最推荐,SQL 层表达清晰且通常比 FINAL 快。
  3. 真正要"精确唯一主键",请用 MySQL/PG 这类 OLTP,或者在业务写入前去重后再写入。

加分项:能补"23.x 起支持 ReplacingMergeTree(ver, is_deleted) 实现软删除";能讲清 "ver 列并列时取最后一条" 是实现定义而非 SQL 标准。

易错点:以为 ORDER BY user_id 就会自动去重,忽视了 Merge 的异步性。


Q2:SELECT ... FINALOPTIMIZE 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"模式。

标准答案(按数据流走一遍):

  1. 明细表 events_raw 用普通 MergeTree 接收原始事件。
  2. 聚合表 events_aggAggregatingMergeTree,列类型是 AggregateFunction(sum, UInt64) / AggregateFunction(uniq, UInt64),存聚合的中间状态(二进制)。
  3. 物化视图 events_mvevents_raw 的 INSERT 转成聚合状态写入 events_agg
    sql
    SELECT event_date, url,
           sumState(toUInt64(1)) AS pv_state,
           uniqState(user_id)    AS uv_state
    FROM events_raw
    GROUP BY event_date, url
  4. 后台 Merge 会把同 (event_date, url) 的状态合并(sumMerge / uniqMerge HyperLogLog 合并)。
  5. 大屏查询时: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:CollapsingMergeTreeVersionedCollapsingMergeTree 的区别?各自适合什么场景?

考察点: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 引擎"减少扫描但不替代聚合"的理解。

标准答案

不能。 原因:

  1. Merge 是异步的,查询瞬间可能还有未合并的 Part,SummingMergeTree 表内仍然可能有"同主键多行"。
  2. 不同分区的数据永远不合并,跨分区查询必然要 GROUP BY。
  3. Summing 的作用是"减少 GROUP BY 要扫的行数"(从明细亿级压到聚合表千万级),而不是"免去 GROUP BY"。

正确用法:在 SummingMergeTree 上依然要写 GROUP BY pk,只是扫的是聚合后的行,非常快。

加分项:能补"非数值列 + 非主键列"的合并行为是不确定的(保留同组任意一条,通常是第一条);建议要么全部进主键、要么进数值列。

易错点:看到 SummingMergeTree 就写 SELECT col FROM t 不加 GROUP BY,以为结果已经是累加完的。


Q6:什么是"聚合状态" AggregateFunction?它解决了什么问题?

考察点:聚合的分布式/增量合并能力。

标准答案

  1. 普通聚合函数(sum, uniq, quantile…)是"一次性算出最终结果";但在分布式 / 增量场景里,我们希望能把"半成品状态"保留下来,再和新状态合并。
  2. ClickHouse 为几乎所有聚合函数都提供了 3 个变种:
    • xxState(col):把行聚合成中间状态(二进制结构:累加器 / HLL / 样本…);
    • xxMerge(state) / xxMerge(state1, state2, …):把多个中间状态合并,得到最终标量;
    • xxMergeState(state):把多个中间状态合并成新的中间状态(用于链式 MV)。
  3. 用途:
    • AggregatingMergeTree 在 Merge 时自动合并同主键的聚合状态;
    • 分布式查询时,各分片返回状态,协调节点做最终 Merge;
    • 物化视图可以存状态,业务查询时再按需 Merge。
  4. 存储的是定长二进制状态,比存原始明细省几个数量级。

加分项:能讲 uniq 是 HyperLogLog、uniqCombined 内部用了 HLL+LL、quantile 是 TDigest(不同算法会影响精度);能讲 argMaxState 这种特殊聚合状态。

易错点:以为 AggregateFunction(sum, UInt64) 就是普通 UInt64 字段 —— 其实它是 blob,不能直接当数用。


Q7:为什么 ClickHouse 没有真正的 UPSERT?它用什么实现这个语义?

考察点:OLAP vs OLTP 的本质差异。

标准答案

  1. 真正的 UPSERT 需要"定位到某一行,原地改"。在列存表里,一行的不同列分散在不同的 .bin 文件里;原地改意味着"解压列块 → 找到行 → 改 → 重压缩 → 写回"——单行改动的 I/O 代价巨大,完全违背 OLAP 的"批量扫描" 本性。
  2. ClickHouse 用 "写入+异步合并" 实现逻辑 UPSERT:
    • ReplacingMergeTree 实现"主键存在则替换",取 ver 最大一条;
    • SummingMergeTree 实现"主键存在则累加数值";
    • CollapsingMergeTree(+Version) 实现"先撤销旧状态再写新状态";
    • 都是"写多行 + 合并时收敛"而非"原地改"。
  3. 副作用:强一致场景要么用 argMax / sum(x*Sign) 在查询时计算、要么接受"最终一致"的语义。
  4. 还有 ALTER TABLE ... UPDATE/DELETEMutation 机制,但本质是"重写整个 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_insertmax_threadsFORMAT 全家桶,解决"怎么把数据高速灌进 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())

family_play.py ↗