ReplacingMergeTree:如何在 ClickHouse 中轻松去重
1. 为什么需要 ReplacingMergeTree——现实中的重复问题
想象你正在开发一个在线赌场。玩家点击“下注”按钮——1000 卢布押黑色。此时,处理请求的服务器突然崩溃(过热、网络故障,谁知道呢)。客户端没有收到响应,心想:“下注没成功。”玩家再次点击。服务器恢复并接受了两个请求。数据库中——两条相同的下注记录。玩家愤怒了:被扣了 2000 卢布而不是 1000。
这是一个典型的幂等性问题(源自拉丁语 idem——相同,potens——能力)。如果一个操作重复执行的结果与执行一次相同,则该操作是幂等的。在数据库世界中,我们需要一种机制来判断:“我已经见过这个下注了,我会忽略第二个版本。”
在 ClickHouse 中,有 ReplacingMergeTree 用于此目的。它是一个表引擎,在数据部分合并期间自动删除重复项。但我提前警告你:这不是魔法——它有我们即将讨论的怪癖。
现实类比: ReplacingMergeTree 就像一位记录会议记录的秘书。人们带着请求来找你。有时同一个客户带来两份相同的申请(例如,错过了火车并请求退款,然后又打电话提出相同请求)。秘书不会在门口扔掉重复件——他们只是把所有文件放进一个文件夹。每天一次,他们翻阅文件夹,只保留每个客户的最新申请。如果在整理之前有人问“伊万诺夫有多少份申请?”——他们会看到两份。之后——一份。
2. ReplacingMergeTree 的工作原理——逐块解析
重复项源于不可靠的投递
ClickHouse 最初是为大规模分析设计的,偶尔的丢失或重复并不关键。但后来人们开始将其用于关键数据——然后遇到了麻烦。ReplacingMergeTree 是对这种痛苦的回应。
为什么会出现重复项?
- 客户端发送了数据,没有收到确认(超时),然后重新发送。
- 队列系统(Kafka、RabbitMQ)提供
at-least-once保证——至少一次投递,可能重复。 - ETL 过程中的错误(提取、转换、加载)——管道运行了两次。
机制:按 ORDER BY 键合并
创建带有 ReplacingMergeTree 的表时,必须指定一个排序键——ORDER BY (column1, column2)。这不是经典意义上的主键(如 PostgreSQL 中),而是一种在磁盘上物理排序数据的方式。ClickHouse 将数据存储在部分中——按此键排序的块。
当两个部分合并为一个时(一个称为合并的后台进程),ReplacingMergeTree 会扫描具有相同 ORDER BY 键值的行,并只保留一个。保留哪一个? 默认情况下——按插入时间最后一个。但你可以指定一个数字 version 列,然后保留版本值最大的行。
Git 类比: 合并时的 ReplacingMergeTree 行为类似于 Git 解决冲突时:对同一文件的两个更改,保留最新的(如果你没有明确指定策略)。只是这里的文件是表中的一行,键是 ORDER BY。
版本控制:ReplacingMergeTree(version) 如何改变规则
语法:ReplacingMergeTree(version_column)。如果 version_column 是整数(UInt* 或 DateTime),则保留最大值的行。这提供了手动控制:你可以明确指定哪个版本“获胜”。
示例:我们发送带有 updated_at = now() 的下注。重新发送时,updated_at 会稍大一些。合并将保留更新的那个。如果不指定 version,ClickHouse 会选择最后到达的那个——这可能不是按业务逻辑最新的,只是最后一次物理插入。区别很重要。
3. 使用 ReplacingMergeTree 创建表——详细解析
-- 创建用于去重的下注表
CREATE TABLE bets_dedup
(
user_id UInt64, -- 玩家 ID(谁的下注)
bet_id String, -- 唯一下注 ID(客户端生成)
amount Decimal(10,2), -- 金额(卢布)
created_at DateTime, -- 下注创建时间
updated_at DateTime -- 最后更新时间(用于版本)
)
ENGINE = ReplacingMergeTree(updated_at) -- 引擎,版本使用 updated_at
ORDER BY (user_id, bet_id) -- 去重键:(user_id, bet_id)
逐行解释:
ENGINE = ReplacingMergeTree(updated_at)—— 指定这是 ReplacingMergeTree,并且updated_at列将用作版本。合并时,对于具有相同ORDER BY的两行,保留updated_at较大(更新)的那一行。如果updated_at相等——则保留最后物理插入的那一行(但最好不要依赖于此)。ORDER BY (user_id, bet_id)—— 最重要的参数!这组列定义了什么是重复项。如果两行在ORDER BY中的所有列具有相同的值,则它们被视为重复项。这里:来自用户 user_id 且 bet_id 为某个值的下注是唯一的。如果两行具有user_id=123, bet_id='abc-456'到达——它们将合并为一行。
为什么是 ORDER BY 而不是 PRIMARY KEY? 在 ClickHouse 中,PRIMARY KEY 不必是唯一的。它是索引的提示,而 ORDER BY 是磁盘上的物理顺序。ReplacingMergeTree 依赖于 ORDER BY,即使 PRIMARY KEY 更短。如果不指定 PRIMARY KEY,则它与 ORDER BY 匹配。
如果 ORDER BY 太宽泛会怎样? 例如,包含 amount。那么两个金额不同的下注(即使 user_id, bet_id 相同)也不会被视为重复项——两者都会保留。去重将不起作用。陷阱 #1(我们最后会回到这一点)。
4. 为什么合并前 SELECT 可能返回重复项——以及如何应对
主要细微差别: ReplacingMergeTree 仅在数据部分合并期间删除重复项。这是一个后台进程,不会立即发生。在插入重复项和物理删除它们之间,可能需要几秒到几小时(取决于设置和负载)。
这在实践中意味着什么?
让我们插入两个重复项:
-- 第一次插入
INSERT INTO bets_dedup VALUES (123, 'bet-001', 1000, now(), now());
-- 5 秒后——第二次(服务器未收到确认并重新发送)
INSERT INTO bets_dedup VALUES (123, 'bet-001', 1000, now(), now() + interval 5 second);
现在运行常规的 SELECT * FROM bets_dedup WHERE user_id = 123。我们会看到什么?两行。 因为合并尚未发生。数据位于不同的部分。每个部分内部按 ORDER BY 排序,但重复项可能位于不同的部分。
如何保证只得到一行? 使用 FINAL:
SELECT * FROM bets_dedup FINAL WHERE user_id = 123;
FINAL 强制 ClickHouse 即时合并此查询的所有部分,应用 ReplacingMergeTree 逻辑。你将得到一行——具有最大 updated_at(如果没有版本,则是按插入时间最后的那一行)。
为什么 FINAL 很慢? ClickHouse 读取表的所有部分,在内存中按 ORDER BY 键排序,删除重复项,然后才返回结果。在大表(数十亿行)上,这可能需要几秒或几分钟。优化器无法有效使用索引——它必须扫描大量数据。
建议: 不要在大表上实时使用 FINAL。将其用于:
- 针对单个
user_id的点查询(索引仍然有帮助)。 - 时间不关键的后台任务(夜间报告)。
- 小表(最多几百万行)。
对于生产负载,有一个更好的模式——不带 FINAL 的物化视图。
5. FINAL 的性能——何时可接受,何时不可接受
何时 FINAL 可以:
- 表很小(每台服务器最多 1000–2000 万行)。
- 你通过索引查询单个用户(WHERE user_id = 特定值)。
- 你有一个每小时运行一次的后台聚合,等待 10 秒没问题。
- 每天导出一次数据用于报告。
何时 FINAL 是杀手:
- 表超过 1 亿行。
- 无过滤的查询(SELECT * FROM table FINAL)——ClickHouse 将读取所有内容。
- 高负载的类 OLTP 场景(每秒数十次带 FINAL 的查询)。
- 频繁更新相同键——许多部分累积,FINAL 读取所有部分。
类比: SELECT ... FINAL 就像手动翻阅档案中的所有文件以找到文档的最新版本,而不是查看专门的“当前版本日志”。它有效,但不适用于每个客户端请求。
如何检查查询是否使用了 FINAL?
ClickHouse 有 EXPLAIN 命令:
EXPLAIN SELECT * FROM bets_dedup FINAL WHERE user_id = 123;
查找带有 final 标志的 ReadFromMergeTree。如果看到它——查询正在老老实实地遍历所有部分。
6. 模式:通过物化视图实现无需 FINAL 的后台聚合
这是我最喜欢的绕过 FINAL 的方法。思路:让 ReplacingMergeTree 过自己的生活,重复项在后台逐渐合并。对于读取,我们创建一个物化视图,定期重建并包含已经“干净”的无重复数据。
实现方式:
-- 1. 基础表——脏数据,有重复项
CREATE TABLE bets_raw
(
user_id UInt64,
bet_id String,
amount Decimal(10,2),
created_at DateTime,
updated_at DateTime
)
ENGINE = ReplacingMergeTree(updated_at)
ORDER BY (user_id, bet_id);
-- 2. 目标表——干净数据,无重复项
CREATE TABLE bets_clean
(
user_id UInt64,
bet_id String,
amount Decimal(10,2),
created_at DateTime,
updated_at DateTime
)
ENGINE = MergeTree() -- 常规 MergeTree,无去重
ORDER BY (user_id, bet_id);
-- 3. 物化视图——插入时传输数据
CREATE MATERIALIZED VIEW bets_mv TO bets_clean AS
SELECT
user_id,
argMax(amount, updated_at) AS amount, -- 从 updated_at 最大的行中取 amount
argMax(created_at, updated_at) AS created_at,
max(updated_at) AS updated_at
FROM bets_raw
GROUP BY user_id, bet_id; -- 按去重键分组
关键点解释:
argMax(amount, updated_at)—— 一个聚合函数,返回具有最大updated_at的行中的amount值。如果我们有不同updated_at(和不同amount——例如下注金额改变)的重复项,则保留最新的金额。这类似于手动版本控制。GROUP BY user_id, bet_id—— 这里我们明确说:“将用户+下注 ID 的组合视为重复项。”现在无需等待合并——每次向bets_raw插入数据都会(几乎)立即通过bets_mv触发bets_clean中的重新计算。重要限制: ClickHouse 中的物化视图批量处理数据——每次插入单独处理。如果一次插入包含两个重复的
(user_id, bet_id)——它们会在批次内合并。如果重复项来自不同的插入——bets_clean可能包含临时重复项,直到bets_raw合并。为了完美清洁,你需要要么在从bets_raw读取时使用FINAL,要么定期运行OPTIMIZE TABLE bets_raw(强制合并)。
类比: 这就像有一个草稿(bets_raw),你放入所有更正,而秘书每 5 分钟重新打印一份干净的副本(bets_clean),没有错误。读者只看干净的副本——快速且无重复。
7. 使用单调递增版本的 ReplacingMergeTree(version)——更新语义
常规的 ReplacingMergeTree 只是保留“最后到达”的行。如果旧数据可能在新数据之后到达(例如由于网络延迟),这很糟糕。解决方案:使用单调递增的 version 列(例如时间戳或序列 ID)。
示例:玩家余额表,包含存款历史
CREATE TABLE player_balance
(
user_id UInt64,
transaction_id String, -- 唯一交易 ID (UUID)
amount Int64, -- 余额变化(可为负)
balance_after Int64, -- 交易后余额
event_time DateTime, -- 客户端事件时间
ingestion_time DateTime -- 插入 ClickHouse 的时间(版本)
)
ENGINE = ReplacingMergeTree(ingestion_time) -- 版本 = 插入时间
ORDER BY (user_id, transaction_id);
现在即使交易 tx-001 到达两次,但具有不同的 ingestion_time,后插入的那个(具有更大的 ingestion_time)将被保留。这可以防止“延迟重复”——当第一次插入在 12:00,第二次在 12:05(重复),但由于网络故障,第二次在服务器上先到达。如果没有版本,将保留较早的那个(按插入时间)——这可能是错误的那个。
“单调递增”是什么意思? 每次新插入时,ingestion_time 的值必须大于或等于之前的值。使用 now()(ClickHouse 服务器上的当前时间)或原子计数器(例如来自 ZooKeeper)。不要依赖客户端时间——时钟可能跳跃。
8. 完整示例:按 transaction_id 对余额充值去重
让我们把所有内容整合起来。我们有一个微服务,接受来自支付系统的余额充值。支付系统发送 webhook(HTTP 调用)——有时会重复。
-- 步骤 1:创建原始事件表
CREATE TABLE balance_events
(
user_id UInt64,
transaction_id String, -- 支付系统的唯一 ID
amount Int64, -- +1000 卢布
event_time DateTime, -- 用户扣款时间
inserted_at DateTime DEFAULT now() -- 插入时自动设置
)
ENGINE = ReplacingMergeTree(inserted_at)
ORDER BY (user_id, transaction_id); -- 按 (用户, 交易) 对去重
-- 步骤 2:插入数据(假设重复到达)
INSERT INTO balance_events (user_id, transaction_id, amount, event_time)
VALUES (1, 'pay_001', 1000, '2025-06-01 10:00:00');
-- 一分钟后,重复到达(inserted_at 将自动设置为 now() + 60 秒)
INSERT INTO balance_events (user_id, transaction_id, amount, event_time)
VALUES (1, 'pay_001', 1000, '2025-06-01 10:00:00');
-- 步骤 3:不带 FINAL 读取——我们会看到 2 行(但仅当它们尚未合并时)
SELECT * FROM balance_events WHERE user_id = 1;
-- 结果:两行,具有相同的 user_id, transaction_id, amount
-- 步骤 4:带 FINAL 读取——我们看到一行(具有最大 inserted_at)
SELECT * FROM balance_events FINAL WHERE user_id = 1;
-- 结果:一行
为什么 ORDER BY 中只有 transaction_id 不够? 因为两个不同的用户可能有相同的 transaction_id(例如每个支付系统有自己的计数器)。添加 user_id 保证了用户内的唯一性。如果系统生成全局 UUID(550e8400-e29b-41d4-a716-446655440000)——你可以单独使用 ORDER BY transaction_id,一个 UUID 就足够了。
9. 与 CollapsingMergeTree 的比较
CollapsingMergeTree 是另一个用于处理变更的引擎。它存储“加”和“减”对,并在合并时折叠它们。
主要区别:
| 特性 | ReplacingMergeTree | CollapsingMergeTree |
|---|---|---|
| 机制 | 从重复项中保留一行 | 折叠对(+1 和 -1) |
| 目的 | 插入去重 | 更新聚合(例如购物车) |
| 是否需要版本 | 可选(版本列) | 必须的 Sign 标志(+1/-1) |
| 能否存储历史 | 可以,直到合并前的所有版本 | 不能,对会被销毁 |
| 读取时使用 FINAL | 是,否则重复项可见 | 是,否则未折叠的对可见 |
何时选择 ReplacingMergeTree:
- 你只需要删除重复行。
- 你有自然的去重键(交易 ID)。
- 数据很少更改(主要是插入)。
何时选择 CollapsingMergeTree:
- 你频繁更新聚合指标(例如“购物车中的商品数量”)。
- 你只需要存储结果,而不是变更历史。
CollapsingMergeTree 示例:
CREATE TABLE cart_items
(
user_id UInt64,
product_id UInt64,
quantity Int16,
sign Int8 -- +1(添加),-1(移除)
) ENGINE = CollapsingMergeTree(sign)
ORDER BY (user_id, product_id);
使用 ReplacingMergeTree,你只需用新的 quantity 版本覆盖行——但你会丢失变更历史。CollapsingMergeTree 允许你计算总数(SUM(quantity * sign)),即使没有 FINAL。
10. 常见陷阱——以及如何避免
陷阱 #1:ORDER BY 未包含所有唯一字段
-- 错误:只使用 user_id
CREATE TABLE bets_bad ENGINE = ReplacingMergeTree ORDER BY user_id;
-- 为同一用户插入两个不同 bet_id 的下注
INSERT INTO bets_bad VALUES (1, 'bet_001', 100);
INSERT INTO bets_bad VALUES (1, 'bet_002', 200);
-- 合并时它们会合并为一行——因为 ORDER BY (user_id) 相同!
-- 丢失了 bet_002。
正确: 在 ORDER BY 中包含所有使行唯一的列——通常是代理 ID(transaction_id)或组合(user_id, bet_id)。
陷阱 #2:天真地期望即时去重
新手插入重复项后立即执行不带 FINAL 的 SELECT——看到重复项。他们对 ClickHouse 感到失望。记住: 去重是异步的。如果你需要即时一致性——使用 FINAL 或物化视图模式。
陷阱 #3:使用非单调的版本
-- 错误:版本是客户端时间
CREATE TABLE events ENGINE = ReplacingMergeTree(client_time) ORDER BY (id);
-- 客户端时钟落后,他们在新版本之后发送旧版本
-- 合并时,错误的(旧)行将被保留
解决方案: 使用 ClickHouse 端的 now() 或硬件计数器。
陷阱 #4:对大数据乐观使用 FINAL
我遇到过一个案例:一个开发者在 20 亿行表的所有报告中启用了 FINAL。查询开始超时,超过 300 秒。不得不重写为使用 GROUP BY 和 argMax 的聚合。
黄金法则: 如果你通过 FINAL 读取超过 10% 的表——你做错了。使用物化视图或重新思考架构。
陷阱 #5:没有 ORDER BY 的 ReplacingMergeTree
ClickHouse 不允许你创建没有 ORDER BY 的表。但你可以指定 ORDER BY tuple()(空元组)。那么表中的所有行都被视为重复项——第一次合并后只会保留一行。几乎从不需要。
下一步——相关文章链接
现在你已经掌握了 ReplacingMergeTree,以下是接下来要探索的主题:
如何优化合并——设置如
merge_with_ttl_timeout、number_of_free_entries_in_pool_to_lower_max_size_of_merge(听起来吓人但有用)。INSERT 级别的去重——带有 ZooKeeper 的
ReplicatedReplacingMergeTree引擎。这是另一个层次:重复项在插入时立即被截断,但代价是延迟和复杂性。替代方案:
VersionedCollapsingMergeTree——一个同时支持版本控制和折叠的混合体。物化视图详解——如何构建多级聚合以完全避免
FINAL。
最后: ReplacingMergeTree 是一个强大的工具,但它不是关于“立即删除重复项”。而是关于“数据最终会变得干净,在此期间你与之一起工作。”如果你需要严格的唯一性(如 PostgreSQL 中的 PRIMARY KEY)——ClickHouse 不是最佳选择。但对于 99% 的具有重复插入的分析任务——它是救星。
← 上一篇: ClickHouse 配置:我是如何搭建生产环境且不踩坑的
→ 下一篇: SummingMergeTree 与 AggregatingMergeTree:轻松实现增量聚合
— Editorial Team
暂无评论。