这是《ClickHouse 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 CollapsingMergeTree 用正负记录抵消旧状态,适合把状态变更表达为事件流。它比“不断更新一行”更贴合 ClickHouse 的追加模型,但要求 sign 语义严格。
10.1 抵消模型
INSERT +1: user 1001 level = 1
INSERT -1: user 1001 level = 1
INSERT +1: user 1001 level = 2
排序键相同的 +1 和 -1 在后台合并时抵消,最终留下:
user 1001 level = 2
正负记录必须是同一个旧状态的镜像,不能一边 level=1,一边 level=2。
10.2 建表
CREATE TABLE analytics.user_state_local
(
user_id UInt64,
level UInt8,
vip UInt8,
updated_at DateTime,
sign Int8
)
ENGINE = CollapsingMergeTree(sign)
PARTITION BY toYYYYMM(updated_at)
ORDER BY (user_id, level, vip, updated_at);
如果只按 user_id 查询,排序键可以简化:
ORDER BY user_id
但所有参与抵消的业务字段必须保证正负行排序键一致。
10.3 写入示例
INSERT INTO analytics.user_state_local VALUES
(1001, 1, 0, '2026-08-25 10:00:00', 1);
-- 抵消旧状态
INSERT INTO analytics.user_state_local VALUES
(1001, 1, 0, '2026-08-25 10:00:00', -1);
-- 新状态
INSERT INTO analytics.user_state_local VALUES
(1001, 2, 1, '2026-08-25 10:05:00', 1);
查询当前状态:
SELECT
user_id,
level,
vip,
updated_at
FROM analytics.user_state_local
GROUP BY user_id, level, vip, updated_at
HAVING sum(sign) > 0;
更简单的聚合方式:
SELECT user_id, sum(sign) AS state_count
FROM analytics.user_state_local
GROUP BY user_id;
复杂字段可以配合 argMax(field, updated_at) 使用,但必须先确认正负记录语义。
10.4 分区问题
如果按 updated_at 分区,旧状态和新状态可能落入不同分区。不同分区之间不会互相抵消,查询时必须跨分区处理。
更稳妥的设计:
使用固定业务日期或首次创建日期分区
例如:
CREATE TABLE analytics.user_state_fixed
(
user_id UInt64,
create_date Date,
level UInt8,
updated_at DateTime,
sign Int8
)
ENGINE = CollapsingMergeTree(sign)
PARTITION BY toYYYYMM(create_date)
ORDER BY (user_id, create_date, level, updated_at);
10.5 乱序与 VersionedCollapsingMergeTree
普通 CollapsingMergeTree 对写入顺序有要求。Kafka 分区重放、并发写多副本、ETL 重跑都可能导致负记录晚到或乱序。
带版本表:
CREATE TABLE analytics.user_state_versioned
(
user_id UInt64,
level UInt8,
updated_at DateTime,
version UInt64,
sign Int8
)
ENGINE = VersionedCollapsingMergeTree(sign, version)
ORDER BY user_id;
引擎先按版本整理,再按 sign 抵消。版本仍然必须来自可信来源,例如业务版本号或有序事件序号。
10.6 生成正负记录
CDC 示例:
旧状态:level=1, version=10
变更后:level=2, version=11
输出:
-1, level=1, version=10 抵消旧状态
+1, level=2, version=11 写入新状态
Flink 伪代码:
if (value.isUpdate()) {
State old = stateStore.get(value.getKey());
if (old != null) {
out.collect(old.withSign(-1));
}
out.collect(value.withSign(1));
stateStore.put(value.getKey(), value);
}
如果使用无状态作业,需要先从数据库、状态主题或 Lookup 表获取旧值。
10.7 查询模式
当前状态:
SELECT user_id, sum(sign) AS cnt
FROM analytics.user_state_local
GROUP BY user_id
HAVING cnt > 0;
统计各等级人数:
SELECT level, sum(sign) AS users
FROM analytics.user_state_local
GROUP BY level;
使用 FINAL:
SELECT *
FROM analytics.user_state_local
FINAL;
大表 FINAL 成本高。查询前先缩小分区或键范围。
10.8 与 ReplacingMergeTree 对比
| 维度 | Replacing | Collapsing |
|---|---|---|
| 数据形态 | 多版本快照 | 正负记录 |
| 合并结果 | 保留最新版本 | 抵消正负 |
| 查询 | argMax / FINAL | sum(sign) / FINAL |
| 删除 | 写入墓碑版本 | 写 -1 |
| 乱序 | 版本可判断 | Versioned 变体 |
如果上游能提供完整最新行,ReplacingMergeTree 更简单。如果上游只有状态变化事件,Collapsing 更贴近事件语义。
10.9 常见错误
| 错误 | 结果 |
|---|---|
| 负记录字段与旧状态不一致 | 无法抵消 |
| 忘写 -1 | 状态重复 |
| 只写 -1 | 状态丢失 |
| sign 不在 Int8 中表达 | 类型错误 |
| 分区随状态变化 | 跨分区不抵消 |
| 查询不用 sign | 结果不正确 |
排查:
SELECT
user_id,
level,
sign,
count()
FROM analytics.user_state_local
WHERE user_id = 1001
ORDER BY updated_at, sign;
本章小结
CollapsingMergeTree 把状态更新转换为追加正负记录,适合事件化状态流。设计关键是镜像旧状态、保证 sign 正确、稳定版本和分区。乱序场景优先选择 VersionedCollapsingMergeTree,查询层必须显式处理 sign。
思考题
- 正负记录为什么必须完全镜像?
- 普通 Collapsing 对写入顺序有什么要求?
- VersionedCollapsing 解决什么问题?
- 为什么分区字段不能随状态变化?
- Replacing 和 Collapsing 如何选择?