同一条「业务服务 → Kafka → Kafka Connect sink → ClickHouse」链路又出了一次重复行,这次的入口在生产端:broker 有一段时间反复返回 NOT_LEADER_OR_FOLLOWER,生产者重试时把一部分记录在 topic 里存成了两份,sink 照常把两份都写进了 ClickHouse。一个一千多万行的日分区里,45 分钟的窗口多出一百多行,每个键正好两份。

这是这条链路上第四种重复入口。清理手段对它们都一样(REPLACE PARTITION runbook),能不能从源头挡住却各不相同。这篇把四个入口放在一起,再拿 sink 自带的 exactlyOnce 逐个对一遍。

重复从哪个入口进来#

入口 怎么产生 见过的情形
生产端重试 请求其实写进了 broker,应答在抖动里丢了,重试又写一份,topic 里是两条 offset 不同的记录 本篇这次
sink 插入超时重投 INSERT 超过 client 超时但服务端已经提交,框架把同一批原样再交一次 Kafka Connect 超时重投那篇里的 234 行
sink task 重启后重放 task 被重启前没来得及提交 offset,接替的 task 从上一个提交点重读 另一次托管 connector 的 task 集体重启,重放出几万行
人工补发 发送超时的记录先存进重试表,事后补发时其中一部分其实早已写进去 这次补发前先按键查了 ClickHouse,没有补出新的重复

服务端的块级去重对这几种基本帮不上:重放批次的边界和原批不同,生产端的两份本来就是两条记录,只有超时重投时批次原样不变,可它又输在窗口太短(ClickHouse 的块级去重窗口)。

开着幂等的 producer 也会写出两份#

生产者用的 kafka-clients 是 3.4.0,配置里没有覆盖 acks 和 enable.idempotence。按 Confluent 的 producer 配置参考,没有冲突配置时幂等默认打开,delivery.timeout.ms 默认 120000。重复还是出现了。

AutoMQ 对 producer 的整理(2025-02)写到了这个边界:幂等只在一个 producer 会话内成立,超过 delivery.timeout.ms 报失败时,这条消息到底写没写进 broker 是不确定的。这次的样本记录正是在发出 120 秒后报发送失败、被转进重试表的。broker 具体怎么存下了两份,没拿到当时的客户端日志之前只能推断,好在不影响这里的结论:topic 里的数据要按至少一次对待。

exactlyOnce 在 sink 里怎么判断#

ClickHouse 官方的 Kafka Connect sink 有一个 exactlyOnce 模式。按设计文档,每个 topic-partition 在一张 KeeperMap 表里记下上一批的 offset 区间和一个状态:插入前写 BEFORE_PROCESSING,ClickHouse 确认之后写 AFTER_PROCESSING,新来的一批先和记录比一下再决定插不插。插件主干 Processing.java 里按状态分两支:

  • 状态是 AFTER_PROCESSING 时,区间相同或者被记录区间包含的批次直接跳过;整批都早于记录区间(重放常见的情况)默认抛异常让 task 停下,打开 tolerateStateMismatch 才改成跳过。
  • 状态是 BEFORE_PROCESSING,也就是上一次插入结果不明时,同一批原样重插,代码注释是「Dedupe in clickhouse will fix it」,交给服务端的块级去重窗口。

这个模式在托管服务上能不能用,之前一直没确认,两边的文档现在都写了。Confluent Cloud 全托管 ClickHouse Sink 的配置页列出了 exactlyOnce(默认 false)、keeperOnCluster 和 tolerateStateMismatch,Aiven 的支持引擎列表里 KeeperMap 在 25.3、25.8、26.3 上都可用。另外按 ClickHouse 的 sink 文档,开了它就不能用 buffering。这条链路上的 sink 目前都没开。

四个入口逐条对照#

入口 exactlyOnce 能不能挡 剩下的靠什么
sink task 重启后重放 能,落在已确认区间里的跳过,整批更早时默认停下,不会写出第二份 无
sink 插入超时重投 部分,按 BEFORE_PROCESSING 原样重插,仍要靠去重窗口,那张表的窗口折合 8 秒 落库后按分区清理
生产端重试 挡不住,两份是 offset 不同的两条记录,对 connector 来说都是新数据 补发前按键查,落库后按分区清理
人工补发 挡不住,理由同上 补发前按键查

