如何在Confluent平台运行Debezium?CDC任务超时排查
Debezium MySQL CDC连接器任务超时获取Kafka元数据故障排查
问题背景
EC2实例上部署了非容器化MySQL,通过Docker Compose启动Kafka Connect平台连接Confluent托管的Kafka集群。Connect本身运行正常,配置、偏移量、状态元数据主题已正常填充,但创建Debezium MySQL CDC连接器后,任务启动失败,报错:
{"name":"mysql-connector","connector":{"state":"RUNNING","worker_id":"connect:8083"},"tasks":[{"id":0,"state":"FAILED","worker_id":"connect:8083","trace":"org.apache.kafka.common.errors.TimeoutException: Timeout expired while fetching topic metadata\n"}],"type":"source"}
已排除凭证问题(错误凭证会导致Connect启动阶段直接报错),聚焦Docker网络与连接器配置问题。
现有配置
Docker Compose文件
services: connect: image: confluentinc/cp-kafka-connect:latest network_mode: host env_file: - .env ports: - "8083:8083" command: - bash - -c - | echo "Installing connector plugins" confluent-hub install --no-prompt debezium/debezium-connector-mysql:2.4.2 echo "Launching Kafka Connect worker" /etc/confluent/docker/run & sleep infinity environment: CONNECT_BOOTSTRAP_SERVERS: "${KAFKA_BOOTSTRAP_SERVER}" CONNECT_TOPIC_CREATION_ENABLE: "false" CONNECT_SECURITY_PROTOCOL: "SASL_SSL" CONNECT_SASL_MECHANISM: "PLAIN" CONNECT_SASL_USERNAME: "${KAFKA_API_KEY}" CONNECT_SASL_PASSWORD: "${KAFKA_API_SECRET}" CONNECT_SASL_JAAS_CONFIG: "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"${KAFKA_API_KEY}\" password=\"${KAFKA_API_SECRET}\";" CONNECT_SSL_ENDPOINT_IDENTIFICATION_ALGORITHM: "HTTPS" CONNECT_REQUEST_TIMEOUT_MS: "20000" CONNECT_RETRY_BACKOFF_MS: "500" CONNECT_CONSUMER_SASL_JAAS_CONFIG: "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"${KAFKA_API_KEY}\" password=\"${KAFKA_API_SECRET}\";" CONNECT_CONSUMER_SECURITY_PROTOCOL: "SASL_SSL" CONNECT_CONSUMER_SASL_MECHANISM: "PLAIN" CONNECT_CONSUMER_REQUEST_TIMEOUT_MS: "20000" CONNECT_CONSUMER_RETRY_BACKOFF_MS: "500" CONNECT_CONSUMER_SSL_ENDPOINT_IDENTIFICATION_ALGORITHM: "https" CONNECT_PRODUCER_SASL_JAAS_CONFIG: "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"${KAFKA_API_KEY}\" password=\"${KAFKA_API_SECRET}\";" CONNECT_PRODUCER_SECURITY_PROTOCOL: "SASL_SSL" CONNECT_PRODUCER_SASL_MECHANISM: "PLAIN" CONNECT_PRODUCER_REQUEST_TIMEOUT_MS: "20000" CONNECT_PRODUCER_RETRY_BACKOFF_MS: "500" CONNECT_PRODUCER_SSL_ENDPOINT_IDENTIFICATION_ALGORITHM: "https" CONNECT_REST_ADVERTISED_HOST_NAME: connect CONNECT_GROUP_ID: 1 CONNECT_OFFSET_FLUSH_INTERVAL_MS: 10000 CONNECT_CONFIG_STORAGE_TOPIC: "${CONFIG_STORAGE_TOPIC}" CONNECT_OFFSET_STORAGE_TOPIC: "${OFFSET_STORAGE_TOPIC}" CONNECT_STATUS_STORAGE_TOPIC: "${STATUS_STORAGE_TOPIC}" CONNECT_KEY_CONVERTER: "org.apache.kafka.connect.storage.StringConverter" CONNECT_VALUE_CONVERTER: "org.apache.kafka.connect.json.JsonConverter" CONNECT_CONFIG_STORAGE_REPLICATION_FACTOR: '3' CONNECT_OFFSET_STORAGE_REPLICATION_FACTOR: '3' CONNECT_STATUS_STORAGE_REPLICATION_FACTOR: '3' CONNECT_CONFIG_PROVIDERS: 'file' CONNECT_CONFIG_PROVIDERS_FILE_CLASS: 'org.apache.kafka.common.config.provider.FileConfigProvider' CONNECT_LOG4J_ROOT_LOGLEVEL: 'INFO' volumes: - .env:/data/credentials.properties
CDC连接器配置
{ "name": "mysql-connector", "config": { "connector.class": "io.debezium.connector.mysql.MySqlConnector", "database.hostname": "localhost", "database.port": "3306", "database.user": "debezium", "database.password": "debezium", "database.server.id": 1, "database.include.list": "mydb", "snapshot.mode": "schema_only", "topic.prefix": "analytics01", "table.include.list": "mydb.table1,mydb.table2", "schema.history.internal.kafka.bootstrap.servers": "${file:/data/credentials.properties:KAFKA_BOOTSTRAP_SERVER}", "schema.history.internal.kafka.topic": "${file:/data/credentials.properties:SERVER_NAME}.schemaChanges", "schema.history.consumer.security.protocol": "SASL_SSL", "schema.history.consumer.ssl.endpoint.identification.algorithm": "https", "schema.history.consumer.sasl.mechanism": "PLAIN", "schema.history.consumer.sasl.jaas.config": "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"${file:/data/credentials.properties:KAFKA_API_KEY}\" password=\"${file:/data/credentials.properties:KAFKA_API_SECRET}\";", "schema.history.producer.security.protocol": "SASL_SSL", "schema.history.producer.ssl.endpoint.identification.algorithm": "https", "schema.history.producer.sasl.mechanism": "PLAIN", "schema.history.producer.sasl.jaas.config": "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"${file:/data/credentials.properties:KAFKA_API_KEY}\" password=\"${file:/data/credentials.properties:KAFKA_API_SECRET}\";", "decimal.handling.mode":"double", "transforms": "InsertField", "transforms.InsertField.type": "org.apache.kafka.connect.transforms.InsertField$Value", "transforms.InsertField.static.field": "licenseId", "transforms.InsertField.static.value": "${file:/data/credentials.properties:LICENSE_ID}" } }
.env文件内容
KAFKA_API_KEY=xxx KAFKA_API_SECRET=xxx KAFKA_BOOTSTRAP_SERVER=xxx.eu-central-1.aws.confluent.cloud:9092 LICENSE_ID=xxx CONFIG_STORAGE_TOPIC=xxx-config OFFSET_STORAGE_TOPIC=xxx-offset STATUS_STORAGE_TOPIC=xxx-status SERVER_NAME=analytics01
故障排查与解决方案
1. 手动创建Schema History主题
由于CONNECT_TOPIC_CREATION_ENABLE设为false,Debezium无法自动创建Schema History主题,这会导致任务启动时无法获取主题元数据:
- 在Confluent Cloud控制台手动创建主题
analytics01.schemaChanges,设置至少3个副本(符合Confluent Cloud要求),分区数可设为1。
2. 补充Schema History客户端超时配置
连接器的Schema History生产者/消费者缺少超时配置,容易触发超时异常,在连接器配置中添加:
"schema.history.producer.request.timeout.ms": "20000", "schema.history.producer.retry.backoff.ms": "500", "schema.history.consumer.request.timeout.ms": "20000", "schema.history.consumer.retry.backoff.ms": "500"
3. 修正Connect Worker的主机名配置
使用network_mode: host时,CONNECT_REST_ADVERTISED_HOST_NAME设为connect会导致内部解析失败,修改Docker Compose中的该配置为EC2实例内网IP或localhost:
CONNECT_REST_ADVERTISED_HOST_NAME: localhost
4. 验证配置变量解析有效性
创建连接器后,通过API查看实际生效的配置,确认${file:...}变量是否正确解析:
curl http://localhost:8083/connectors/mysql-connector/config
如果变量未解析,检查Connect容器日志,确认FileConfigProvider是否正常加载,或直接硬编码Schema History相关的Kafka配置进行测试。
5. 验证MySQL连接与Binlog配置
虽然报错指向Kafka,但需确保MySQL连接正常:
- 在Connect容器内执行
mysql -h localhost -P3306 -u debezium -pdebezium,确认能成功连接。 - 检查MySQL配置:确保
binlog_format=ROW,且server_id与连接器的database.server.id(1)不冲突(MySQL自身的server_id需设为不同值)。
内容的提问来源于stack exchange,提问作者m_m
相关产品推荐
相关产品推荐

