ClickHouse 里的重复行来自 Kafka Connect 超时重投
目录
背景#
一条很常见的链路:业务服务往 Kafka 写事件,ClickHouse 官方的 Kafka Connect Sink Connector(跑在 Confluent 上)消费这个 topic,批量 INSERT 进 ClickHouse 的一张日分区表。ClickHouse 是 Aiven 托管的,1 shard × 3 replicas。表是 ReplicatedMergeTree,排序键 (settle_time, id, version),按 toYYYYMMDD(settle_time) 分区,没有 TTL,历史分区一直留着。
Kafka Connect 是 at-least-once 的。sink 的 put() 抛出 RetriableException,框架不会把这批记录丢掉,也不会回 Kafka 重新拉一份:它暂停所有分区,把手里那批原样再交给 put() 一次,直到成功。下游写入不是幂等的,就会写第二份。这条链路上两个条件都满足:插件自己没有一层兜底(下面会核到源码),表也没有唯一约束。
现象#
对账时发现,同一个业务键 (id, version) 在表里有两行。两行的所有业务列逐字节相同:金额、状态、时间戳、外部单号全都一样。
先确认它不是上游双发。最先发现的那几组,上游日志里请求只有一次,ClickHouse 里却是两行,第二行是写入侧多写的。234 组没有逐笔回查上游日志,撑住这个判断的是间隔只集中在四档,以及下一节的源码。
再看两行的写入时间差:
create_time 间隔 |
组数 |
|---|---|
| 32 秒 | 39 |
| 60 秒 | 13 |
| 90 秒 | 158 |
| 124 秒 | 24 |
间隔集中在四个离散值上,每个值下面的行多半来自同一次重投,所以这里其实只有四个样本。从 ClickHouse 侧的 http_user_agent 认出插件版本是 clickhouse-kafka-connect/1.3.9,它依赖的 clickhouse-java 是 0.9.5。四个值都在 30 秒以上、看着像 30 的整数倍,30 秒是谁的、间隔为什么是这几档,得到源码里找。
重投是谁发起的,30 秒是谁的#
INSERT 外面没有重试循环。 插件配置里那个 retryCount(默认 3)不作用在 INSERT 上。v1.3.9 里它只喂给 ClickHouseHelperClient 的四个 while (retryCount < retry) 循环(:189、:205、:268、:294),那几个循环所在的方法是 ping 和 query,日志文案也写着 Ping retry %d out of %d、Query retry %d out of %d。写数据的路径在 ClickHouseWriter:clientVersion 默认是 V1(ClickHouseSinkConfig :291),走 doInsert*V1,一次请求加一个 future.get()(:1086、:1172、:1337),外面一层循环都没有;clientVersion=V2 那条路是三处 client.insert(...).get()(:1033、:1248、:1412),同样裸。超时后由 Utils.handleException 把 SocketTimeoutException 包成 RetriableException(:105)交给框架。所以重发这件事只有一种形态:框架层重投。
这个 30 秒来自 clickhouse-java 的默认值。 插件那个 timeoutSeconds(默认 30)在源码里只喂给 ping()(ClickHouseHelperClient :190、:206),建 V1 节点时(:101–133)传给 client 的选项只有用户名、密码和 SNI,没有 socket 超时。于是生效的是 clickhouse-java 0.9.5 V1 client 自己的默认值:ClickHouseClientOption.SOCKET_TIMEOUT 取 ClickHouseDataConfig.DEFAULT_TIMEOUT,即 30 * 1000 毫秒。两个 30 碰巧相等,很容易把功劳记到配置项头上。V2 client 在 0.9.5 里默认 socket_timeout=0,也就是不设应答超时。在 1.3.9 上切到 V2,挂住的 INSERT 就没有一个确定的返回时间了。上游后来在 PR #756(v1.3.10 起)给 V2 路径加了 clickhouseClientInsertTimeoutMs(默认 240 秒),套在 future.get() 外面,只管 V2。
想改这 30 秒,插件配置里唯一能碰到它的键是 jdbcConnectionProperties:它被拼进 V1 的 URL query,ClickHouseNode.of(url, options) 会把 socket_timeout=90000 这类参数解析成 client 选项。Confluent 全托管版也暴露了这个键,那一页 Connection details 下列着 jdbcConnectionProperties(Type: string,默认空,值要以 ? 开头、参数之间用 & 连),所以这条路在全托管上也是通的。
重投要等到下一个 offset 提交点。 WorkerSinkTask(Kafka 3.7)里 put() 抛 RetriableException 后:pausedForRedelivery = true、pauseAll(),messageBatch 不清(只有成功之后 :612 才清),下一轮循环 put(new ArrayList<>(messageBatch))(:601)把同一批再交出去(:620–628,注释原话是「the batch will be reprocessed on the next loop」)。但下一轮循环先做的是 poll(nextCommit - now)(:249):分区全部暂停、没有新数据,这个 poll 要等满才返回,等的时间就是到下一次 offset 提交点的距离,最长一个 offset.flush.interval.ms,默认 60 秒。sink 可以用 context.timeout() 让框架早点回来(:338–341),这个插件没调。
所以每失败一轮的代价是 30 秒超时加 0 到 60 秒的等待,一轮之内两份的间隔落在 (30, 90] 秒里:32、60、90 都是一轮失败;124 秒超出这个范围,至少两轮。把 32 / 60 / 90 / 124 读成一到四轮超时是不对的,那是把重投当成了紧接着发生的事。
框架重投的是保留在内存里的同一批。 因为 messageBatch 被保留,重投出去的记录、顺序、边界和第一次一模一样,插入块逐字节相同。再加一层:插件给 INSERT 带 insert_deduplication_token,值是 topic-partition-minOffset-maxOffset(QueryIdentifier :56–60;V1 路径在 getMutationRequest 里设,:1449,V2 路径 :1002)。批次挂在某个 partition 上时这个 token 一定拼得出来,getDeduplicationToken() 只在 partition == -1 时返回 null;重投的这几批是普通消费批次,token 相同。ClickHouse 服务端凭 token(没有 token 时凭块哈希)认重复,这两样重投前后都没变。去重失败就不会是「认不出来」,只剩「记不住了」这一种,也就是服务端那个去重窗口太短(窗口有多长、为什么没调大,写在ClickHouse 的块级去重窗口)。批次边界会变、哈希不稳那类论证在这条路径上用不上:边界只在 task 重启或 rebalance 之后才会变,这次没有这方面的证据。
每组都是 2 行,没有 3 行 4 行的组。 32 / 60 / 90 那三批只失败了一轮,两次尝试两份落地,没有多余的轮次要解释。124 那批至少三次尝试,中间那次没留下行:要么它压根没走到提交,要么它和第一次一样迟到提交、但提交时第一份的 token 还在窗口里,被去重掉了。事故期的 Connect 日志只留 7 天、query_log 只留 4 天,现在分不出来,也不影响结论。
errors.retry.timeout 与此无关。 Confluent 托管版暴露的 errors.retry.timeout(默认 30000ms,可填到 300000ms),那一页的原话是「Specifies the retry budget, in milliseconds, for failed record inserts」,字面读像在管 INSERT。它属于 errors.* 那套容错,落到 WorkerSinkTask 里是 retryWithToleranceOperator 包住 key / value / header 的转换(:533–541),put() 抛出的 RetriableException 走不到它,按文档超时之后是进 DLQ 或者让 task 失败。所以 124 秒的间隔并不说明这个值被调高过:put() 层面的重投本来就没有上限,报障单里也确实没记过任何调整。
排除 merge 背压#
sink 侧的症状是 Read timed out,插件对以 Read timed out after 开头的错误消息也转成 RetriableException(Utils :113–115),报障单里记了 542 条,这是当时提单人的一手观察;事故期的 Connect 日志和 Confluent 指标只留 7 天,现在都取不到了。
这个报错最容易联想到 part 太多、merge 排不过来,ClickHouse 对插入施加背压。这个方向查两个现成的指标就能否掉(下面这些 clickhouse_* 名字来自我这套部署的 Prometheus,Aiven 导出的那一套,换部署前缀未必相同):
clickhouse_metrics_delayed_inserts:被节流的插入数clickhouse_events_rejected_inserts:因 too many parts 被拒的插入数
把事故窗口前后逐小时拉一遍,两个指标全程是 0,说明慢的原因不在背压。
持续变慢的那半天没有重复,瞬时停摆的三小时才有#
对照要再拉三个指标:
clickhouse_metrics_readonly_replica:处于只读状态的 Replicated 表数量clickhouse_replica_queue_size:复制队列长度clickhouse_events_failed_query:失败的查询数
按小时数一遍重复行(这个数怎么低成本地数出来,写在清理 ClickHouse 重复行的 REPLACE PARTITION runbook),44 个小时窗口里 234 行多余数据全在 07-30 的 3 个小时里。前一天有一段更长、监控上更难看的窗口,两段放在一起:
| 窗口(UTC) | 集群状态 | failed/s | INSERT/s | 多余行 |
|---|---|---|---|---|
| 07-29 09–20Z | readonly_replica 最高 43,复制队列峰值 18,274 |
≈0 | 130 降到 24–54,最低 9 | 0 |
| 07-30 00–03Z | readonly_replica 瞬时采样为 0,30 分钟窗 max_over_time 抓到 01:00 和 02:00–02:30 两次短暂闪断;复制队列 118 → 9,867 → 1,105 |
0 | 102–137,接近常态 | 234 |
07-29 那段 failed/s 全程接近 0,sink 一次都没被拒过。看起来是负载均衡把转成只读的副本摘出了路由,写请求由健康副本接,只是慢。Aiven 没有给出这一层的指标,这是从 failed≈0 反推的。这 11 个小时里副本转只读持续了 7.5 小时,INSERT 一直慢着但都成功,一行重复都没有;按上一节的机制,这意味着没有一条 INSERT 既卡过 30 秒又最终提交。吞吐谷底只剩常态的 7%,是一次可用性事故,数据没问题。
07-30 那 3 小时不一样。吞吐看着是常态,readonly_replica 用瞬时值看是 0,把采样换成 30 分钟窗的 max_over_time 才看到两次短暂闪断。Keeper 的在途请求则明显涨了上去。clickhouse_metrics_zoo_keeper_request 常态是 5 到 15,那几个时间点是:
| 时刻(UTC) | Keeper 在途请求 |
|---|---|
| 07-30 00:00 | 3,415 |
| 07-30 00:10 | 6,256 |
| 07-30 00:40 | 2,157 |
| 07-30 02:35 | 1,748 |
几百倍的在途堆积表示请求发出去了、应答没回来。吞吐没有跟着下降,停摆的时间很短,多数批次照常走完,只有赶上尖峰的那几批卡过了 30 秒。
这里能排除的只有 merge 背压。磁盘、S3、fsync、CPU 同样能让一条 INSERT 变慢,Prometheus 里没有这几个 ClickHouse 指标,Aiven 控制面的月粒度也够不到 7 月,都排除不掉。只能说卡在 Keeper 或者它背后的存储路径上。
两段的区别在于慢的方式。持续变慢时,每条 INSERT 都比平时久,但没有一条超过 client 的 30 秒超时线,框架层没有重投过,也就写不出第二份。
瞬时停摆时,吞吐整体正常,个别批次卡在 Keeper 应答上超过 30 秒。sink 只看到超时,框架等到下一个提交点把同一批再投一次,而服务端把第一次那份提交了,于是有了第二份。
上面说的服务端提交成功、只是 ack 没回来,这一句是推断,不是观测。query_log 只留 4 天,事故当天的 INSERT 时长已经查不到;Prometheus 的 clickhouse_query_duration_ms{quantile} 顶不上,它是在途查询的快照,被长 SELECT 主导(同一个 24 小时里它的 p50 是 47ms,query_log 里 INSERT 的 p50 只有 8ms)。
支撑它的直接证据只有一条:failed_query 在那 3 小时是 0,INSERT 没有报错返回,那它大概率是跑完了。这条和上一节排除背压用的 delayed_inserts / rejected_inserts 出自同一批服务端指标,合起来只排除掉背压和显式报错两类失败,彼此算不上独立。
写出重复的那 3 小时里,readonly_replica 只闪了两下、每次不到一个采样周期,瞬时采样干脆看不见,带 for: 5m 的告警规则也不会响。要看的是 Keeper 在途请求、复制队列长度,还有 sink 侧插入耗时的尾部(connector 维度的 sink_task_put_batch_max_time_milliseconds,具体名字以 Metrics API 的 descriptors 端点返回为准)。同一套系统在另一个 region 的那条链路,插入耗时尾部已经到了 24.4 秒,接近 30 秒超时的八成,下一次 Keeper 抖动大概率是它。
小结#
Read timed out要看的是慢的方式:每条都慢、但都没过 30 秒,写不出重复;吞吐正常、个别批次过了 30 秒而服务端仍提交成功,才会写出第二行。两种慢法的failed_query都接近 0,区分靠吞吐和 Keeper 在途请求。- 30 秒是 clickhouse-java V1 client 的默认
socket_timeout,不是插件的timeoutSeconds(那只管 ping),retryCount也不管 INSERT。想改这个值,插件配置里唯一能碰到它的键是jdbcConnectionProperties。 - 框架重投的是保留在内存里的同一批,而且要等到下一个 offset 提交点才投,所以一轮失败的间隔落在 (30, 90] 秒之间。把 32 / 60 / 90 / 124 读成一到四轮超时是不对的。
- 超时最容易联想到 part 太多、merge 排不过来,这个方向查
delayed_inserts和rejected_inserts两个指标就能否掉,事故窗口里两个全程是 0。 - 告警不能只看副本只读这类硬故障指标,写出重复的那 3 小时里它只闪了两下,瞬时采样看不见。Keeper 在途请求、复制队列长度和写入耗时的尾部更接近重复行的前兆。
相关文章#
这四篇都在同一个 Aiven 托管的 ClickHouse 集群上。(一)是一次补数方案评审,(二)到(四)是一次重复行事故,从定位到清理。
- 用只读权限评审 ClickHouse 补数方案 — 同一个集群上的 mutation 评审:列级重写的成本怎么算、回滚为什么要按 id 名单而不是按状态
- ClickHouse 里的重复行来自 Kafka Connect 超时重投(本篇)
- ClickHouse 的块级去重窗口 — 服务端那层为什么没兜住:窗口按块数算,这张表的建块速率折合 8 秒
- 清理 ClickHouse 重复行的 REPLACE PARTITION runbook — 怎么低成本数出有多少、怎么清掉:临时表加 REPLACE PARTITION,以及跑之前要验的四件事
参考资料#
- ClickHouse Kafka Connect Sink — 官方文档这一页没有
timeoutSeconds/retryCount ClickHouseSinkConfig.java@ v1.3.9 —clientVersion默认V1(:291)、timeoutSeconds默认 30、retryCount默认 3、exactlyOnce默认 false、bufferCount默认 0 的出处;build.gradle.kts@ v1.3.9 钉的 clickhouse-java 是 0.9.5ClickHouseHelperClient.java/ClickHouseWriter.java@ v1.3.9 —retryCount的四个重试循环只包住 ping / query(:189、:205、:268、:294);timeoutSeconds只进ping()(:190、:206);建 V1 节点不传 socket 超时(:101–133);INSERT 在 V1 路径是future.get()(:1086、:1172、:1337)、V2 路径是client.insert(...).get()(:1033、:1248、:1412),都没有重试循环;insert_deduplication_token在 V1 由getMutationRequest设(:1449)、V2 由InsertSettings设(:1002)QueryIdentifier.java@ v1.3.9 — token 的拼法topic-partition-minOffset-maxOffset(:56–60)Utils.java@ v1.3.9 —handleException把哪些错误码转成RetriableException(含 159 TIMEOUT_EXCEEDED、209 SOCKET_TIMEOUT、210 NETWORK_ERROR、999 KEEPER_EXCEPTION),SocketTimeoutException也在内(:105)- clickhouse-java v0.9.5:
ClickHouseClientOption.java(V1SOCKET_TIMEOUT默认取ClickHouseDataConfig.DEFAULT_TIMEOUT,:289)、ClickHouseDataConfig.java(DEFAULT_TIMEOUT = 30 * 1000,:161)、client-v2ClientConfigProperties.java(socket_timeout默认"0",:66) - Kafka 3.7
WorkerSinkTask.java—deliverMessages()里RetriableException的处理(:620–628)、messageBatch只在成功后清空(:612)、下一轮poll(nextCommit - now)(:249)、context.timeout()(:338–341)、retryWithToleranceOperator只包转换(:533–541);WorkerConfig.javaoffset.flush.interval.ms默认 60000(:107) - issue #801 — 同一个 connector 上另一条写重的路径:开了 buffer(
bufferCount > 0、exactlyOnce=false)时flushBuffer()先 flush 后清空,重试异常让清空被跳过,Connect 重投的记录被再 append 一次。2026-07-30 提出,截至 2026-09-16 未修,PR #835 待合入。这条链路bufferCount是默认 0,不沾 - PR #756 — 给 V2 client 的 insert 加上限时(
clickhouseClientInsertTimeoutMs,默认 240000),套在future.get()外面、只包 V2 路径。2026-06-16 合入,晚于 v1.3.9(2026-05-22),首个含它的 release 是 v1.3.10(2026-06-25);这套跑的版本里没有,默认走的也是 V1 - ClickHouse Sink Connector for Confluent Cloud — 全托管版暴露
errors.retry.timeout(默认 30000ms,可填 0–300000ms),页面上写的是 “the retry budget, in milliseconds, for failed record inserts”,超时后进 DLQ 或让 task 失败;本篇按源码判断它只覆盖转换阶段(WorkerSinkTask:533–541),说法以那段代码为准。同一页 Connection details 下有jdbcConnectionProperties,exactlyOnce默认 false,没有timeoutSeconds/retryCount;FAQ 写明默认 at-least-once、重放会产生重复行 - Confluent Cloud Metrics API — connector 维度的指标要从
descriptors/datasets端点取,页面正文没有逐个列出指标名;正文那个sink_task_put_batch_max_time_milliseconds待与 descriptors 返回再核一次[need manual confirm] - Aiven for ClickHouse 服务架构 — 单 shard 三节点、无主从、连接随机落到任一节点,
MergeTree自动改写成ReplicatedMergeTree