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

使用Debezium与JDBC Sink同步Postgres时目标表未创建问题

Postgres跨库同步:目标表未自动创建的修复方案

核心问题定位

你的配置存在两个关键错误,导致JDBC Sink无法自动创建目标表:

  1. Topic名称不匹配:Debezium Postgres连接器默认生成的Topic格式为{database.server.name}.{database.dbname}.{schema}.{table},你的场景下应为dbserver1.kafka_test.public.events,但Sink配置的topics字段写的是events,Sink根本监听不到数据,自然不会触发表创建。
  2. 无数据触发同步: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使用的数据库账号拥有创建表的权限,避免因权限不足导致表创建失败。

验证方法

  1. 重启修改后的源、Sink连接器
  2. 登录目标库,执行\dt命令查看是否存在events表
  3. 查询表数据:SELECT * FROM events;,确认数据已同步

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 23:42:47