返回首页

ClickHouse中的ReplacingMergeTree:完整指南

本文详细介绍了ClickHouse中的ReplacingMergeTree引擎:为何需要它来应对不可靠交付中的重复项,如何通过ORDER BY进行合并,版本控制和FINAL修饰符的作用。讨论了典型陷阱、与CollapsingMergeTree的对比以及绕过FINAL的物化视图模式。

ReplacingMergeTree:ClickHouse中的数据去重
Advertisement 728x90

ReplacingMergeTree:如何在 ClickHouse 中轻松去重

1. 为什么需要 ReplacingMergeTree——现实中的重复问题

想象你正在开发一个在线赌场。玩家点击“下注”按钮——1000 卢布押黑色。此时,处理请求的服务器突然崩溃(过热、网络故障,谁知道呢)。客户端没有收到响应,心想:“下注没成功。”玩家再次点击。服务器恢复并接受了两个请求。数据库中——两条相同的下注记录。玩家愤怒了:被扣了 2000 卢布而不是 1000。

这是一个典型的幂等性问题(源自拉丁语 idem——相同,potens——能力)。如果一个操作重复执行的结果与执行一次相同,则该操作是幂等的。在数据库世界中,我们需要一种机制来判断:“我已经见过这个下注了,我会忽略第二个版本。”

在 ClickHouse 中,有 ReplacingMergeTree 用于此目的。它是一个表引擎,在数据部分合并期间自动删除重复项。但我提前警告你:这不是魔法——它有我们即将讨论的怪癖。

Google AdInline article slot

现实类比: ReplacingMergeTree 就像一位记录会议记录的秘书。人们带着请求来找你。有时同一个客户带来两份相同的申请(例如,错过了火车并请求退款,然后又打电话提出相同请求)。秘书不会在门口扔掉重复件——他们只是把所有文件放进一个文件夹。每天一次,他们翻阅文件夹,只保留每个客户的最新申请。如果在整理之前有人问“伊万诺夫有多少份申请?”——他们会看到两份。之后——一份。

2. ReplacingMergeTree 的工作原理——逐块解析

重复项源于不可靠的投递

ClickHouse 最初是为大规模分析设计的,偶尔的丢失或重复并不关键。但后来人们开始将其用于关键数据——然后遇到了麻烦。ReplacingMergeTree 是对这种痛苦的回应。

为什么会出现重复项?

Google AdInline article slot
  • 客户端发送了数据,没有收到确认(超时),然后重新发送。
  • 队列系统(Kafka、RabbitMQ)提供 at-least-once 保证——至少一次投递,可能重复。
  • ETL 过程中的错误(提取、转换、加载)——管道运行了两次。

机制:按 ORDER BY 键合并

创建带有 ReplacingMergeTree 的表时,必须指定一个排序键——ORDER BY (column1, column2)。这不是经典意义上的主键(如 PostgreSQL 中),而是一种在磁盘上物理排序数据的方式。ClickHouse 将数据存储在部分中——按此键排序的块。

当两个部分合并为一个时(一个称为合并的后台进程),ReplacingMergeTree 会扫描具有相同 ORDER BY 键值的行,并只保留一个。保留哪一个? 默认情况下——按插入时间最后一个。但你可以指定一个数字 version 列,然后保留版本值最大的行。

Git 类比: 合并时的 ReplacingMergeTree 行为类似于 Git 解决冲突时:对同一文件的两个更改,保留最新的(如果你没有明确指定策略)。只是这里的文件是表中的一行,键是 ORDER BY

Google AdInline article slot

版本控制: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:天真地期望即时去重

新手插入重复项后立即执行不带 FINALSELECT——看到重复项。他们对 ClickHouse 感到失望。记住: 去重是异步的。如果你需要即时一致性——使用 FINAL 或物化视图模式。

陷阱 #3:使用非单调的版本

-- 错误:版本是客户端时间
CREATE TABLE events ENGINE = ReplacingMergeTree(client_time) ORDER BY (id);

-- 客户端时钟落后,他们在新版本之后发送旧版本
-- 合并时,错误的(旧)行将被保留

解决方案: 使用 ClickHouse 端的 now() 或硬件计数器。

陷阱 #4:对大数据乐观使用 FINAL

我遇到过一个案例:一个开发者在 20 亿行表的所有报告中启用了 FINAL。查询开始超时,超过 300 秒。不得不重写为使用 GROUP BYargMax 的聚合。

黄金法则: 如果你通过 FINAL 读取超过 10% 的表——你做错了。使用物化视图或重新思考架构。

陷阱 #5:没有 ORDER BY 的 ReplacingMergeTree

ClickHouse 不允许你创建没有 ORDER BY 的表。但你可以指定 ORDER BY tuple()(空元组)。那么表中的所有行都被视为重复项——第一次合并后只会保留一行。几乎从不需要。

下一步——相关文章链接

现在你已经掌握了 ReplacingMergeTree,以下是接下来要探索的主题:

  1. 如何优化合并——设置如 merge_with_ttl_timeoutnumber_of_free_entries_in_pool_to_lower_max_size_of_merge(听起来吓人但有用)。

  2. INSERT 级别的去重——带有 ZooKeeper 的 ReplicatedReplacingMergeTree 引擎。这是另一个层次:重复项在插入时立即被截断,但代价是延迟和复杂性。

  3. 替代方案:VersionedCollapsingMergeTree——一个同时支持版本控制和折叠的混合体。

  4. 物化视图详解——如何构建多级聚合以完全避免 FINAL

最后: ReplacingMergeTree 是一个强大的工具,但它不是关于“立即删除重复项”。而是关于“数据最终会变得干净,在此期间你与之一起工作。”如果你需要严格的唯一性(如 PostgreSQL 中的 PRIMARY KEY)——ClickHouse 不是最佳选择。但对于 99% 的具有重复插入的分析任务——它是救星。


上一篇:
下一篇: SummingMergeTree 与 AggregatingMergeTree:轻松实现增量聚合

— Editorial Team

Advertisement 728x90

继续阅读