Learn
ClickHouse/07-replacing-collapsing

Replacing 与 Collapsing 引擎

MySQL 里一句 UPDATE orders SET status='paid' WHERE order_id=1 就完事。ClickHouse 没有真正的原地更新——列存的物理结构决定了改一行要重写整个数据块。

这一章讲 ClickHouse 的解法:用追加代替更新,在 merge 时消除重复。

1. 为什么不能直接 UPDATE

ClickHouse 确实有 ALTER TABLE ... UPDATE 语法(叫 mutation):

ALTER TABLE orders UPDATE status = 'paid' WHERE order_id = 12345;

但它的执行过程是:

1. 找出所有可能包含 order_id = 12345 的 part
2. 对每个 part,读出所有列的所有行
3. 逐行判断,改掉命中的行
4. 写出一个全新的 part
5. 把旧 part 标记为 inactive
 
改 1 行 → 可能重写 100 GB

而且它是异步的:语句立刻返回,实际执行在后台,通过 system.mutations 查进度。执行期间查询可能读到旧值。

⚠️mutation 绝不能用于高频更新

生产上见过的事故:应用把 ClickHouse 当 MySQL 用,每秒几十个 ALTER TABLE ... UPDATE。结果 mutation 队列堆积上万条,后台线程 100% CPU 重写 part,磁盘 I/O 打满,集群彻底不可用,且 mutation 队列无法快速清空(KILL MUTATION 只能一条条杀)。

mutation 的正确用法:低频的批量数据修正,比如 GDPR 删除某个用户的全部数据、修复一次导入错误。频率应该是「一周几次」而不是「一秒几次」。

2. ReplacingMergeTree:用追加实现更新

思路转变:不改老行,插入新行,让引擎在 merge 时把同一个主键的旧行丢掉。

2.1 基本用法

CREATE TABLE orders
(
    order_id   UInt64,
    user_id    UInt64,
    order_time DateTime,
    amount     Decimal(18, 2),
    status     LowCardinality(String),
    country    LowCardinality(String),
    updated_at DateTime DEFAULT now()      -- 版本列
)
ENGINE = ReplacingMergeTree(updated_at)    -- 保留 updated_at 最大的那行
PARTITION BY toYYYYMM(order_time)
ORDER BY (order_id);                       -- 去重依据 = ORDER BY 的全部列

关键点:

  • 去重依据是 ORDER BY 的全部列,不是 PARTITION BY,也不是某个"主键"声明
  • 引擎参数 updated_at 是版本列:同一个 ORDER BY 键的多行,保留版本列最大的那行
  • 不写版本列参数时,保留最后插入的那行(依赖 part 的块号,不完全可靠)

2.2 演示

-- 第一次写入订单
INSERT INTO orders VALUES
    (1001, 88, '2024-06-15 10:00:00', 299.00, 'created', 'CN', '2024-06-15 10:00:00'),
    (1002, 91, '2024-06-15 10:05:00', 158.50, 'created', 'US', '2024-06-15 10:05:00');
 
-- 订单 1001 支付了,插入新版本(不是 UPDATE)
INSERT INTO orders VALUES
    (1001, 88, '2024-06-15 10:00:00', 299.00, 'paid', 'CN', '2024-06-15 10:12:33');
 
-- 订单 1001 发货了
INSERT INTO orders VALUES
    (1001, 88, '2024-06-15 10:00:00', 299.00, 'shipped', 'CN', '2024-06-16 09:00:00');
 
SELECT * FROM orders ORDER BY order_id, updated_at;
┌─order_id─┬─user_id─┬──────────order_time─┬─amount─┬─status──┬─country─┬──────────updated_at─┐
│     1001 │      88 │ 2024-06-15 10:00:00 │ 299.00 │ created │ CN      │ 2024-06-15 10:00:00 │
│     1001 │      88 │ 2024-06-15 10:00:00 │ 299.00 │ paid    │ CN      │ 2024-06-15 10:12:33 │
│     1001 │      88 │ 2024-06-15 10:00:00 │ 299.00 │ shipped │ CN      │ 2024-06-16 09:00:00 │
│     1002 │      91 │ 2024-06-15 10:05:00 │ 158.50 │ created │ US      │ 2024-06-15 10:05:00 │
└──────────┴─────────┴─────────────────────┴────────┴─────────┴─────────┴─────────────────────┘

三行都在! 因为它们在不同的 part 里,还没合并。

强制合并后:

