CollapsingMergeTree:如何在ClickHouse中无需UPDATE更新聚合数据
1. 为什么需要CollapsingMergeTree——更新聚合数据的问题
让我们回到在线赌场的场景。每个玩家都有一个余额。当玩家下注时,余额减少;当玩家赢钱时,余额增加;当管理员取消一笔可疑交易时,余额再次变化。
在传统数据库(如PostgreSQL)中,你只需执行 UPDATE players SET balance = balance - 100 WHERE user_id = 123。简单明了。
但ClickHouse无法更新数据。完全不行。为什么?因为ClickHouse是为分析而构建的,数据只追加。更新对列式存储来说很痛苦,因为数据存储在压缩块中。要更改一个单元格,需要重写整个块。
那么如何更改余额? 你不修改旧记录。而是添加一条新记录,表示“取消之前的更改”并“添加新更改”。这称为通过取消实现变更物化。
现实类比: 想象一个会计分类账,所有条目都用墨水书写。你不能擦除和更正。相反,你在底部写上:“第45行——错误,已取消。新第46行——正确金额。”然后在计算总额时,你读取所有行,考虑取消项。
CollapsingMergeTree 是一个ClickHouse引擎,它在后台合并期间自动折叠“+1”和“-1”对。它就像那个会计:看到一对取消行,就丢弃这两行。
2. sign列的原理——取消的数学原理
CollapsingMergeTree 的主要思想是一个特殊的 sign 列,有两个可能的值:
+1— “添加”(当前版本)-1— “取消”(过时版本)
在数据部分合并期间,ClickHouse会查找具有相同排序键(ORDER BY)且一个 sign = +1、另一个 sign = -1 的行对。当找到这样的对时,两行都会被删除。只有没有配对的行保留——即那些没有被取消的行。
为什么这样有效: 任何数据更改都表示为取消旧版本并添加新版本。对 (+1, -1) 的和为零。当你执行 SUM(amount * sign) 时,旧版本相互抵消,新版本保留。
与复式记账的类比: 在会计中,每笔交易记录两次:借方和贷方。CollapsingMergeTree 做同样的事情——每个更改都有其对立面。求和时,它们相互抵消。
3. CREATE TABLE——分解说明
-- 创建玩家余额变更历史表
CREATE TABLE player_balance
(
user_id UInt64, -- 玩家ID
date Date, -- 余额变更日期
amount Int64, -- 余额变更(+100, -50等)
balance_after Int64, -- 操作后余额(可选)
sign Int8, -- +1 — 添加, -1 — 取消
updated_at DateTime DEFAULT now() -- 操作时间戳
)
ENGINE = CollapsingMergeTree(sign) -- 折叠引擎,指定sign列
ORDER BY (user_id, date) -- 分组和折叠的键
重要事项:
ENGINE = CollapsingMergeTree(sign)— 唯一必需的参数是sign列的名称(通常是sign或is_active)。该列必须是Int8类型(-128到127的整数),但实际只使用+1和-1。ORDER BY (user_id, date)— 此键中的列决定了哪些行被视为“对”。两行在ORDER BY所有列上具有相同值且sign相反(+1和-1)时,将被折叠。
如果ORDER BY不包含所有必要字段会怎样? 例如,如果不包含 user_id,来自不同用户的行可能会被折叠——这将是一场灾难。所有用于区分操作的字段都必须放在 ORDER BY 中。
为什么不用PRIMARY KEY? 与其他MergeTree引擎相同的原因——ORDER BY 控制物理顺序和合并,而 PRIMARY KEY(如果指定)只控制索引。
4. 插入——如何正确更新数据
假设玩家余额为1000卢布。他们下注100卢布。你不在单行中更改余额,而是执行两次插入:
-- 步骤1:取消旧版本的余额(原来是1000,现在应为900)
-- 旧版本:user_id=123, date='2025-06-01', amount = 1000(操作前余额)
-- 要取消,插入一行 sign = -1
INSERT INTO player_balance VALUES
(123, '2025-06-01', 1000, 1000, -1, now()); -- 取消旧余额
-- 步骤2:添加新版本的余额(下注后为900)
INSERT INTO player_balance VALUES
(123, '2025-06-01', -100, 900, +1, now()); -- 新版本:变更-100,结果900
但这不方便。 实际上,用“变更”而非“完整余额”来思考更容易。以下是一个更典型的模式:
-- 玩家下注100卢布(余额减少)
-- 只插入一行 sign = +1,其中 amount 是余额变更(-100)
INSERT INTO player_balance VALUES
(123, '2025-06-01', -100, 900, +1, now());
-- 如果需要取消此下注(例如由于技术错误)
-- 插入一对取消行
INSERT INTO player_balance VALUES
(123, '2025-06-01', -100, 900, -1, now()), -- 取消下注
(123, '2025-06-01', +100, 1000, +1, now()); -- 恢复余额
为什么这样有效: 当合并发生时,ClickHouse会找到具有相同 user_id、date 和 amount(如果amount在ORDER BY中)且 sign 不同的行对——并删除它们。只有当前余额保留。
重要: 你必须自行确保对的正确性。ClickHouse不会验证变更总和是否平衡。它只是折叠具有相反sign和相同ORDER BY键的行。
5. 使用SUM of sign的SELECT——如何正确读取
读取数据时,你需要考虑sign进行聚合。主要模式:
-- 获取每个玩家的当前余额
SELECT
user_id,
SUM(amount * sign) AS current_balance
FROM player_balance
WHERE sign != 0 -- 过滤掉意外的0(不应存在)
GROUP BY user_id;
逻辑分解:
amount * sign— 如果 sign = +1,项为amount;如果 sign = -1,项为-amount(取消前一个)SUM(...)— 所有取消对在求和中相互抵消GROUP BY user_id— 按用户聚合
为什么不用FINAL? 与 ReplacingMergeTree 不同,对于 CollapsingMergeTree,你始终使用 SUM(amount * sign) 进行聚合。这在合并前后都能正确工作,因为sign数学不依赖于行是否被物理折叠。
现实示例:
-- 初始数据(合并前):
-- (123, -100, +1) — 下注100卢布
-- (123, -100, -1) — 取消下注
-- (123, +100, +1) — 恢复余额
-- 查询:SUM(amount * sign) = (-100*1) + (-100*-1) + (100*1) = -100 + 100 + 100 = 100
-- 正确:余额增加了100(取消下注退还了钱)
如果你想查看不带聚合的历史记录(例如,按时间顺序的所有操作),直接 SELECT * 将显示所有行,包括已取消的。这很正常——这是设计使然。
6. 陷阱——行顺序问题
最大的陷阱: CollapsingMergeTree 要求具有相同键的行按正确顺序到达——先+1,后-1(还是反过来?我们来弄清楚)。
ClickHouse不检查时间戳。它在合并期间查看每个部分内行的顺序。如果一个部分包含一对 (+1, -1),它将折叠它们。但如果+1在一个部分,-1在另一个部分,它们不会折叠,直到这些部分合并为一个(这可能需要一段时间)。
为什么在分布式系统中这是一个问题?
想象一下,你的数据通过Kafka(消息队列)来自三个不同的服务器。服务器#1发送了“下注100卢布”(+1)。服务器#2发送了“取消下注”(-1)。服务器#3发送了“恢复余额”(+1)。它们可能以不同的顺序进入不同的ClickHouse部分。
如果一个部分只包含+1,另一个部分包含-1,余额将暂时不正确(总和会显示额外的钱)。当部分合并时,数据会自行纠正,但这可能需要一个小时。
类比: 就像把信件放进两个不同的文件夹。一个文件夹有“债务100卢布”(+1),另一个有“债务免除”(-1)。在合并文件夹之前,你的会计系统会认为你被欠了100卢布。
7. VersionedCollapsingMergeTree——解决顺序问题
为了解决顺序问题,ClickHouse开发者添加了 VersionedCollapsingMergeTree。它增加了第三列——版本(通常是 version 或 timestamp)。
CREATE TABLE player_balance_versioned
(
user_id UInt64,
date Date,
amount Int64,
version UInt64, -- 单调递增的版本号
sign Int8
)
ENGINE = VersionedCollapsingMergeTree(sign, version) -- 两个参数!
ORDER BY (user_id, date);
工作原理:
- 在合并期间,ClickHouse查找具有相同
ORDER BY键和相同版本的对 (+1, -1)。 - 如果行具有相同键但不同版本,则不会折叠。相反,保留版本最高的行(最新状态)。
- 版本允许即使行乱序到达也能正确处理——只要取消操作与原始操作具有相同版本。
为什么这解决了顺序问题: 即使+1在-1之后到达,ClickHouse会看到它们具有不同版本(或相同——则折叠)。如果版本相同,无论部分中的物理顺序如何,对都会折叠。如果版本不同,较新的版本保留。
使用版本的示例:
-- 操作1:下注100卢布(版本1001)
INSERT INTO player_balance_versioned VALUES (123, '2025-06-01', -100, 1001, +1);
-- 取消同一笔下注(相同版本1001,sign = -1)
INSERT INTO player_balance_versioned VALUES (123, '2025-06-01', -100, 1001, -1);
-- 新的正确下注50卢布(版本1002)
INSERT INTO player_balance_versioned VALUES (123, '2025-06-01', -50, 1002, +1);
现在,即使所有三行以不同顺序到达不同部分,在合并时,具有相同版本的对会折叠。只有版本1002的50卢布下注保留。
8. 实际用例:赌场中实时玩家余额
想象你有一个“余额”微服务,必须显示当前玩家余额精确到分,延迟不超过5秒。
工作流程:
-- 所有余额操作的表
CREATE TABLE balance_operations
(
user_id UInt64,
operation_id String, -- 唯一操作ID(下注、支付、取消)
amount Int64, -- 变更(+1000赢,-500下注)
version UInt64, -- 单调版本号
sign Int8, -- +1 = 新操作,-1 = 取消
created_at DateTime DEFAULT now()
)
ENGINE = VersionedCollapsingMergeTree(sign, version)
ORDER BY (user_id, operation_id, version);
场景1:玩家下注100卢布
-- 插入一行(sign = +1)
INSERT INTO balance_operations VALUES (123, 'bet_001', -100, 1001, +1, now());
场景2:玩家赢500卢布(支付)
INSERT INTO balance_operations VALUES (123, 'win_001', +500, 1002, +1, now());
场景3:管理员取消下注bet_001(玩家作弊?)
-- 取消旧下注(相同operation_id,sign = -1,相同版本1001)
INSERT INTO balance_operations VALUES (123, 'bet_001', -100, 1001, -1, now());
-- 添加纠正操作(退还100卢布)
INSERT INTO balance_operations VALUES (123, 'admin_correction_bet_001', +100, 1003, +1, now());
如何实时读取余额:
-- 仪表盘查询(每3秒运行一次)
SELECT
user_id,
SUM(amount * sign) AS current_balance
FROM balance_operations
WHERE user_id = 123 AND created_at > now() - interval 1 day -- 按时间限制
GROUP BY user_id;
正确使用版本后,即使插入顺序混乱,此查询也会返回正确的余额。
9. 方法比较:CollapsingMergeTree vs ReplacingMergeTree 用于余额
许多初学者问:“为什么不用 ReplacingMergeTree 并将余额作为行版本来更新?”
让我们比较一下。
| 特性 | CollapsingMergeTree | ReplacingMergeTree |
|---|---|---|
| 如何表示变更 | 两行:取消(-1)和新(+1) | 一行新行,版本更高 |
| 是否需要存储完整历史 | 是,直到折叠 | 是,直到合并 |
| 读取方式 | SUM(amount * sign) |
argMax(amount, version) 或 FINAL |
| 插入复杂度 | 更高(需要考虑对) | 更低(只需新版本) |
| 读取复杂度 | 更低(简单聚合) | 更高(FINAL慢或argMax) |
| 错误风险 | 对不匹配(逻辑错误) | 版本不单调(客户端错误) |
何时选择CollapsingMergeTree:
- 你需要频繁更改相同键(例如,玩家余额每小时变化100次)。
- 你想使用简单聚合
SUM(amount * sign)并且不想依赖FINAL。 - 你控制插入顺序或使用
VersionedCollapsingMergeTree。 - 你需要回滚操作(取消下注)——在
ReplacingMergeTree中,这需要插入一个版本更高的新行,不能明确反映“取消”。
何时ReplacingMergeTree更好:
- 你不频繁更新(例如,订单状态:已创建→已支付→已送达)。
- 你存储非数值的可变属性而不是数值聚合。
- 你需要查看每条记录的版本历史。
余额示例——哪个更好? 对于高负载账户(每秒数千次下注),VersionedCollapsingMergeTree 更好。它提供可预测的性能和在混乱顺序下的正确行为。
10. 性能及何时优于PostgreSQL的UPDATE
CollapsingMergeTree的性能
- 插入: 非常快(常规INSERT,无锁)。你为存储两行而不是更新一行付出代价——但在列式数据库中,这并不算太糟。
- 带聚合的读取: ClickHouse只读取
amount和sign列(列式存储!),执行快速向量化计算。对于十亿行——不到一秒。 - 合并: 后台工作。不影响插入。
与PostgreSQL的比较
在PostgreSQL中更新余额:
-- 带行锁的原子更新
UPDATE players SET balance = balance - 100 WHERE user_id = 123;
优点:简单,ACID保证(原子性、一致性、隔离性、持久性),即时一致性。
缺点:每秒10,000次更新时——锁(行锁)、WAL(预写日志)、vacuum。你会遇到IO限制。
在ClickHouse中使用CollapsingMergeTree:
优点:单台服务器每秒100,000+次插入,数据压缩(10:1),无锁,线性扩展。
缺点:无即时一致性(合并前需要聚合),逻辑更复杂(sign、版本),最终一致性——系统会达到正确状态,但不是立即。
何时CollapsingMergeTree优于PostgreSQL:
- 你需要非常多的更新(每秒数千到数万次)。
- 几秒的延迟(用于折叠)是可接受的。
- 你已经在使用ClickHouse进行分析。
何时PostgreSQL仍然更好:
- 你需要严格的即时一致性(银行账户间转账)。
- 更新很少(<1000次/秒)。
- 你不想使架构复杂化。
下一步
现在你了解了CollapsingMergeTree及其大哥VersionedCollapsingMergeTree。接下来要探索的主题:
- 如何在CollapsingMergeTree和ReplacingMergeTree之间选择——每个任务的检查清单。
- 优化合并——设置如
merge_with_ttl_timeout使对更快折叠。 - 模式:物化视图 + CollapsingMergeTree——用于多级聚合。
总结: CollapsingMergeTree是一个强大但需要纪律的工具。它不会原谅插入顺序或对正确性方面的错误。但如果正确设置(尤其是使用版本),它能提供传统数据库无法企及的性能。记住黄金法则:始终用 SUM(amount * sign) 检查查询,并在分布式系统中使用 VersionedCollapsingMergeTree。
← 上一篇: SummingMergeTree 与 AggregatingMergeTree:轻松实现增量聚合
→ 下一篇: ClickHouse 中的分区:如何在文件夹级别管理数据
— Editorial Team
暂无评论。