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" } }
验证步骤
- 启动所有服务后,检查连接器状态:
curl -X GET http://localhost:8083/connectors/inventory-connector/status curl -X GET http://localhost:8083/connectors/jdbc-sink/status - 在PostgreSQL的
customers表中插入/更新/删除记录,检查MySQL的inventory.customers表是否同步变更。
内容的提问来源于stack exchange,提问作者Alex

