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

MirrorSourceConnector配置异常:Kafka跨集群复制任务未启动排查

Kafka Mirror Maker 2.0 跨集群数据复制失败排查与修复

问题背景

基于Docker Compose搭建了两个独立Kafka集群(broker1+zookeeper、broker2+zookeeper2),通过Kafka Connect部署MirrorSourceConnector实现跨集群主题数据复制,但向broker1发送消息后,broker2无同步数据,且Kafka Connect无运行任务。

当前Docker Compose配置

networks:
  local_kafka:
    name: local_kafka

services:
  zookeeper:
    image: confluentinc/cp-zookeeper:7.5.2
    container_name: zookeeper
    networks:
      - local_kafka
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000
      KAFKA_OPTS: "-Dzookeeper.4lw.commands.whitelist=*" # 本地调试用

  zookeeper2:
    image: confluentinc/cp-zookeeper:7.5.2
    container_name: zookeeper2
    networks:
      - local_kafka
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000
      KAFKA_OPTS: "-Dzookeeper.4lw.commands.whitelist=*" # 本地调试用

  broker1:
    image: confluentinc/cp-kafka:7.5.2
    container_name: broker1
    networks:
      - local_kafka
    ports:
      - "19092:19092"
    depends_on:
      - zookeeper
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_BOOTSTRAP_SERVERS: broker1:9092
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,CONNECTIONS_FROM_HOST:PLAINTEXT
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://broker1:9092,CONNECTIONS_FROM_HOST://localhost:19092
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      KAFKA_MIN_INSYNC_REPLICAS: 1
      ALLOW_PLAINTEXT_LISTENERS: "true"
      KAFKA_AUTO_CREATE_TOPICS_ENABLE: "false"
      KAFKA_RETRIES: 3

  broker2:
    image: confluentinc/cp-kafka:7.5.2
    container_name: broker2
    networks:
      - local_kafka
    ports:
      - "19093:19093"
    depends_on:
      - zookeeper2
    environment:
      KAFKA_BROKER_ID: 2
      KAFKA_BOOTSTRAP_SERVERS: broker2:9092
      KAFKA_ZOOKEEPER_CONNECT: zookeeper2:2181
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,CONNECTIONS_FROM_HOST:PLAINTEXT
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://broker2:9092,CONNECTIONS_FROM_HOST://localhost:19093
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      KAFKA_MIN_INSYNC_REPLICAS: 1
      KAFKA_SCHEMA_REGISTRY_URL: schemaregistry:8081
      ALLOW_PLAINTEXT_LISTENERS: "true"
      KAFKA_AUTO_CREATE_TOPICS_ENABLE: "false"
      KAFKA_RETRIES: 3

  schemaregistry:
    image: confluentinc/cp-schema-registry:7.5.2
    container_name: schemaregistry
    networks:
      - local_kafka
    environment:
      SCHEMA_REGISTRY_HOST_NAME: schemaregistry
      SCHEMA_REGISTRY_LISTENERS: "http://0.0.0.0:8081"
      SCHEMA_REGISTRY_KAFKASTORE_SECURITY_PROTOCOL: PLAINTEXT
      SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: PLAINTEXT://broker2:9092
      SCHEMA_REGISTRY_DEBUG: "true"
      SCHEMA_REGISTRY_KAFKASTORE_TOPIC: _registry-schemas
    ports:
      - "8081:8081"

  kafka-connect:
      image: confluentinc/cp-kafka-connect-base:7.6.1
      container_name: kafka-connect
      networks:
        - local_kafka
      depends_on:
        - broker1
        - broker2
        - schemaregistry
        - zookeeper
        - zookeeper2
      environment:
        CONNECT_BOOTSTRAP_SERVERS: PLAINTEXT://broker2:9092
        CONNECT_SCHEMA_REGISTRY_URL: http://schemaregistry:8081
        CONNECT_REST_PORT: 8083
        CONNECT_GROUP_ID: kafka-connect
        CONNECT_CONFIG_STORAGE_TOPIC: _connect-configs
        CONNECT_OFFSET_STORAGE_TOPIC: _connect-offsets
        CONNECT_STATUS_STORAGE_TOPIC: _connect-status
        CONNECT_KEY_CONVERTER: org.apache.kafka.connect.storage.StringConverter
        CONNECT_VALUE_CONVERTER: io.confluent.connect.protobuf.ProtobufConverter
        CONNECT_VALUE_CONVERTER_SCHEMA_REGISTRY_URL: http://schemaregistry:8081
        CONNECT_REST_ADVERTISED_HOST_NAME: "kafka-connect"
        CONNECT_CONFIG_STORAGE_REPLICATION_FACTOR: "2"
        CONNECT_OFFSET_STORAGE_REPLICATION_FACTOR: "2"
        CONNECT_STATUS_STORAGE_REPLICATION_FACTOR: "2"
        CONNECT_PLUGIN_PATH: /usr/local/share/kafka/plugins,/usr/share/filestream-connectors
        KAFKA_LOG4J_OPTS: -Dlog4j.configuration=file:/etc/kafka/log4j.properties

      volumes:
        - ./log4j.properties:/etc/kafka/log4j.properties
      command:
        - bash
        - -c
        - |
          echo "Installing Protobuf Connector"
          confluent-hub install --no-prompt confluentinc/kafka-connect-protobuf-converter:7.3.3
          echo "Launching Kafka Connect worker"
          /etc/confluent/docker/run &
          sleep infinity
      ports:
        - "8083:8083"

