MirrorSourceConnector配置异常:Kafka跨集群复制任务未启动排查
问题背景
基于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

