Kafka分布式Connect运行FTP源连接器产生重复消息问题咨询
- 共3台服务器
- 每台服务器部署1个Kafka broker、Connect、Schema Registry组件,使用版本为confluent-7.1.0
- 部署1个FTP测试源连接器,配置最大任务数为3
- Connect运行过程中产生重复消息,预期效果为FTP连接器针对每个文件仅生成1条消息
分布式Connect同一条消息重复生成3次(每个Connect任务各生成1次)
- 连接器任务生成消息时的运行日志在每个Connect进程中均有打印,日志内容如下:
[2022-06-26 15:23:12,839] INFO [ftp-test-conn|task-0] poll (com.datamountaineer.streamreactor.connect.ftp.source.FtpSourcePoller:77) [2022-06-26 15:23:12,839] INFO [ftp-test-conn|task-0] connect 10.0.0.138:None (com.datamountaineer.streamreactor.connect.ftp.source.FtpMonitor:294) [2022-06-26 15:23:12,862] INFO [ftp-test-conn|task-0] successfully connected to the ftp server and logged in (com.datamountaineer.streamreactor.connect.ftp.source.FtpMonitor:311) [2022-06-26 15:23:12,863] INFO [ftp-test-conn|task-0] passive we are (com.datamountaineer.streamreactor.connect.ftp.source.FtpMonitor:318) [2022-06-26 15:23:12,870] INFO [ftp-test-conn|task-0] Found 4 items in /home/smheo/ftp-dir/* (com.datamountaineer.streamreactor.connect.ftp.source.FtpMonitor:245) [2022-06-26 15:23:12,877] INFO [ftp-test-conn|task-0] meta store storage HASN'T /home/smheo/ftp-dir/msg-4 (com.datamountaineer.streamreactor.connect.ftp.source.ConnectFileMetaDataStore:48) [2022-06-26 15:23:12,878] INFO [ftp-test-conn|task-0] fetching /home/smheo/ftp-dir/msg-4 (com.datamountaineer.streamreactor.connect.ftp.source.FtpMonitor:102) [2022-06-26 15:23:12,881] INFO [ftp-test-conn|task-0] fetched /home/smheo/ftp-dir/msg-4, wasn't known before (com.datamountaineer.streamreactor.connect.ftp.source.FtpMonitor:218) [2022-06-26 15:23:12,881] INFO [ftp-test-conn|task-0] dump entire /home/smheo/ftp-dir/msg-4 (com.datamountaineer.streamreactor.connect.ftp.source.FtpMonitor:219) [2022-06-26 15:23:12,881] INFO [ftp-test-conn|task-0] got some fileChanges: /home/smheo/ftp-dir/msg-4, offset = -1 (com.datamountaineer.streamreactor.connect.ftp.source.FtpSourcePoller:96)
消费者端消费到重复消息
执行kafka-console-consumer命令消费对应topic时,同一条内容重复出现3次,执行命令及消费结果如下:
(base) ubuntu@ubuntu:~/distributed-pipeline/confluent-7.1.0$ ./bin/kafka-console-consumer --bootstrap-server <BROKER_IP>:9092 --topic default-topic-1 hello hello hello
FTP连接器配置与状态信息
当前FTP连接器的配置参数及运行状态如下:
{ "ftp-test-conn": { "info": { "name": "ftp-test-conn", "config": { "connector.class": "com.datamountaineer.streamreactor.connect.ftp.source.FtpSourceConnector", "connect.ftp.address": "<FTP HOST IP>", "connect.ftp.keystyle": "string", "compression.type": "gzip", "connect.ftp.user": "ftpusername", "connect.ftp.refresh": "PT1M", "tasks.max": "3", "connect.ftp.file.maxage": "P7D", "name": "ftp-test-conn", "connect.ftp.monitor.update": "/home/username/ftp-dir/:default-topic-1", "connect.ftp.timeout": "3000000", "connect.ftp.password": "<PASSWORD>" }, "tasks": [ { "connector": "ftp-test-conn", "task": 0 }, { "connector": "ftp-test-conn", "task": 1 }, { "connector": "ftp-test-conn", "task": 2 } ], "type": "source" }, "status": { "name": "ftp-test-conn", "connector": { "state": "RUNNING", "worker_id": "<BROKER 1 IP>:8083" }, "tasks": [ { "id": 0, "state": "RUNNING", "worker_id": "<BROKER 1 IP>:8083" }, { "id": 1, "state": "RUNNING", "worker_id": "<BROKER 2 IP>:8083" }, { "id": 2, "state": "RUNNING", "worker_id": "<BROKER 3 IP>:8083" } ], "type": "source" } } }
当前使用的StreamReactor FTP源连接器本身不支持单监控路径下的多任务并行处理。
从日志打印的meta store storage HASN'T /home/smheo/ftp-dir/msg-4可以直接定位问题:该连接器的已处理文件元数据存储是任务级本地存储,不支持跨任务共享。配置tasks.max=3后,3个分散在不同节点的任务会独立建立FTP连接、独立全量扫描配置的/home/username/ftp-dir/路径,每个任务都判定目标文件未被自己处理过,因此各自生成1条消息推送到Kafka,最终同一份文件的内容重复出现3次。
本质上这个连接器没有实现多任务下的负载拆分逻辑(比如按文件哈希分配、按子目录分配),配置多少个任务,每个任务就会完整扫描所有配置的监控路径,任务数和消息重复倍数完全相等。
直接将连接器配置中的tasks.max参数修改为1,重启连接器即可。
该FTP连接器单任务即可覆盖所有配置的监控路径,多任务配置不仅无法提升处理性能,反而会直接导致重复消息问题。调整后仅会有1个任务执行FTP路径扫描,每个文件仅会被处理1次,符合每个FTP文件生成1条消息的预期。
如果后续需要提升FTP文件处理吞吐量,只能通过拆分监控目录、为不同目录创建独立FTP连接器实例的方式横向扩展,不能通过调大单连接器的任务数实现。
内容的提问来源于stack exchange,提问作者Seongmin Heo

