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
相关产品推荐
相关产品推荐

