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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 07:06:07