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

Docker Compose部署Debezium CDC架构的持久化问题咨询

问题描述

我正尝试通过Docker Compose搭建基于Debezium、Kafka及Kafka-Connect的CDC(变更数据捕获)架构,需要实现容器停止或重启后,连接器偏移量(connector offsets)与数据仍能保留,但目前Debezium连接器的数据会被重置。

原Docker Compose配置如下:

version: '2'
services:
  mysql:
    image: debezium/example-mysql:1.8
    container_name: mysql
    hostname: mysql
    ports:
      - '3306:3306'
    environment:
      MYSQL_ROOT_PASSWORD: debezium
      MYSQL_USER: mysqluser
      MYSQL_PASSWORD: mysqlpw

  zookeeper:
    image: confluentinc/cp-zookeeper:5.5.3
    container_name: zookeerper
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000
    ports:
      - '2181:2181'

  kafka:
    image: confluentinc/cp-enterprise-kafka:5.5.3
    container_name: kafka
    depends_on:
      - zookeeper
    environment:
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_MAX_REQUEST_SIZE : 200000000
      KAFKA_MESSAGE_MAX_BYTES : 200000000
      KAFKA_MAX_PARTITION_FETCH_BYTES : 200000000
      KAFKA_MAX_REQUEST_SIZE : 200000000
      KAFKA_MESSAGE_MAX_BYTES : 200000000
      KAFKA_MAX_PARTITION_FETCH_BYTES : 200000000
      KAFKA_REPLICA_FETCH_MAX_BYTES: 5048576
      KAFKA_PRODUCER_MAX_REQUEST_SIZE: 5048576
      KAFKA_CONSUMER_MAX_PARTITION_FETCH_BYTES: 5048576
      KAFKA_CONSUMER_MAX_POLL_RECORDS : 300
    links:
      - zookeeper
    ports:
      - '9092:9092'

  schema-registry:
    container_name: schema-registry
    image: confluentinc/cp-schema-registry:4.0.3
    environment:
      - SCHEMA_REGISTRY_KAFKASTORE_CONNECTION_URL=zookeeper:2181
      - SCHEMA_REGISTRY_HOST_NAME=schema-registry
      - SCHEMA_REGISTRY_LISTENERS=http://schema-registry:8081
    ports:
      - '8081:8081'

  kafka-connect:
      container_name: kafka-connect
      image: confluentinc/cp-kafka-connect-base:latest
      ports:
        - '8083:8083'
      links:
        - kafka
        - zookeeper
      environment:
         - CONNECT_BOOTSTRAP_SERVERS=kafka:9092
         - CONNECT_GROUP_ID=medium_debezium
         - CONNECT_CONFIG_STORAGE_TOPIC=my_connect_configs
         - CONNECT_CONFIG_STORAGE_REPLICATION_FACTOR=1
         - CONNECT_OFFSET_STORAGE_TOPIC=my_connect_offsets
         - CONNECT_OFFSET_STORAGE_REPLICATION_FACTOR=1
         - CONNECT_STATUS_STORAGE_TOPIC=my_connect_statuses
         - CONNECT_STATUS_STORAGE_REPLICATION_FACTOR=1
         - CONNECT_REST_ADVERTISED_HOST_NAME=medium_debezium
         - KAFKA_MAX_REQUEST_SIZE=200000000
         - KAFKA_MESSAGE_MAX_BYTES=200000000
         - KAFKA_MAX_PARTITION_FETCH_BYTES=200000000
         - CONNECT_PRODUCER_MAX_REQUEST_SIZE=5048576
         - CONNECT_CONSUMER_MAX_PARTITION_FETCH_BYTES=5048576
         - CONNECT_REST_PORT=8083 
         - CONNECT_KEY_CONVERTER=org.apache.kafka.connect.json.JsonConverter
         - CONNECT_VALUE_CONVERTER=org.apache.kafka.connect.json.JsonConverter
         - CONNECT_INTERNAL_KEY_CONVERTER=org.apache.kafka.connect.json.JsonConverter
         - CONNECT_INTERNAL_VALUE_CONVERTER=org.apache.kafka.connect.json.JsonConverter
         - CONNECT_PLUGIN_PATH=/usr/share/java,/connectors,/usr/share/confluent-hub-components/
      command:
        - bash
        - -c
        - |
          echo "Installing Connector"
          confluent-hub install --no-prompt debezium/debezium-connector-mysql:2.2.1
          confluent-hub install --no-prompt confluentinc/kafka-connect-s3:10.5.1
          #
          echo "Launching Kafka Connect worker"
          /etc/confluent/docker/run &
          #
          sleep infinity
  kafka-ui:
    container_name: kafka-ui
    image: provectuslabs/kafka-ui:latest
    ports:
      - 8080:8080
    environment:
      DYNAMIC_CONFIG_ENABLED: 'true'
  aws-s3:
    image: localstack/localstack:latest
    container_name: aws-s3
    environment:
      SERVICES : s3:5000
      HOSTNAME : aws-s3
      AWS_ACCESS_KEY_ID : "KEY"
      AWS_SECRET_ACCESS_KEY :  "KEY"
    ports:
      - "5000:5000"

  connector-sender:
    container_name: connector-sender
    image: confluentinc/cp-kafka-connect:5.5.3
    volumes:
      - ./conf:/conf
    depends_on:
      - kafka-connect
    command:
        - bash
        - -c
        - |
          echo "Wainting for kafka connect to start..."
          until [[ "$$(curl -s -o /dev/null -w %{http_code} kafka-connect:8083/connectors)" -eq 200 ]]; do
            sleep 1
          done
          echo -e "Sending connector"
          curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" kafka-connect:8083/connectors/ -d @/conf/mysql_config.json