OPTIMIZE TABLE orders FINAL;
SELECT * FROM orders ORDER BY order_id;
┌─order_id─┬─user_id─┬──────────order_time─┬─amount─┬─status──┬─country─┬──────────updated_at─┐
│     1001 │      88 │ 2024-06-15 10:00:00 │ 299.00 │ shipped │ CN      │ 2024-06-16 09:00:00 │
│     1002 │      91 │ 2024-06-15 10:05:00 │ 158.50 │ created │ US      │ 2024-06-15 10:05:00 │
└──────────┴─────────┴─────────────────────┴────────┴─────────┴─────────┴─────────────────────┘

2.3 最关键的事实:去重是"最终一致"的

时刻     part 状态                        SELECT * 看到的
─────────────────────────────────────────────────────────
T0      part1: [1001,created]            1 行
T1      part1, part2: [1001,paid]        2 行(重复!)
T2      part1, part2, part3              3 行(重复!)
T3      后台 merge → part1_3: [shipped]  1 行

在 merge 发生之前,查询会看到重复行。 merge 何时发生完全不可预测——可能几秒,可能几小时,可能永远不会(如果两个 part 大小差异太大,合并策略不选它们)。

而且:不同分区的 part 永远不会合并。如果同一个 order_id 的两个版本落在不同月份的分区(比如 order_time 被改过),它们永远不会去重。

⚠️ReplacingMergeTree 不保证去重,只保证最终会去重

这是最常见的误解。永远不要假设「查出来就是去重的」。必须在查询层面处理。

3. 查询时正确去重的三种方式

3.1 FINAL 修饰符

SELECT * FROM orders FINAL WHERE order_id = 1001;

FINAL 让 ClickHouse 在查询时临时做一次归并,返回去重后的结果。

代价:

  • 需要读取所有相关 part 的全部 ORDER BY 列 + 版本列做归并
  • 归并过程串行度高,早期版本几乎完全单线程
  • 无法使用某些查询优化(比如 count() 直接读 count.txt)

24.x 版本对 FINAL 做了大量优化(do_not_merge_across_partitions_select_final、并行 FINAL),但仍然比不带 FINAL 慢 2~10 倍。

-- 关键优化设置
SELECT * FROM orders FINAL
SETTINGS do_not_merge_across_partitions_select_final = 1,   -- 分区内独立归并,可并行
         max_final_threads = 8;
💡FINAL 也遵守 WHERE 下推

SELECT * FROM orders FINAL WHERE order_id = 1001 会先用主键索引缩小范围,只对命中的 granule 做归并,比全表 FINAL 快得多。FINAL 加上高选择性 WHERE 是可以接受的。

3.2 用 argMax 手动去重(推荐)

在很多场景下,用聚合手动实现去重比 FINAL 更快、更可控:

SELECT
    order_id,
    argMax(user_id,    updated_at) AS user_id,
    argMax(status,     updated_at) AS status,
    argMax(amount,     updated_at) AS amount,
    argMax(country,    updated_at) AS country,
    max(updated_at)                AS updated_at
FROM orders
WHERE order_time >= '2024-06-01'
GROUP BY order_id;

argMax(x, y) 返回 y 最大那行的 x 值。这就是"取最新版本"。

优势:

  • 走标准的向量化聚合路径,可以完全并行
  • 可以和其他聚合、过滤自由组合
  • 不受 part 分布影响

劣势:列多时 SQL 很长(可以用视图封装)。

-- 用普通视图封装,业务方直接查视图
CREATE VIEW orders_latest AS
SELECT
    order_id,
    argMax(user_id,    updated_at) AS user_id,
    argMax(order_time, updated_at) AS order_time,
    argMax(amount,     updated_at) AS amount,
    argMax(status,     updated_at) AS status,
    argMax(country,    updated_at) AS country,
    max(updated_at)                AS updated_at
FROM orders
GROUP BY order_id;
 
SELECT * FROM orders_latest WHERE country = 'CN' AND status = 'paid';

3.3 LIMIT 1 BY

SELECT *
FROM orders
ORDER BY order_id, updated_at DESC
LIMIT 1 BY order_id;

LIMIT n BY expr 是 ClickHouse 特有语法:对每个 expr 分组只取 n 行。配合 ORDER BY ... DESC 就是取最新。

它比 FINAL 灵活(可以取每组前 3 条),但需要全量排序,大数据量时不如 argMax。

3.4 三种方式对比

方式1 亿行耗时适合
FINAL4.2 s需要 SELECT *、有高选择性 WHERE
argMax + GROUP BY1.1 s大批量分析查询,只需部分列
LIMIT 1 BY6.8 s需要每组取前 N 条

4. is_deleted:ReplacingMergeTree 实现删除

ClickHouse 23.2 起,ReplacingMergeTree 支持第二个参数表示删除标记:

