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

3VM Kafka集群中MongoDB源汇连接器仅同步单一VM的问题咨询

问题原因及解决办法

核心问题:消费者组冲突

你的第二、第三台VM上的Sink连接器配置中,name字段都设为了sink1。Kafka Connect默认会将连接器的name作为消费者组ID(group.id)的一部分,当两个Sink连接器属于同一个消费者组时,Kafka的消费者组机制会把目标主题Ecomm1.orders的分区分配给组内的单个消费实例,因此只会有一台VM的MongoDB收到同步数据。

解决步骤

  1. 修改Sink连接器名称:给第二、第三台VM的Sink连接器设置不同的name,确保它们属于独立的消费者组。修改后的配置示例:
    • VM2的Sink配置:
      {
        "name": "sink-vm2",
        "config": {
          "connector.class": "com.mongodb.kafka.connect.MongoSinkConnector",
          "topics": "Ecomm1.orders",
          "connection.uri": "mongodb://localhost:27017",
          "key.converter": "org.apache.kafka.connect.storage.StringConverter",
          "value.converter": "org.apache.kafka.connect.json.JsonConverter",
          "value.converter.schemas.enable": false,
          "database": "Ecomm",
          "collection": "orders"
        }
      }
      
    • VM3的Sink配置:
      {
        "name": "sink-vm3",
        "config": {
          "connector.class": "com.mongodb.kafka.connect.MongoSinkConnector",
          "topics": "Ecomm1.orders",
          "connection.uri": "mongodb://localhost:27017",
          "key.converter": "org.apache.kafka.connect.storage.StringConverter",
          "value.converter": "org.apache.kafka.connect.json.JsonConverter",
          "value.converter.schemas.enable": false,
          "database": "Ecomm",
          "collection": "orders"
        }
      }
      
  2. 重启连接器生效:重启两台VM上的Sink连接器后,它们会以独立消费者组身份订阅主题,此时每个连接器都会收到主题的全部分区消息,实现第一台VM的插入操作同步到另外两台MongoDB。

额外说明

  • 若想显式控制消费者组ID,可在Sink连接器的config中添加consumer.group.id字段,给两个Sink设置不同值,效果与修改name一致。
  • 即使目标主题只有1个分区,不同消费者组的实例也能同时收到消息,因为消费者组之间是独立订阅关系。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 04:10:20