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
相关产品推荐
相关产品推荐

