You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Kafka Producer flush方法是否原子?JDBC源连接器数据重复场景确认

JDBC Source Connector数据重复与Flush超时问题解析

问题背景

我正在使用JDBC Source Connector,近期在Kafka中遇到数据重复问题。经调试发现Producer抛出如下错误:

WorkerSourceTask{id=connect-name-0} Failed to flush, timed out while waiting for producer to flush outstanding 6590 messages (org.apache.kafka.connect.runtime.WorkerSourceTask:509)

修改flush超时时间后,当前错误已解决,但仍不理解数据重复的原因。我的场景设定如下:
源数据库有20k条数据,连接器从offset 0开始,batch size为1000,poll间隔2小时,flush间隔1分钟,flush超时5秒:

  1. 连接器/任务执行查询,从数据库获取1000条数据;
  2. 缓冲区内存为32MB,任务缓存数据并每分钟尝试flush,假设已缓存15k条数据;
  3. 因数据量过大,在5秒内无法完成flush及offset提交,操作失败。

现存在三种场景假设:

  • 场景1:若flush是原子操作,则不会提交任何offset,任务重启后不会出现重复数据;
  • 场景2:连接器已将部分数据(如15k中的10k)写入Topic,但因flush超时未提交offset,offset仍为0;
  • 场景3:flush超时后,已写入的数据会回滚,offset提交也会回滚。

请问哪种场景符合实际情况?另外,Kafka Producer的flush方法是否为原子操作?


解答

实际发生的场景

场景2完全符合实际情况。当flush超时触发时,WorkerSourceTask会终止当前flush流程,但此时Kafka Producer可能已经将部分消息成功发送到集群并获得Broker确认——因为Producer的消息发送是异步批量执行的,并非要等所有消息处理完成才推进流程。而offset提交的触发条件是所有缓存消息都成功完成flush,所以一旦flush超时,offset不会被更新,仍停留在之前的位置(如示例中的0)。当任务重启后,会从这个旧offset重新拉取数据,导致已经写入Topic的那部分数据被重复推送,这就是数据重复的根源。

Kafka Producer flush方法的原子性

Kafka Producer的flush()不是原子操作。它的核心逻辑是等待所有未完成的发送请求得到Broker确认,同时触发缓冲区中待发送的消息立即推送。在flush过程中,消息是分批次处理的:已经完成确认的消息会永久写入Topic,若中途发生超时,未完成的发送请求会被终止,但已成功写入的消息不会回滚。简言之,flush是一个“尽力完成所有请求”的操作,不具备全成或全败的原子性。


内容的提问来源于stack exchange,提问作者aditya81070

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.01 00:55:19