使用Debezium与JDBC Sink同步Postgres时目标表未创建问题
Postgres跨库同步:目标表未自动创建的修复方案
核心问题定位
你的配置存在两个关键错误,导致JDBC Sink无法自动创建目标表:
- Topic名称不匹配:Debezium Postgres连接器默认生成的Topic格式为
{database.server.name}.{database.dbname}.{schema}.{table},你的场景下应为dbserver1.kafka_test.public.events,但Sink配置的topics字段写的是events,Sink根本监听不到数据,自然不会触发表创建。 - 无数据触发同步:JDBC Sink的
auto.create仅在有数据流入Topic时才会执行表创建操作,若源表无变更,Sink不会主动建表。
修复步骤
1. 修正Sink连接器的Topic配置
将Sink的topics字段修改为Debezium生成的正确Topic名称:
{"name": "lead-sink", "config": {"connector.class":"io.confluent.connect.jdbc.JdbcSinkConnector", "tasks.max":"10", "topics": "dbserver1.kafka_test.public.events", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "value.converter": "io.confluent.connect.avro.AvroConverter", "value.converter.schema.registry.url": "http://schema-registry:8081", "connection.url": "jdbc:postgresql://postgres:5432/kafka_test?user=postgres&password=postgres", "key.converter.schemas.enable": "false", "value.converter.schemas.enable": "true", "auto.create": "true", "auto.evolve": "true", "insert.mode": "upsert", "pk.fields": "identity_id", "pk.mode": "record_key" } }
2. (可选)自定义Topic名称
若希望Sink使用简洁的events作为Topic名,可在源连接器中添加Topic重命名转换:
{"name": "postgres-source", "config": {"connector.class":"io.debezium.connector.postgresql.PostgresConnector", "tasks.max":"1", "database.hostname": "postgres", "database.port": "5432", "database.user": "postgres", "database.password": "postgres", "database.dbname" : "kafka_test", "database.server.name": "dbserver1", "database.whitelist": "kafka_test", "database.history.kafka.bootstrap.servers": "kafka:9092", "database.history.kafka.topic": "schema-changes.kafka_test", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "value.converter": "io.confluent.connect.avro.AvroConverter", "key.converter.schemas.enable": "false", "value.converter.schemas.enable": "true", "value.converter.schema.registry.url": "http://schema-registry:8081", "transforms": "unwrap,renameTopic", "transforms.unwrap.type": "io.debezium.transforms.UnwrapFromEnvelope", "transforms.renameTopic.type": "org.apache.kafka.connect.transforms.RegexRouter", "transforms.renameTopic.regex": "dbserver1.kafka_test.public.(.*)", "transforms.renameTopic.replacement": "$1" } }
修改后Sink的topics字段保持events即可。
3. 触发数据同步
在源库执行插入操作,触发Debezium捕获变更并发送到Topic:
INSERT INTO events(identity_id, name, family_name, event_type) VALUES('test_001', '张三', '张', 'login');
4. 验证数据库权限
确认Sink使用的数据库账号拥有创建表的权限,避免因权限不足导致表创建失败。
验证方法
- 重启修改后的源、Sink连接器
- 登录目标库,执行
\dt命令查看是否存在events表 - 查询表数据:
SELECT * FROM events;,确认数据已同步
内容的提问来源于stack exchange,提问作者hishAg
相关产品推荐
相关产品推荐

