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

Debezium对接MongoDB无法创建Kafka主题,求排查docker-compose配置

Debezium无法创建Kafka主题的配置排查与修正

核心问题点及修正方案

1. Connect环境变量前缀缺失

Debezium Connect容器的核心配置环境变量必须以CONNECT_为前缀,否则无法被Connect进程识别。原配置中以下变量缺少前缀:

  • BOOTSTRAP_SERVERS → 改为CONNECT_BOOTSTRAP_SERVERS
  • GROUP_ID → 改为CONNECT_GROUP_ID
  • CONFIG_STORAGE_TOPIC → 改为CONNECT_CONFIG_STORAGE_TOPIC
  • OFFSET_STORAGE_TOPIC → 改为CONNECT_OFFSET_STORAGE_TOPIC

修正后的Debezium环境变量片段:

environment:
  CONNECT_BOOTSTRAP_SERVERS: kafka:9092
  CONNECT_GROUP_ID: 1
  CONNECT_CONFIG_STORAGE_TOPIC: connect_configs
  CONNECT_OFFSET_STORAGE_TOPIC: connect_offsets
  CONNECT_KEY_CONVERTER: io.confluent.connect.avro.AvroConverter
  CONNECT_VALUE_CONVERTER: io.confluent.connect.avro.AvroConverter
  CONNECT_KEY_CONVERTER_SCHEMA_REGISTRY_URL: http://schema-registry:8081
  CONNECT_VALUE_CONVERTER_SCHEMA_REGISTRY_URL: http://schema-registry:8081
  CONNECT_TOPIC_CREATION_ENABLE: true

2. Debezium版本与MongoDB版本不兼容

原配置使用debezium/connect:1.4,该版本仅支持MongoDB 4.4及以下版本,无法适配MongoDB 7.0.5。需升级Debezium版本至2.4或更高(推荐2.4.2.Final,兼容MongoDB 7.x):

image: debezium/connect:2.4.2.Final

3. MongoDB副本集未初始化

Debezium依赖MongoDB副本集捕获变更日志,但原配置仅启动了两个副本节点,未完成副本集初始化。需在mongo-1服务中添加初始化逻辑:

mongo-1:
  image: mongo:7.0.5
  container_name: mongo1
  ports:
    - "27017:27017"
  command: >
    bash -c "
      mongod --replSet rs0 --port 27017 --bind_ip_all;
      sleep 10;
      mongo --port 27017 --eval 'rs.initiate({_id: \"rs0\", members: [{_id: 0, host: \"mongo1:27017\"}, {_id: 1, host: \"mongo2:27018\"}]});'
    "
  depends_on:
    - mongo-2
  volumes:
    - ./mongo-data:/data/db
  healthcheck:
    test: ["CMD", "mongo", "--eval", "db.adminCommand('ping')"]
    interval: 10s
    timeout: 5s
    retries: 5

4. Schema Registry冗余配置清理

Schema Registry容器中添加了属于Kafka的环境变量,这些配置无效且易引发混淆,需删除:

  • KAFKA_ADVERTISED_LISTENERS
  • KAFKA_LISTENER_SECURITY_PROTOCOL_MAP

5. Debezium MongoDB连接器配置位置修正

DEBEZIUM_MONGODB_URI和DEBEZIUM_MONGODB_DB是MongoDB连接器的专属参数,不是Connect容器的环境变量,需通过REST API创建连接器时指定。示例请求:

curl -X POST -H "Content-Type: application/json" http://localhost:8085/connectors -d '{
  "name": "mongodb-source-connector",
  "config": {
    "connector.class": "io.debezium.connector.mongodb.MongoDbConnector",
    "mongodb.hosts": "rs0/mongo1:27017,mongo2:27018",
    "mongodb.database": "carbon_emission_db",
    "mongodb.name": "mongodb-source",
    "topic.prefix": "mongodb",
    "schema.history.internal.kafka.bootstrap.servers": "kafka:9092",
    "schema.history.internal.kafka.topic": "mongodb-schema-history"
  }
}'