问题原因

容器重启或停止后数据丢失,核心原因是:

  • MySQL、Zookeeper、Kafka的默认数据存储路径都在容器内部,容器销毁后数据随之丢失
  • Kafka Connect的连接器配置和偏移量存储在Kafka的my_connect_configs、my_connect_offsets等topic中,如果Kafka本身没有持久化,这些topic也会被删除
  • MySQL的binlog是Debezium捕获变更的依据,若MySQL数据未持久化,binlog也会丢失,导致Debezium重启后无法恢复之前的偏移量
修改后的Docker Compose配置
version: '2'
services:
  mysql:
    image: debezium/example-mysql:1.8
    container_name: mysql
    hostname: mysql
    ports:
      - '3306:3306'
    environment:
      MYSQL_ROOT_PASSWORD: debezium
      MYSQL_USER: mysqluser
      MYSQL_PASSWORD: mysqlpw
    volumes:
      # 持久化MySQL数据和binlog
      - mysql_data:/var/lib/mysql
      - mysql_binlog:/var/lib/mysql-files

  zookeeper:
    image: confluentinc/cp-zookeeper:5.5.3
    container_name: zookeeper
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000
    ports:
      - '2181:2181'
    volumes:
      # 持久化Zookeeper元数据
      - zookeeper_data:/var/lib/zookeeper/data
      - zookeeper_log:/var/lib/zookeeper/log

  kafka:
    image: confluentinc/cp-enterprise-kafka:5.5.3
    container_name: kafka
    depends_on:
      - zookeeper
    environment:
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_MAX_REQUEST_SIZE : 200000000
      KAFKA_MESSAGE_MAX_BYTES : 200000000
      KAFKA_MAX_PARTITION_FETCH_BYTES : 200000000
      KAFKA_REPLICA_FETCH_MAX_BYTES: 5048576
      KAFKA_PRODUCER_MAX_REQUEST_SIZE: 5048576
      KAFKA_CONSUMER_MAX_PARTITION_FETCH_BYTES: 5048576
      KAFKA_CONSUMER_MAX_POLL_RECORDS : 300
      # 指定Kafka数据存储路径
      KAFKA_LOG_DIRS: /var/lib/kafka/data
    links:
      - zookeeper
    ports:
      - '9092:9092'
    volumes:
      # 持久化Kafka消息数据
      - kafka_data:/var/lib/kafka/data

  schema-registry:
    container_name: schema-registry
    image: confluentinc/cp-schema-registry:4.0.3
    environment:
      - SCHEMA_REGISTRY_KAFKASTORE_CONNECTION_URL=zookeeper:2181
      - SCHEMA_REGISTRY_HOST_NAME=schema-registry
      - SCHEMA_REGISTRY_LISTENERS=http://schema-registry:8081
    ports:
      - '8081:8081'
    volumes:
      # 持久化Schema Registry数据
      - schema_registry_data:/var/lib/schema-registry/data

  kafka-connect:
      container_name: kafka-connect
      image: confluentinc/cp-kafka-connect-base:latest
      ports:
        - '8083:8083'
      links:
        - kafka
        - zookeeper
      environment:
         - CONNECT_BOOTSTRAP_SERVERS=kafka:9092
         - CONNECT_GROUP_ID=medium_debezium
         - CONNECT_CONFIG_STORAGE_TOPIC=my_connect_configs
         - CONNECT_CONFIG_STORAGE_REPLICATION_FACTOR=1
         - CONNECT_OFFSET_STORAGE_TOPIC=my_connect_offsets
         - CONNECT_OFFSET_STORAGE_REPLICATION_FACTOR=1
         - CONNECT_STATUS_STORAGE_TOPIC=my_connect_statuses
         - CONNECT_STATUS_STORAGE_REPLICATION_FACTOR=1
         - CONNECT_REST_ADVERTISED_HOST_NAME=medium_debezium
         - CONNECT_PRODUCER_MAX_REQUEST_SIZE=5048576
         - CONNECT_CONSUMER_MAX_PARTITION_FETCH_BYTES=5048576
         - CONNECT_REST_PORT=8083 
         - CONNECT_KEY_CONVERTER=org.apache.kafka.connect.json.JsonConverter
         - CONNECT_VALUE_CONVERTER=org.apache.kafka.connect.json.JsonConverter
         - CONNECT_INTERNAL_KEY_CONVERTER=org.apache.kafka.connect.json.JsonConverter
         - CONNECT_INTERNAL_VALUE_CONVERTER=org.apache.kafka.connect.json.JsonConverter
         - CONNECT_PLUGIN_PATH=/usr/share/java,/connectors,/usr/share/confluent-hub-components/
      volumes:
        # 持久化已安装的连接器插件,避免每次启动重新安装
        - kafka_connect_plugins:/usr/share/confluent-hub-components
      command:
        - bash
        - -c
        - |
          # 仅在插件目录为空时安装连接器
          if [ -z "$(ls -A /usr/share/confluent-hub-components)" ]; then
            echo "Installing Connector"
            confluent-hub install --no-prompt debezium/debezium-connector-mysql:2.2.1
            confluent-hub install --no-prompt confluentinc/kafka-connect-s3:10.5.1
          fi
          #
          echo "Launching Kafka Connect worker"
          /etc/confluent/docker/run &
          #
          sleep infinity
  kafka-ui:
    container_name: kafka-ui
    image: provectuslabs/kafka-ui:latest
    ports:
      - 8080:8080
    environment:
      DYNAMIC_CONFIG_ENABLED: 'true'
    volumes:
      # 持久化Kafka UI配置
      - kafka_ui_config:/config

  aws-s3:
    image: localstack/localstack:latest
    container_name: aws-s3
    environment:
      SERVICES : s3:5000
      HOSTNAME : aws-s3
      AWS_ACCESS_KEY_ID : "KEY"
      AWS_SECRET_ACCESS_KEY :  "KEY"
    ports:
      - "5000:5000"
    volumes:
      # 持久化LocalStack S3数据
      - localstack_data:/var/lib/localstack

  connector-sender:
    container_name: connector-sender
    image: confluentinc/cp-kafka-connect:5.5.3
    volumes:
      - ./conf:/conf
    depends_on:
      - kafka-connect
    command:
        - bash
        - -c
        - |
          echo "Waiting for kafka connect to start..."
          until [[ "$$(curl -s -o /dev/null -w %{http_code} kafka-connect:8083/connectors)" -eq 200 ]]; do
            sleep 1
          done
          # 检查连接器是否已存在,避免重复创建
          if [ -z "$$(curl -s kafka-connect:8083/connectors/mysql-connector)" ]; then
            echo -e "Sending connector"
            curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" kafka-connect:8083/connectors/ -d @/conf/mysql_config.json
          else
            echo -e "Connector already exists"
          fi

