Kafka Oracle JDBC Sink连接器主键列冗余数据处理咨询
问题根源
你设置了pk.mode: record_key,但Debezium生成的Record Key是包含schema和payload的结构化对象,而你的key.converter用的是StringConverter,导致整个结构化对象被序列化为字符串写入主键列,而非提取其中的ID数值。
解决方案1(推荐,配置更简洁)
修改主键模式为从处理后的Value中获取ID——你已经通过ExtractNewRecordState将Value转换成了纯业务数据结构,直接用这里的ID即可:
- 将配置中的
pk.mode从record_key修改为record_value - 保留现有
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:
- 将
key.converter从org.apache.kafka.connect.storage.StringConverter改为org.apache.kafka.connect.json.JsonConverter - 设置
key.converter.schemas.enable为true(因为Debezium的Record Key包含Schema信息) - 添加额外的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
相关产品推荐
相关产品推荐

