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

Kafka Oracle JDBC Sink连接器主键列冗余数据处理咨询

问题根源

你设置了pk.mode: record_key,但Debezium生成的Record Key是包含schema和payload的结构化对象,而你的key.converter用的是StringConverter,导致整个结构化对象被序列化为字符串写入主键列,而非提取其中的ID数值。


解决方案1(推荐,配置更简洁)

修改主键模式为从处理后的Value中获取ID——你已经通过ExtractNewRecordState将Value转换成了纯业务数据结构,直接用这里的ID即可:

  1. 将配置中的pk.mode从record_key修改为record_value
  2. 保留现有transforms.unwrap的配置(它已经帮你把Value中的payload提取出来,里面包含纯数值的ID字段)

修改后的关键配置片段:

"pk.fields": "ID",
"pk.mode": "record_value",
"transforms": "unwrap",
"transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
"transforms.unwrap.drop.tombstones": "false"

解决方案2(若必须使用Record Key作为主键)

如果业务要求必须从Record Key获取主键,需要调整Key的转换器并添加转换逻辑提取ID:

  1. 将key.converter从org.apache.kafka.connect.storage.StringConverter改为org.apache.kafka.connect.json.JsonConverter
  2. 设置key.converter.schemas.enable为true(因为Debezium的Record Key包含Schema信息)
  3. 添加额外的Transform来提取Key中payload里的ID字段

修改后的关键配置片段:

"key.converter": "org.apache.kafka.connect.json.JsonConverter",
"key.converter.schemas.enable": "true",
"transforms": "unwrap,extractKeyId",
"transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
"transforms.unwrap.drop.tombstones": "false",
"transforms.extractKeyId.type": "org.apache.kafka.connect.transforms.ExtractField$Key",
"transforms.extractKeyId.field": "ID"

验证

修改配置后重启JDBC Sink连接器,插入或更新数据后,目标库的ID列会显示纯数值2,而非结构化字符串。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 15:26:15