Embulk采集SFTP文件至Kafka出现重复消息问题排查咨询
排查Embulk推送Kafka时的重复消息问题
针对你使用Embulk v0.10.12 + SFTP输入插件 + embulk-output-kafka出现的重复消息问题,结合你提供的配置,从以下几个方向排查:
1. Embulk增量模式的状态跟踪异常
你启用了incremental: true,Embulk依赖状态存储记录已处理的文件列表以避免重复采集。若状态存储异常,会导致Embulk重新处理已完成的文件:
- 检查默认状态存储路径
~/.embulk/state下的任务状态文件,确认文件权限正常、无损坏或丢失; - 若为集群部署Embulk,需确认是否使用分布式状态存储(如JDBC状态插件),单节点本地状态会导致多节点重复处理;
- 查看Embulk运行日志,确认每次任务启动时加载的状态是否正确,是否存在状态回滚情况。
2. Kafka幂等性配置的兼容性问题
虽然你设置了enable.idempotence=true和acks=all,但仍需验证以下细节:
- Broker版本兼容性:Kafka幂等性要求Broker版本≥0.11.0.0,版本过低会导致幂等性失效;
- Broker副本配置:确保
transaction.state.log.replication.factor≥3且min.insync.replicas≥2,否则acks=all无法保证消息不丢失,且幂等性依赖的事务日志无法正常同步; - 生产者参数冲突:
enable.idempotence=true时,Kafka会自动将max.in.flight.requests.per.connection限制为5(你的配置符合要求),但需确认Broker是否开启事务日志支持(transaction.state.log.enabled=true为默认值); - 重试次数配置:你设置了
retries=1,若重试发生在幂等性未生效的场景(如Broker临时不可用),可能导致重复。可尝试将retries设为Integer.MAX_VALUE(配合delivery.timeout.ms控制总超时),让Kafka自动处理重试逻辑。
3. embulk-output-kafka插件的事务逻辑缺陷
部分旧版本的embulk-output-kafka插件可能存在事务处理逻辑问题:
- 检查插件版本,建议升级至最新稳定版(如v1.1.0+),旧版本可能存在"标记文件为已处理"早于"确认Kafka消息写入成功"的逻辑错误;
- 确认插件是否支持**精确一次(Exactly-Once)**语义,若仅支持至少一次(At-Least-Once),即使Kafka幂等性开启,仍可能因Embulk重试导致重复。
4. 文件源或传输过程的重复
- 验证SFTP源文件:检查是否存在重复上传的文件,或文件本身包含重复记录;
- 查看SFTP服务器日志,确认文件是否被多次读取或修改,导致Embulk重复采集。
5. 网络分区导致的隐性重试
即使开启幂等性,极端网络场景下仍可能出现重复:
- 查看Kafka Broker日志,搜索"Duplicate sequence number"关键字,确认是否有幂等性去重的记录;
- 检查网络监控数据,确认是否存在间歇性网络分区,导致生产者未收到Broker确认而触发重试(此时幂等性应生效,若仍出现重复,需排查Broker的幂等性配置)。
6. Embulk任务的意外重试
- 查看Embulk任务的运行历史,确认是否存在任务意外终止(如OOM、进程被杀)后自动重启的情况,此时若状态未及时持久化,会导致重新处理文件。
内容的提问来源于stack exchange,提问作者user14392764
相关产品推荐
相关产品推荐