打开它能把量最大的重放挡在门外,落库后的清理仍然要留着。这个模式本身我还没实测:lab 里没有 Kafka 和 Connect,它在 Keeper 抖动时是停住还是出别的问题、吞吐会降多少,要在测试环境里重启 task、掐连接试过才知道。它把状态存在 KeeperMap 里,而 Kafka Connect 超时重投那次的起点正是 Keeper 抖动,评估时这一点要算进去。

其他团队把去重放在哪一层#

2025 到 2026 年能查到的几篇公开做法,按去重放在哪一层分开看。

在数据进 ClickHouse 之前按键去掉重复,GlassFlow(2026-01)和它与 Altinity 合写的一篇(2025-05)走的是这条,能连生产端的重复一起挡掉,代价是 topic 和 ClickHouse 之间要多一个有状态的组件。Fiddler AI(2025-08)把 API 网关、应用、插入 token 和 ReplacingMergeTree 四层叠在一起。

ReplacingMergeTree 加 FINAL 这条路,Yandex Cloud 的方案(2025-03)把需要 FINAL 的部分限制在最近一段数据上,文中近期数据带 FINAL 的扫描速度约为老数据的十五分之一。对这条链路来说它要换引擎、改所有读方,也救不了下游按写入时间增量累加的汇总表(清理那篇最后一节)。

按已封口的分区定期重写,和 Altinity 知识库里「isolated dirty data fragments」的思路接近,落到这里就是把 runbook 做成定时作业。几篇的共同点是没有指望某一层单独兜住。

小结#

  • 这条链路上的重复有四个入口:生产端重试、sink 插入超时重投、task 重启后重放、人工补发。清理手段相同,能从源头挡住的不同。
  • 生产端开着幂等也会出重复:超过 delivery.timeout.ms 报失败的那条,写没写进 broker 不确定,topic 要按至少一次对待。
  • sink 的 exactlyOnce 按 offset 区间和 KeeperMap 里的状态判断:能挡住重放;超时重投时状态停在 BEFORE_PROCESSING,照样重插,还得靠块级去重窗口;生产端和人工补发的两份 offset 不同,挡不住。
  • Confluent 全托管 sink 和 Aiven 的 KeeperMap 都支持这个模式,这条链路没开,也还没实测;它的状态和块级去重、复制元数据共用同一个 Keeper。

相关文章#

这几篇都在同一个 Aiven 托管的 ClickHouse 集群上。(一)是一次补数方案评审,其余六篇是重复行:(二)(三)(五)讲重复从哪来、为什么没被挡住,(四)(六)(七)讲怎么清。

  1. 用只读权限评审 ClickHouse 补数方案 — 同一个集群上的 mutation 评审:列级重写的成本怎么算、回滚为什么要按 id 名单而不是按状态
  2. ClickHouse 里的重复行来自 Kafka Connect 超时重投 — 重复行怎么产生的:30 秒超时是谁的默认值、框架什么时候把同一批再投一次
  3. ClickHouse 的块级去重窗口 — 服务端那层为什么没兜住:窗口按块数算,这张表的建块速率折合 8 秒
  4. 清理 ClickHouse 重复行的 REPLACE PARTITION runbook — 怎么低成本数出有多少、怎么清掉:临时表加 REPLACE PARTITION,以及跑之前要验的四件事
  5. ClickHouse 重复行的四个入口和 sink exactlyOnce 的覆盖范围(本篇)
  6. 清理 ClickHouse 重复行时行数校验拦不住的情形 — 按位置对列、argMin 跳过 NULL、迟到写入和落后副本,以及快照和 part_log 怎么补上
  7. ClickHouse 轻量删除清重复行的代价 — 带子查询的 DELETE 在复制表上按 part 数重复执行,IN PARTITION 和同秒两份的问题

参考资料#