ClickHouseNotes

第 10 章:CollapsingMergeTree

zjc 于 2026-01-10 发布

这是《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。

思考题

  1. 正负记录为什么必须完全镜像?
  2. 普通 Collapsing 对写入顺序有什么要求?
  3. VersionedCollapsing 解决什么问题?
  4. 为什么分区字段不能随状态变化?
  5. Replacing 和 Collapsing 如何选择?