Oracle到Postgres Kafka数据同步:插入正常但更新删除失效求助
问题分析与解决方案
核心问题原因
- 源端模式限制:当前用
mode: incrementing,仅能捕获新增数据,无法检测更新、删除——它只跟踪递增列(ID)的最大值,不会扫描已有数据的变更。 - 消息Key缺失主键:源端未把主键(ID)设为Kafka消息的Key,导致目标端
pk.mode: record_key找不到对应的键架构,无法匹配数据做更新/删除。 - 错误转换配置:目标端用了Debezium的
ExtractNewRecordState转换,但源端是Confluent JDBC Source(非Debezium CDC连接器),该转换不适用于当前消息格式,会破坏结构。
具体修复步骤
1. 修改源端连接器配置
切换到支持更新的模式,并配置主键作为消息Key:
{ "name": "source", "config": { "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector", "connection.url": "jdbc:oracle:thin:@192.168.91.139:1521/orcl1", "connection.user": "sys as sysdba", "connection.password": "oracle", "topic.prefix": "person", "mode": "timestamp+incrementing", // 替换为支持更新的模式 "poll.interval.ms": "1000", "incrementing.column.name": "ID", "timestamp.column.name": "UPDATE_TIME", // 需确保Oracle表存在记录更新时间的字段,无则新增 "table.whitelist": "person", // 替换原query配置,确保主键规则生效 "numeric.mapping": "none", "include.schema.changes": "true", "validate.non.null": "false", "value.converter.schemas.enable": "true", "key.converter": "io.confluent.connect.avro.AvroConverter", "key.converter.schema.registry.url": "http://localhost:8081", "value.converter": "io.confluent.connect.avro.AvroConverter", "value.converter.schema.registry.url": "http://localhost:8081", "pk.mode": "record_key", // 指定主键从消息Key读取 "pk.fields": "ID" // 设置主键字段为ID } }
2. 修改目标端连接器配置
移除不适用的Debezium转换,明确主键配置:
{ "name": "jdbc-sink", "config": { "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "tasks.max": "1", "topics": "person", "connection.url": "jdbc:postgresql://192.168.91.229:5432/postgres?user=postgres&password=postgres ", "auto.create": "true", "insert.mode": "upsert", "pk.mode": "record_key", "pk.fields": "ID", // 明确指定目标表主键字段 "delete.enabled": "true", "value.converter.schemas.enable": "true", "key.converter": "io.confluent.connect.avro.AvroConverter", "key.converter.schema.registry.url": "http://localhost:8081", "value.converter": "io.confluent.connect.avro.AvroConverter", "value.converter.schema.registry.url": "http://localhost:8081" } }
关于删除操作的补充说明
Confluent JDBC Source是轮询机制,无法捕获删除操作——只能通过时间戳/递增列检测新增和更新。如果需要同步删除,建议替换为Debezium Oracle CDC连接器,它读取Oracle的redo log捕获所有变更(包括删除),更适合实时CDC场景。
内容的提问来源于stack exchange,提问作者saad
相关产品推荐
相关产品推荐

