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

Debezium-JDBC Sink Connector非主键列处理操作配置咨询

解决方案

针对你的场景,需要调整以下JDBC Sink连接器配置项,确保以record_id作为更新、删除的依据:

1. 修正主键识别配置

保持pk.mode: "record_value"(指定从消息value中提取主键字段),确认pk.fields: "record_id"(明确主键字段为record_id)。

注意:需保证Debezium同步的消息value中确实包含record_id字段,且值为源端主键id。如果当前消息里只有源端的id字段,需添加字段重命名转换,将源端id映射为record_id,配置如下:

"transforms": "unwrap,renameId",
"transforms.renameId.type": "org.apache.kafka.connect.transforms.ReplaceField$Value",
"transforms.renameId.renames": "id:record_id"

2. 调整删除事件处理逻辑

当前删除相关配置会导致删除操作失效,需修改:

  • 将delete.handling.mode从none改为rewrite:该配置会把删除事件转换为仅保留主键的UPSERT操作,触发目标库基于record_id的删除
  • 调整transforms.unwrap.delete.handling.mode为rewrite:让Debezium的unwrap转换将删除事件重写为包含record_id的消息,供Sink连接器识别
  • 保留delete.enabled: "true",开启删除支持

3. 清理干扰性配置

删除delete.tombstone.handling.mode: "drop"和transforms.unwrap.drop.tombstones: "false",避免墓碑消息干扰删除逻辑。

最终调整后的完整配置示例

{
    "name": "test_table_sink-connector",
    "config": {
        "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
        "tasks.max": "1",
        "topics": "utanga_dev.public.test_table",
        "key.converter": "io.confluent.connect.avro.AvroConverter",
        "value.converter": "io.confluent.connect.avro.AvroConverter",
        "value.converter.schema.registry.url": "http://host.docker.internal:8081",
        "key.converter.schema.registry.url": "http://host.docker.internal:8081",
        "connection.url": "jdbc:postgresql://db:5432/utanga_dev?user=utanga&password=changeme",
        "key.converter.schemas.enable": "true",
        "value.converter.schemas.enable": "true",
        "auto.create": "false",
        "auto.evolve": "true",
        "insert.mode": "upsert",
        "pk.mode": "record_value",
        "pk.fields": "record_id",
        "transforms": "unwrap,renameId",
        "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
        "transforms.unwrap.delete.handling.mode": "rewrite",
        "transforms.renameId.type": "org.apache.kafka.connect.transforms.ReplaceField$Value",
        "transforms.renameId.renames": "id:record_id",
        "table.name.format": "${topic}",
        "delete.enabled": "true",
        "delete.handling.mode": "rewrite"
    }
}

关键说明

  • 如果你的同步流程已经能自动将源端id填充到Sink端record_id,可去掉renameId相关转换配置
  • 测试时可先发送更新/删除消息,验证Kafka消息是否包含record_id,以及Sink端数据库是否执行了对应操作

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 16:33:13