当前MirrorSourceConnector配置

{
    "name": "local-to-remote-mirror",
    "config":{
        "connector.class": "org.apache.kafka.connect.mirror.MirrorSourceConnector",
        "tasks.max": "2",
        "topics": ".*",
        "replication.policy.separator": "",
        "source.cluster.alias": "",
        "target.cluster.alias": "",
        "source.cluster.bootstrap.servers": "http://broker1:9092",
        "target.cluster.bootstrap.servers": "http://broker2:9092",
        "offset-syncs.topic.location": "target",
        "offset-syncs.topic.replication.factor": "2",
        "sync.topic.acls.enabled": false,
        "sync.topic.configs.enabled": false,
        "replication.policy.class": "org.apache.kafka.connect.mirror.IdentityReplicationPolicy",

        "_comment": "Kafka Connect converter used to deserialize keys",
        "key.converter": "org.apache.kafka.connect.converters.ByteArrayConverter",

        "_comment": "Kafka Connect converter used to deserialize values",
        "value.converter": "org.apache.kafka.connect.converters.ByteArrayConverter"
    }
}

关键问题与修复方案

1. 集群别名不能为空

MM2依赖集群别名区分源和目标集群,空别名会导致Connector无法识别集群身份,无法初始化任务。

  • 修正:设置唯一别名,例如:
"source.cluster.alias": "source",
"target.cluster.alias": "target"

2. Bootstrap地址协议错误

Kafka Broker的Bootstrap地址不需要http://前缀,直接使用Broker内部端口即可(默认PLAINTEXT协议):

  • 修正:
"source.cluster.bootstrap.servers": "broker1:9092",
"target.cluster.bootstrap.servers": "broker2:9092"

3. 副本因子配置不匹配

目标集群broker2仅1个节点,但offset-syncs.topic.replication.factor设为2,同时Kafka Connect的存储主题副本因子也设为2,这会导致主题无法创建,Connector启动失败。

  • 修正Connector配置:
"offset-syncs.topic.replication.factor": "1"
  • 修正Docker Compose中kafka-connect的环境变量:
CONNECT_CONFIG_STORAGE_REPLICATION_FACTOR: "1"
CONNECT_OFFSET_STORAGE_REPLICATION_FACTOR: "1"
CONNECT_STATUS_STORAGE_REPLICATION_FACTOR: "1"

4. 缺少MM2插件安装

当前Kafka Connect仅安装了Protobuf转换器,未安装Mirror Maker 2.0插件,导致无法识别MirrorSourceConnector类。

  • 修正kafka-connect的启动脚本,添加MM2插件安装(版本需与Connect版本7.6.1匹配):
command:
  - bash
  - -c
  - |
    echo "Installing Protobuf Connector"
    confluent-hub install --no-prompt confluentinc/kafka-connect-protobuf-converter:7.3.3
    echo "Installing Mirror Maker 2.0 Connector"
    confluent-hub install --no-prompt confluentinc/kafka-connect-mirror-maker:7.6.1
    echo "Launching Kafka Connect worker"
    /etc/confluent/docker/run &
    sleep infinity

5. 源集群主题自动创建关闭

源集群broker1设置了KAFKA_AUTO_CREATE_TOPICS_ENABLE: "false",如果测试用的源主题未提前创建,Connector无法读取数据。

  • 临时开启自动创建用于测试:
    修改broker1的环境变量:
KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true"

测试完成后可重新关闭。

6. 验证Connector状态

修复上述配置后,重启Kafka Connect服务并重新提交Connector配置,通过以下命令查看状态:

curl -X GET http://localhost:8083/connectors/local-to-remote-mirror/status

若任务仍未启动,查看Connect日志排查具体错误:

docker logs kafka-connect

内容的提问来源于stack exchange,提问作者J. Snow

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 01:54:50