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

Oracle到Postgres Kafka数据同步:插入正常但更新删除失效求助

问题分析与解决方案

核心问题原因

  1. 源端模式限制:当前用mode: incrementing,仅能捕获新增数据,无法检测更新、删除——它只跟踪递增列(ID)的最大值,不会扫描已有数据的变更。
  2. 消息Key缺失主键:源端未把主键(ID)设为Kafka消息的Key,导致目标端pk.mode: record_key找不到对应的键架构,无法匹配数据做更新/删除。
  3. 错误转换配置:目标端用了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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 13:03:32