Kafka Producer flush方法是否原子?JDBC源连接器数据重复场景确认
问题背景
我正在使用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秒:
- 连接器/任务执行查询,从数据库获取1000条数据;
- 缓冲区内存为32MB,任务缓存数据并每分钟尝试flush,假设已缓存15k条数据;
- 因数据量过大,在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

