Kafka同步Oracle至PostgreSQL时主键NULL值冲突问题求助
解决方案
1. 修正源连接器的主键传递配置
问题核心是源连接器未将主键ID写入Kafka消息的Key中,导致Sink连接器从Key中获取到NULL作为主键。调整源连接器配置如下:
- 添加
"pk.mode": "record_key",让源连接器将主键字段放入消息Key - 将
key.converter改为与value.converter一致的Avro转换器,保证主键类型正确传递 - 开启
validate.non.null过滤源表中主键为NULL的记录,避免无效数据流入Kafka
修改后的源连接器配置片段:
{ "name": "source", "config": { // 保留原有其他配置 "pk.mode": "record_key", "validate.non.null": "true", "key.converter":"io.confluent.connect.avro.AvroConverter", "key.converter.schema.registry.url":"http://localhost:8081" } }
2. 调整Sink连接器的主键读取逻辑(源端无法修改时可选)
如果源端配置暂时不能变更,可修改Sink连接器的主键读取方式,直接从消息Value中获取ID:
"pk.mode": "record_value"
同时建议手动在Postgres中创建目标表,明确主键非空约束:
CREATE TABLE person ( ID INT NOT NULL PRIMARY KEY, NAME VARCHAR(255) );
3. 使用Transform过滤无效记录
在Sink连接器中添加Transform插件,直接过滤掉ID为NULL的记录:
"transforms": "FilterNullID", "transforms.FilterNullID.type": "org.apache.kafka.connect.transforms.Filter$Value", "transforms.FilterNullID.filter.condition": "$.ID IS NOT NULL", "transforms.FilterNullID.filter.type": "include"
4. 源表数据校验
检查Oracle源表person的ID字段是否为非空主键,若未设置,在源端修正表结构:
ALTER TABLE person MODIFY ID NOT NULL; ALTER TABLE person ADD CONSTRAINT pk_person_id PRIMARY KEY (ID);
内容的提问来源于stack exchange,提问作者saad
相关产品推荐
相关产品推荐