# 定义Docker卷,用于持久化各服务数据
volumes:
  mysql_data:
  mysql_binlog:
  zookeeper_data:
  zookeeper_log:
  kafka_data:
  schema_registry_data:
  kafka_connect_plugins:
  kafka_ui_config:
  localstack_data:
关键改动说明
  • MySQL:添加mysql_data和mysql_binlog卷,分别持久化数据库数据和binlog文件,确保Debezium能读取历史变更记录
  • Zookeeper:添加zookeeper_data和zookeeper_log卷,持久化Zookeeper的元数据和日志,保证Kafka的集群状态不丢失
  • Kafka:添加kafka_data卷,持久化Kafka的消息数据,包括Connect的配置、偏移量和状态topic;同时指定KAFKA_LOG_DIRS配置明确数据存储路径
  • Kafka Connect:添加kafka_connect_plugins卷,持久化已安装的连接器插件,避免每次启动重复下载安装;修改启动命令,仅在插件目录为空时执行安装
  • Connector Sender:添加连接器存在性检查,避免容器重启后重复创建连接器
  • LocalStack:添加localstack_data卷,持久化S3存储的数据
  • 统一卷管理:在配置底部定义所有Docker卷,便于数据管理和备份

内容的提问来源于stack exchange,提问作者Varun Kumar Soni

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 00:54:51