修正后的完整docker-compose.yml

services:
  # DEBEZIUM CONFIGURATION
  debezium:
    image: debezium/connect:2.4.2.Final
    container_name: debezium
    environment:
      CONNECT_BOOTSTRAP_SERVERS: kafka:9092
      CONNECT_GROUP_ID: 1
      CONNECT_CONFIG_STORAGE_TOPIC: connect_configs
      CONNECT_OFFSET_STORAGE_TOPIC: connect_offsets
      CONNECT_KEY_CONVERTER: io.confluent.connect.avro.AvroConverter
      CONNECT_VALUE_CONVERTER: io.confluent.connect.avro.AvroConverter
      CONNECT_KEY_CONVERTER_SCHEMA_REGISTRY_URL: http://schema-registry:8081
      CONNECT_VALUE_CONVERTER_SCHEMA_REGISTRY_URL: http://schema-registry:8081
      CONNECT_TOPIC_CREATION_ENABLE: true
    depends_on:
      - kafka
      - mongo-1
    ports:
      - 8085:8083

  # KAFKA SERVICES CONFIGURATION
  schema-registry:
    container_name: schema-registry
    image: confluentinc/cp-schema-registry:7.4.4
    environment:
      SCHEMA_REGISTRY_KAFKASTORE_CONNECTION_URL: zookeeper:2181
      SCHEMA_REGISTRY_HOST_NAME: schema-registry
      SCHEMA_REGISTRY_LISTENERS: http://schema-registry:8081,http://localhost:8081
      SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: kafka:9092
      SCHEMA_REGISTRY_LISTENER_HTTP_PORT: 8081
      SCHEMA_REGISTRY_KAFKASTORE_SECURITY_PROTOCOL: PLAINTEXT
    ports:
      - 8081:8081
    depends_on:
      - zookeeper
      - kafka

  zookeeper:
    image: quay.io/debezium/zookeeper
    container_name: zookeeper
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000
    ports:
      - 2181:2181

  kafka:
    image: confluentinc/cp-kafka:7.4.4
    container_name: kafka
    depends_on:
      - zookeeper
    ports:
      - 29092:29092
      - 9092:9092
    environment:
      KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true"
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:29092
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1

  kafka_manager:
    container_name: kafka-mngr
    image: kafkamanager/kafka-manager
    restart: always
    ports:
      - "9010:9010"
    depends_on:
      - zookeeper
      - kafka
    environment:
      ZK_HOSTS: "zookeeper:2181"
      APPLICATION_SECRET: "random-secret"

  # MONGODB CONFIGURATION
  mongo-1:
    image: mongo:7.0.5
    container_name: mongo1
    ports:
      - "27017:27017"
    command: >
      bash -c "
        mongod --replSet rs0 --port 27017 --bind_ip_all;
        sleep 10;
        mongo --port 27017 --eval 'rs.initiate({_id: \"rs0\", members: [{_id: 0, host: \"mongo1:27017\"}, {_id: 1, host: \"mongo2:27018\"}]});'
      "
    depends_on:
      - mongo-2
    volumes:
      - ./mongo-data:/data/db
    healthcheck:
      test: ["CMD", "mongo", "--eval", "db.adminCommand('ping')"]
      interval: 10s
      timeout: 5s
      retries: 5

  mongo-2:
    image: mongo:7.0.5
    container_name: mongo2
    ports:
      - "27018:27018"
    command: mongod --replSet rs0 --port 27018 --bind_ip_all

  # GRAFANA CONFIGURATION
  grafana:
    image: grafana/grafana:latest
    container_name: grafana
    environment:
      - GF_SECURITY_ADMIN_PASSWORD=admin
      - GF_INSTALL_PLUGINS=grafana-piechart-panel,grafana-clock-panel
    ports:
      - "3010:3000"
    volumes:
      - grafana-storage:/var/lib/grafana

volumes:
  grafana-storage:

内容的提问来源于stack exchange,提问作者Rahil Khan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 02:25:54