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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 05:04:50