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

PostgreSQL到MySQL流式同步:Debezium配置问题求助

问题分析与解决方案

问题1:PostgreSQL变更无法同步到MySQL

以下是你的配置中存在的核心问题及修复方法:

1. Topic名称不匹配

PostgreSQL源连接器配置了database.server.name: psql和topic.prefix: psql,因此生成的变更事件topic为psql.public.customers,但JDBC Sink连接器的topics参数配置为postgres.public.customers,导致Sink无法接收任何事件。
修复:将Sink配置的topics修改为psql.public.customers。

2. PostgreSQL源连接器的无效参数

配置中的mode: bulk是JDBC源连接器的专属参数,Debezium PostgreSQL Connector不支持该参数,保留会导致连接器初始化失败。
修复:直接删除这一行配置。

3. 主键与Upsert配置校验

确保PostgreSQL的customers表主键为id,当前Sink配置的pk.fields: id和pk.mode: record_key是正确的。若发现主键未被正确提取,可添加transforms.unwrap.add.fields: id确保主键被包含在处理后的记录中。

4. Avro转换器配置修正

若使用本地Docker环境,Schema Registry通常使用HTTP而非HTTPS,需将key.converter.schema.registry.url和value.converter.schema.registry.url修改为http://schema-registry:8081(假设Schema Registry服务名为schema-registry),并确保用户名密码配置正确。


问题2:Docker启动时Kafka DNS解析失败

该错误是因为Connect容器无法通过DNS解析kafka主机名,修复方法如下:

1. 统一Docker网络

确保Kafka、Connect、PostgreSQL、MySQL等服务都加入同一个自定义Docker网络,示例docker-compose配置片段:

networks:
  debezium-network:
    driver: bridge

services:
  kafka:
    # ...其他配置
    networks:
      - debezium-network
  connect:
    # ...其他配置
    networks:
      - debezium-network

2. 校验Kafka服务名

确认docker-compose中Kafka的服务名确实为kafka,若服务名为其他(如kafka-broker),需修改Connect配置中的bootstrap.servers为kafka-broker:9092。

3. 调整Bootstrap Servers配置

在Connect服务的环境变量中设置正确的BOOTSTRAP_SERVERS,示例:

connect:
  environment:
    BOOTSTRAP_SERVERS: kafka:9092
    # ...其他环境变量

若使用Docker Desktop本地测试,也可尝试替换为host.docker.internal:9092访问宿主机的Kafka服务。


修复后的完整配置示例

PostgreSQL源连接器配置

{
"name": "inventory-connector",
"config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "topic.prefix": "psql",
    "database.hostname": "postgres",
    "database.port": "5432",
    "database.user": "postgresuser",
    "database.password": "postgrespw",
    "database.dbname": "inventory",
    "table.include.list": "public.customers",
    "slot.name": "test_slot",
    "plugin.name": "wal2json",
    "database.server.name": "psql",
    "tombstones.on.delete": "true",
    "key.converter": "io.confluent.connect.avro.AvroConverter",
    "key.converter.schema.registry.url": "http://schema-registry:8081",
    "key.converter.basic.auth.credentials.source": "USER_INFO",
    "key.converter.schema.registry.basic.auth.user.info": "[SCHEMA_REGISTRY_USER]:[SCHEMA_REGISTRY_PASSWORD]",
    "value.converter": "io.confluent.connect.avro.AvroConverter",
    "value.converter.schema.registry.url": "http://schema-registry:8081",
    "value.converter.basic.auth.credentials.source": "USER_INFO",
    "value.converter.schema.registry.basic.auth.user.info": "[SCHEMA_REGISTRY_USER]:[SCHEMA_REGISTRY_PASSWORD]"
}
}

JDBC Sink连接器配置

{
"name": "jdbc-sink",
"config": {
    "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
    "tasks.max": "1",
    "topics": "psql.public.customers",
    "connection.url": "jdbc:mysql://mysql:3306/inventory",
    "connection.user": "debezium",
    "connection.password": "dbz",
    "transforms": "unwrap",
    "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
    "transforms.unwrap.drop.tombstones": "false",
    "auto.create": "true",
    "auto.evolve": "true",
    "insert.mode": "upsert",
    "delete.enabled": "true",
    "pk.fields": "id",
    "pk.mode": "record_key"
}
}

验证步骤

  1. 启动所有服务后,检查连接器状态:
    curl -X GET http://localhost:8083/connectors/inventory-connector/status
    curl -X GET http://localhost:8083/connectors/jdbc-sink/status
    
  2. 在PostgreSQL的customers表中插入/更新/删除记录,检查MySQL的inventory.customers表是否同步变更。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 14:54:09