CREATE TABLE orders
(
    order_id   UInt64,
    user_id    UInt64,
    order_time DateTime,
    amount     Decimal(18, 2),
    status     LowCardinality(String),
    country    LowCardinality(String),
    updated_at DateTime,
    is_deleted UInt8 DEFAULT 0            -- 删除标记
)
ENGINE = ReplacingMergeTree(updated_at, is_deleted)
ORDER BY (order_id)
SETTINGS clean_deleted_rows = 'Always';   -- merge 时物理删除
 
-- "删除"订单 1002:插入一条 is_deleted = 1 的行
INSERT INTO orders VALUES
    (1002, 91, '2024-06-15 10:05:00', 158.50, 'cancelled', 'US', now(), 1);
 
SELECT * FROM orders FINAL;   -- 1002 消失了

这让 ClickHouse 可以完整承接 MySQL CDC(binlog)同步:INSERT / UPDATE 都是插入新版本,DELETE 是插入 is_deleted = 1。

ℹ️CDC 同步的标准模式
MySQL binlog          Kafka          ClickHouse ReplacingMergeTree
─────────────────────────────────────────────────────────────────
INSERT order 1001  →  msg  →  INSERT (1001, ..., ts=T1, deleted=0)
UPDATE order 1001  →  msg  →  INSERT (1001, ..., ts=T2, deleted=0)
DELETE order 1001  →  msg  →  INSERT (1001, ..., ts=T3, deleted=1)
                                       ↓ merge 后
                                    该 order_id 完全消失

版本列用 binlog 的 GTID 序号或事件时间戳,保证乱序消息也能正确收敛。

5. CollapsingMergeTree:用正负行折叠

另一条技术路线。每行带一个 sign 列,值为 1(状态行)或 -1(取消行)。merge 时,同一个 ORDER BY 键的 +1 和 -1 会互相抵消。

CREATE TABLE order_stats
(
    order_id UInt64,
    user_id  UInt64,
    amount   Decimal(18, 2),
    status   LowCardinality(String),
    sign     Int8
)
ENGINE = CollapsingMergeTree(sign)
ORDER BY (order_id);

更新流程:

-- 1. 初始状态
INSERT INTO order_stats VALUES (1001, 88, 299.00, 'created', 1);
 
-- 2. 状态变更:先"撤销"旧行(原样但 sign=-1),再写新行
INSERT INTO order_stats VALUES
    (1001, 88, 299.00, 'created', -1),    -- 撤销
    (1001, 88, 299.00, 'paid',     1);    -- 新状态

merge 时 created,+1 和 created,-1 抵消,只剩 paid,+1。

5.1 折叠的真正价值:聚合可以直接算

ReplacingMergeTree 求和必须先去重再求和。CollapsingMergeTree 不需要——乘上 sign 直接求和就是正确答案:

-- 不需要 FINAL,不需要去重,直接算
SELECT
    status,
    sum(amount * sign) AS total_amount,
    sum(sign)          AS order_cnt
FROM order_stats
GROUP BY status
HAVING sum(sign) > 0;
┌─status─┬─total_amount─┬─order_cnt─┐
│ paid   │       299.00 │         1 │
└────────┴──────────────┴───────────┘

即使 created,+1 和 created,-1 还没 merge,它们的和也是 0,不影响结果。这是 Collapsing 相对 Replacing 的核心优势:聚合查询天然正确,不依赖 merge。

代价是应用端必须记住旧行的全部内容才能写出撤销行。这通常要求上游有状态存储。

⚠️Collapsing 对写入顺序敏感

如果 -1 行先于 +1 行到达(消息乱序),merge 时无法正确折叠,会残留脏数据。ClickHouse 只在 +1 在前、-1 在后时才折叠。

解决方案就是下面的 VersionedCollapsingMergeTree。

6. VersionedCollapsingMergeTree

加一个版本列,让引擎能识别乱序:

CREATE TABLE order_stats_v
(
    order_id UInt64,
    user_id  UInt64,
    amount   Decimal(18, 2),
    status   LowCardinality(String),
    sign     Int8,
    version  UInt64
)
ENGINE = VersionedCollapsingMergeTree(sign, version)
ORDER BY (order_id);
 
INSERT INTO order_stats_v VALUES
    (1001, 88, 299.00, 'paid',    1, 2),    -- 后到的新状态先写入
    (1001, 88, 299.00, 'created', -1, 1),   -- 撤销行后到
    (1001, 88, 299.00, 'created', 1, 1);    -- 原始行最后到

引擎按 version 排序后再折叠,无论到达顺序如何都能正确收敛。

7. 三种引擎怎么选

