清理 ClickHouse 重复行的 REPLACE PARTITION runbook
目录
上一次在这个集群上评审补数方案,结论里专门拦下了 REPLACE PARTITION:它会把做快照到换分区之间插进生产表的行一起覆盖掉,而那张表的老分区上确实还在进迟到行(用只读权限评审 ClickHouse 补数方案)。这次反过来用它,因为条件变了:要删的是 234 行内容逐字节相同的重复数据,来自 Kafka Connect sink 在 INSERT 超时后重投同一批(ClickHouse 里的重复行来自 Kafka Connect 超时重投),而且重复会再来、规模事前不知道,要的是一个能反复跑的 runbook。分区旧了并不代表快照风险消失,它变成下面要先验的一条。
这套 runbook 在生产上还没有执行。下面每一步都在本地一套同版本的三副本集群上跑过,实测结果在文末。
只删这一次的 234 行,lightweight delete 更便宜:重写的只是 _row_exists 那一列掩码,不碰业务列的文件,不用把 30GiB 从 S3 拉回来,也没有快照窗口的问题,代价是掩码要一直参与后续查询。按分区重写是一次性的,读侧没有额外开销,也不用改表结构,换来的是整个分区读一遍写一遍,而这个分区已经沉到 S3 层,30GiB 要拉回来,新 part 还会落到热盘上。多余行多到掩码不划算时它才占优。换 ReplacingMergeTree 这条路没走,理由在那次评审里写全了(合并异步、FINAL 是永久读开销),对账口径的表不合适。
低成本地数出有多少重复#
清理之前要先知道有多少、集中在哪几个小时。直接对整个分区 GROUP BY (id, version) HAVING count() > 1 也能数,但这张表单个日分区两亿行、30GiB,七周前的分区已经沉到 S3 层,在生产库上这么扫一次代价太大。
换一个省的办法。重复只需要一个数字就能看出来:
SELECT count() AS rows,
uniqExact(id, version) AS keys
FROM events
WHERE _partition_id = '20260730'
AND settle_time >= {hour_start_ms}
AND settle_time < {hour_start_ms} + 3600000
count() - uniqExact(...) 就是这一小时的多余行数。省在 WHERE 的两层裁剪:_partition_id 裁掉分区,settle_time 是排序键的第一列,按小时切能让主键索引生效。一小时 500 万到 1600 万行,几百 MB 内存,单线程跑得动。24 条这样的查询覆盖一整天,哪一小时有问题也直接看出来。
跑生产库的时候把限制写在 URL 参数里,服务端强制,比在客户端约束可靠:
curl "https://host:port/?readonly=2&max_threads=1&max_memory_usage=1000000000&priority=10" \
--data-binary @query.sql
先把重复的键单独存一张表#
runbook 只碰重复的那几百组,第一步是把它们的键固化下来。上一节已经数出重复集中在哪几个小时,只在那几个小时里找,不扫整个分区:
-- 只在数出有重复的那几个小时里找,不扫整个分区
DROP TABLE IF EXISTS events_dedup_keys_20260730 SYNC;
CREATE TABLE events_dedup_keys_20260730
ENGINE = MergeTree ORDER BY (id, version) AS
SELECT id, version
FROM events
WHERE _partition_id = '20260730'
AND settle_time >= {t0_ms}
AND settle_time < {t0_ms} + 3 * 3600000
GROUP BY id, version
HAVING count() > 1;
这里不用 Memory 引擎也不用临时表。Aiven 的服务架构是单个 URL 指向全部三个节点、「connections going randomly to any of the servers」,后面几条语句未必落在同一个节点上执行,这张表得在三个节点上都读得到。写成 MergeTree 就行,Aiven 会自动改写成 ReplicatedMergeTree。
有了这张表,后面取详情用 (id, version) IN (SELECT id, version FROM events_dedup_keys_20260730) 就只读这几百组的完整内容,整个分区算 cityHash64(*) 没有必要。
runbook 本体#
-- runbook 开头和结尾各做一次,保证重跑幂等
DROP TABLE IF EXISTS events_dedup_tmp_20260730 SYNC;
CREATE TABLE events_dedup_tmp_20260730 AS events;
-- 没有重复的行原样拷贝
INSERT INTO events_dedup_tmp_20260730
SELECT * FROM events
WHERE _partition_id = '20260730'
AND (id, version) NOT IN (SELECT id, version FROM events_dedup_keys_20260730);
-- 有重复的那几百组,每组只留一行
INSERT INTO events_dedup_tmp_20260730
SELECT * FROM events
WHERE _partition_id = '20260730'
AND (id, version) IN (SELECT id, version FROM events_dedup_keys_20260730)
LIMIT 1 BY id, version;
ALTER TABLE events REPLACE PARTITION ID '20260730' FROM events_dedup_tmp_20260730;
分区操作文档对 REPLACE PARTITION 的原话是「The operation is atomic」。重复执行是幂等的,读侧没有额外开销,也不用改引擎或换排序键。文档只说它 copies the data partition,没说底层是否复用文件。同一页提到 hardlink 的只有 FREEZE PARTITION 那条(备份进 shadow/N/),REPLACE 这条没有。这里到底要不要真搬一遍 30GiB,没有实测过。
拆成两条 INSERT 是为了省内存。 直接对整个分区 LIMIT 1 BY id, version 要在内存里存下两亿个键。拆开之后,LIMIT 1 BY 只作用在几百个重复键上,剩下的行走普通拷贝。
用 DROP TABLE IF EXISTS … SYNC 而不是 CREATE TABLE IF NOT EXISTS。 runbook 如果在 CREATE 和 REPLACE 之间挂掉,下次重跑会撞见残留的临时表,IF NOT EXISTS 会复用上次写了一半的内容,两条 INSERT 追加上去就错了。SYNC 也要加:延迟删除是 Atomic 库引擎的性质,文档里 DROP 只是把元数据挪进 metadata_dropped/,真正的删除推迟到 database_atomic_delay_before_drop_table_sec。Replicated 库建在这套语义之上,不加 SYNC 紧接着 CREATE 同名表会撞车。
两张表不能指到同一个 Keeper 路径#
同一页文档对 REPLACE PARTITION 的要求是结构、partition key、ORDER BY、主键、storage policy 全一致,目标表还要包含源表所有的 index 和 projection。CREATE TABLE … AS events 复制表定义能满足这些,但 ReplicatedMergeTree 的 Keeper 路径也跟着表定义走。原表建表时如果把路径写成字面量而不是 {uuid} / {database} / {table} 这类 macro,新表会指到同一个路径上。本地试下来这时候 CREATE TABLE … AS 当场就报 REPLICA_ALREADY_EXISTS,不会悄悄共用一套复制元数据;路径里带 {uuid} 的那种,克隆拿到自己的 uuid,两个路径不同,也不共用。所以查这一步的价值在于动手建表之前就知道会不会撞,不是防静默损坏。
Aiven 上查不到建表语句:SHOW CREATE TABLE 对 avnadmin 是拒绝的,system.tables.engine_full 也被抹掉。改查 system.replicas,两个值必须不同:
SELECT table, zookeeper_path
FROM system.replicas
WHERE database = 'default'
AND table IN ('events', 'events_dedup_tmp_20260730')
分区 ID 要带引号,分区表达式不要#
写 PARTITION ID '20260730',别写 PARTITION '20260730'。分区键是 toYYYYMMDD(settle_time),分区表达式的值是个整数。同一页文档有两条相关规则:引号要看分区表达式的类型(「for the String type, you have to specify its name in quotes … For the Date and Int* types no quotes are needed」),而用分区 ID 时另有要求(「The partition ID must be specified in the PARTITION ID clause, in a single quotes」)。这里按第二条写。第一条说的是 Date 和 Int* 不「需要」引号,不是不许加:本地四种写法试下来,只有 PARTITION ID 20260730(ID 不加引号)被拒,报 Expected one of: string literal,PARTITION '20260730' 照样解析到那个分区。选 PARTITION ID '20260730' 是因为它的规则只有一条、不用再判断表达式类型。前面查询里用的 _partition_id 本来就是字符串形式的 '20260730',PARTITION ID 跟它是同一个口径。
换分区之前先确认三个副本同步#
连接随机落到任一节点,没有主从之分,所以 count() 可能被一个还没收到最新 part 的节点回答。而这批重复行本来就是复制滞后期间写进来的,在滞后期间换分区等于拿旧数据覆盖新数据。那次评审里 mutation 的验收也是同一个理由,查法一样,用 clusterAllReplicas 问遍三台:
SELECT max(queue_size), max(absolute_delay)
FROM clusterAllReplicas('default', system.replicas)
WHERE database = 'default' AND table = 'events'
分区必须是静止的#
从建临时表到 REPLACE PARTITION 之间插进生产表的行不在临时表里,换完就没了。这正是那次评审拦下 REPLACE PARTITION 的理由,这次要靠前置校验把它挡在门外。分区是按 settle_time 分的,迟到事件照样会进七周前的分区,不能靠分区够老来推定。
runbook 开头和 REPLACE 之前各数一次这个分区的行数,两次相等才继续,对不上就中止重跑。后台 merge 不改变行数,变了就是真有新行进来。这个 count() 同样要走 clusterAllReplicas。
跑之前要能拒绝执行#
这个 runbook 只处理 sink 重投产生的、内容完全相同的重复行。动手前先验两件事:多余行占比不超过一个阈值(定在 0.1%,超过说明不是重投事故,该人来判断),以及每一组重复的业务列哈希数是 1。第二条更要紧,LIMIT 1 BY 任选一行,如果两份内容其实不同,它会不报错地丢掉好的那份。
SELECT count() FROM (
SELECT id, version,
uniqExact(cityHash64(amount, status, settle_time, external_id)) AS h
FROM events
WHERE _partition_id = '20260730'
AND (id, version) IN (SELECT id, version FROM events_dedup_keys_20260730)
GROUP BY id, version HAVING h > 1
)
这个数要是 0,不是 0 就中止并报警。
在本地三副本上跑过一遍#
生产上这套还没执行,但上面每一步都在本地验过。用 docker compose 起一套同版本的 ClickHouse(25.3.14.14,生产是 25.3.14.1),1 shard × 3 replicas 配一个 Keeper,照着 Aiven 那个单 shard 三节点的拓扑搭。
造的数据按生产那次的占比来:10000 行正常数据加 5 组逐字节相同的重复,多余行占 0.05%,和两亿行里多 234 行是同一个量级。然后把上面的 SQL 原样跑一遍:
| 步骤 | 结果 |
|---|---|
count() - uniqExact(id, version) |
数出 5 行多余 |
| 键表 | 5 组重复键 |
| 前置校验一 | 占比 0.05%,低于 0.1% 那条线 |
| 前置校验二 | 业务列哈希数大于 1 的组数是 0 |
| 五条 SQL 跑完 | 10000 行,0 组重复,唯一键数仍是 10000 |
| 原样再跑一遍 | 还是 10000 行 |
三个副本各自 count() |
ch1、ch2、ch3 都是 10000 |
CREATE TABLE … AS events 这一步能成功,前提是原表路径里带 {uuid},这跟上面那节说的是一件事。
没覆盖到的有四样,跑生产之前心里要有数:
- 规模。 本地一万行,生产是单个日分区两亿行、30GiB。整个分区读一遍写一遍的耗时,本地量不出来。
- S3。 生产那个分区已经沉到 remote 层,30GiB 要拉回来,新 part 还会落到热盘上。本地没挂对象存储。
- 托管层。
SHOW CREATE TABLE被拒、连接随机落到任一节点,这两条是 Aiven 的行为,本地那套没有。 - 分区静止那条校验的失败路径。 本地没造过「跑到一半有新行进来」,所以只验了它通过的时候不挡路,没验它该中止的时候真会中止。
小结#
- 数重复用
count() - uniqExact(key),按分区加排序键首列切成小时窗,比整分区GROUP BY便宜得多,也能定位到出问题的时段。 - 补数那次评审拦下
REPLACE PARTITION,是因为那张表的老分区还在进迟到行;这次把同一条理由改成前置校验:开头和 REPLACE 之前各数一次分区行数,两次相等才继续。 CREATE TABLE … AS复制的是整份表定义,ReplicatedMergeTree的 Keeper 路径也在里面。路径写死成字面量的表,克隆会撞在同一个路径上,CREATE当场报REPLICA_ALREADY_EXISTS。Aiven 上看不到建表语句,改查system.replicas.zookeeper_path,动手之前就能知道。PARTITION ID '20260730'要带单引号,PARTITION后面跟分区表达式时按表达式类型决定加不加,两条是不同的规则。- 临时表用
DROP … SYNC开头,不用CREATE IF NOT EXISTS:Atomic库的 DROP 是延迟的,重跑时残留的半成品会被当成可复用的表。 LIMIT 1 BY只该作用在重复键上,没有重复的行走普通拷贝,否则要在内存里存下整个分区的键。它任选一行,所以跑之前要先验每组重复的业务列哈希数是 1。- 一次性清几百行用 lightweight delete 更便宜,这个 runbook 要的是可重跑,代价是整个分区读一遍写一遍。
- 整套在本地一套 25.3 的三副本集群上跑通了:重复清掉、行数对得上、重跑幂等、三个副本都换到。没验的是规模、S3 那一层、托管层的行为,以及分区静止校验的失败路径。
相关文章#
这四篇都在同一个 Aiven 托管的 ClickHouse 集群上。(一)是一次补数方案评审,(二)到(四)是一次重复行事故,从定位到清理。
- 用只读权限评审 ClickHouse 补数方案 — 同一个集群上的 mutation 评审:列级重写的成本怎么算、回滚为什么要按 id 名单而不是按状态
- ClickHouse 里的重复行来自 Kafka Connect 超时重投 — 重复行怎么产生的:30 秒超时是谁的默认值、框架什么时候把同一批再投一次
- ClickHouse 的块级去重窗口 — 服务端那层为什么没兜住:窗口按块数算,这张表的建块速率折合 8 秒
- 清理 ClickHouse 重复行的 REPLACE PARTITION runbook(本篇)
参考资料#
- ALTER TABLE … PARTITION —
REPLACE PARTITION的原子性、两张表要满足的条件、分区表达式什么时候要加引号、PARTITION ID要带单引号;hardlink 只出现在FREEZE PARTITION那条(shadow 备份目录),EXCHANGE TABLES不在这一页 - Atomic 库引擎 — DROP 为什么是延迟的,以及
SYNC修饰符 - cluster / clusterAllReplicas — 一条查询问遍所有副本
- Lightweight DELETE —
_row_exists掩码,以及查询时被注入的PREWHERE _row_exists - Aiven for ClickHouse 服务架构 — 单 shard 三节点、无主从、连接随机落到任一节点,
MergeTree自动改写成ReplicatedMergeTree