ReplacingMergeTreeCollapsingMergeTreeVersionedCollapsing
应用端复杂度低(只写新行)高(要写撤销行)高
需要记住旧行否是是
聚合是否需要 FINAL是(或 argMax)否(乘 sign 求和)否
乱序容忍是(有版本列)否是
实现删除is_deleted 参数写 -1 行写 -1 行
典型场景CDC 同步、状态表实时指标累加乱序消息的指标累加

新手建议:先用 ReplacingMergeTree。它心智负担最小,配合 argMax 查询能覆盖 90% 需求。只有当「查询必须免 FINAL 且必须实时正确」时才上 Collapsing。

8. 完整实战:订单状态表

CREATE TABLE demo.orders
(
    order_id   UInt64,
    user_id    UInt64,
    order_time DateTime,
    amount     Decimal(18, 2),
    status     LowCardinality(String),
    country    LowCardinality(String),
    updated_at DateTime,
    is_deleted UInt8 DEFAULT 0
)
ENGINE = ReplacingMergeTree(updated_at, is_deleted)
PARTITION BY toYYYYMM(order_time)
ORDER BY (order_id)
SETTINGS clean_deleted_rows = 'Always';
 
-- 灌入测试数据:5 万订单,每个订单 1~3 个版本
INSERT INTO orders
SELECT
    1000 + number % 50000                                       AS order_id,
    1 + rand(1) % 10000                                         AS user_id,
    toDateTime('2024-06-01') + (number % 50000) % 2592000       AS order_time,
    toDecimal64(10 + rand(2) % 50000, 2) / 100                  AS amount,
    ['created','paid','shipped','cancelled'][1 + number % 4]    AS status,
    ['CN','US','JP','DE','IN'][1 + rand(3) % 5]                 AS country,
    toDateTime('2024-06-01') + number                           AS updated_at,
    0                                                            AS is_deleted
FROM numbers(120000);
-- 错误查法:会把所有版本都算进去
SELECT country, count() AS cnt, sum(amount) AS total
FROM orders GROUP BY country;
┌─country─┬───cnt─┬────────total─┐
│ CN      │ 24135 │ 60379422.31  │   ← 数量和金额都虚高
│ US      │ 23987 │ 59921033.08  │
└─────────┴───────┴──────────────┘
-- 正确查法 A:FINAL
SELECT country, count() AS cnt, sum(amount) AS total
FROM orders FINAL GROUP BY country;
 
-- 正确查法 B:argMax(更快)
SELECT country, count() AS cnt, sum(amount) AS total
FROM (
    SELECT
        order_id,
        argMax(country, updated_at) AS country,
        argMax(amount,  updated_at) AS amount
    FROM orders
    GROUP BY order_id
)
GROUP BY country;
┌─country─┬───cnt─┬────────total─┐
│ CN      │ 10012 │ 25041887.44  │   ← 每个 order_id 只算一次
│ US      │  9988 │ 24983201.17  │
└─────────┴───────┴──────────────┘
🎯练习
  1. 建上面的 orders ReplacingMergeTree 表,插入同一个 order_id 的 3 个版本(不同 updated_at)。不执行 OPTIMIZE,直接 SELECT * 看到几行?再执行 OPTIMIZE TABLE orders FINAL 后看到几行?
  2. 分别用 FINAL、argMax、LIMIT 1 BY 三种方式查询去重结果,用 SET send_logs_level='trace' 或对比 Elapsed,看哪种最快。
  3. 给表加 is_deleted 列,插入一条 is_deleted = 1 的行"删除"某订单,验证 SELECT ... FINAL 里它消失了。
  4. 用 CollapsingMergeTree 重建一张 order_stats,写入一组 +1 / -1 行,在不执行 OPTIMIZE 的情况下用 sum(amount * sign) 查询,验证结果已经正确——体会它和 Replacing 的本质区别。

小结

  • ClickHouse 没有原地更新,ALTER TABLE ... UPDATE(mutation)会重写整个 part,只能低频批量使用
  • ReplacingMergeTree 用「插入新版本 + merge 时按版本列保留最新」模拟更新,去重依据是 ORDER BY 全部列
  • 去重是最终一致的,merge 前查询会看到重复,且跨分区永不去重
  • 查询层去重三选一:FINAL 最简单、argMax 最快、LIMIT 1 BY 最灵活
  • is_deleted 参数让 ReplacingMergeTree 能完整承接 CDC 的增删改
  • CollapsingMergeTree 用 sign 折叠,优势是聚合时乘 sign 求和天然正确、不依赖 merge,代价是应用端要写撤销行
  • 乱序消息场景用 VersionedCollapsingMergeTree
  • 下一章讲 Summing 与 Aggregating,把预聚合做到引擎层 →