无主键场景下SharePlex-Kafka-Oracle的JDBC Sink Connector配置问询
问题场景与错误说明
Oracle数据通过SharePlex同步至Kafka(源表无主键),需将topic sp_test的消息回写至Oracle。该topic的消息内容如下:
{"meta":{"op":"schema","table":""},"schema":{"name":"SPLEX.KAFKA","utcOffset":"9:00","NAME":{"jsonType":"string","num":1,"key":1,"nullable":1,"length":30,"precision":0,"scale":0,"src_name":"NAME"},"ADDRESS":{"jsonType":"string","num":2,"key":1,"nullable":1,"length":60,"precision":0,"scale":0,"src_name":"ADDRESS"},"PHONE":{"jsonType":"string","num":3,"key":1,"nullable":1,"length":12,"precision":0,"scale":0,"src_name":"PHONE"}}} {"meta":{"op":"del","table":"SPLEX.KAFKA"},"data":{"NAME":"111","ADDRESS":"222","PHONE":"333"}} {"meta":{"op":"del","table":"SPLEX.KAFKA"},"data":{"NAME":"444","ADDRESS":"444","PHONE":"444"}} {"meta":{"op":"del","table":"SPLEX.KAFKA"},"data":{"NAME":"aaa","ADDRESS":"aaa","PHONE":"aaa"}} {"meta":{"op":"del","table":"SPLEX.KAFKA"},"data":{"NAME":"bbb","ADDRESS":"bbb","PHONE":"bbb"}} {"meta":{"op":"del","table":"SPLEX.KAFKA"},"data":{"NAME":"eee","ADDRESS":"eee","PHONE":"eee"}} {"meta":{"op":"del","table":"SPLEX.KAFKA"},"data":{"NAME":"ttt","ADDRESS":"ttt","PHONE":"ttt"}} {"meta":{"op":"del","table":"SPLEX.KAFKA"},"data":{"NAME":null,"ADDRESS":"ccc","PHONE":"ccc"}} {"meta":{"op":"del","table":"SPLEX.KAFKA"},"data":{"NAME":"zzz","ADDRESS":"zzz","PHONE":"zzz"}} {"meta":{"op":"del","table":"SPLEX.KAFKA"},"data":{"NAME":"123","ADDRESS":"123","PHONE":"123"}} {"meta":{"op":"ins","table":"SPLEX.KAFKA"},"data":{"NAME":"aaa","ADDRESS":"aaa","PHONE":"aaa"}} {"meta":{"op":"commit","table":""}}
使用以下JDBC Sink Connector配置时返回错误{"error_code": 500,"message": null}:
{ "name": "oracle_jdbc_sink_customers_00", "config": { "connector.class":"io.confluent.connect.jdbc.JdbcSinkConnector", "tasks.max": "1", "topics": "sp_test", "connection.url": "jdbc:oracle:thin:@host:1521/host", "connection.user": "splex", "connection.password": "splex", "insert.mode": "insert", "table.name.format": "KAFKA_JSON_SINK", "auto.create": "true", "auto.evolve": "true", "pk.mode": "none", "delete.enabled": "none", "transforms": "unwrap", "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState", "transforms.unwrap.drop.tombstones": "false", "transforms.unwrap.delete.handling.mode": "rewrite", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "false" } }
正确配置方案
核心问题分析
- 错误使用Debezium的
ExtractNewRecordState转换:该转换仅适配Debezium格式的消息,而SharePlex的消息结构为meta+data的自定义格式,完全不兼容。 delete.enabled参数值错误:合法值为true/false,而非none。- 未过滤非数据操作消息:topic中包含
schema、commit等非业务数据消息,会干扰Sink的正常写入逻辑。 - 无主键场景下的删除操作限制:JDBC Sink需要主键才能精准执行删除,源表无主键时无法直接处理
del类型消息,需调整处理策略。
修正后的配置(仅处理插入操作)
如果仅需回写插入消息、忽略删除操作,可使用以下配置:
{ "name": "oracle_jdbc_sink_customers_00", "config": { "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "tasks.max": "1", "topics": "sp_test", "connection.url": "jdbc:oracle:thin:@host:1521/host", "connection.user": "splex", "connection.password": "splex", "insert.mode": "insert", "table.name.format": "KAFKA_JSON_SINK", "auto.create": "true", "auto.evolve": "true", "pk.mode": "none", "delete.enabled": "false", "transforms": "filter,extractData", // 过滤掉schema和commit类型的消息,只保留ins操作 "transforms.filter.type": "org.apache.kafka.connect.transforms.Filter$Value", "transforms.filter.condition": "$.meta.op == 'ins'", "transforms.filter.language": "jsonscript", // 提取data字段作为写入Oracle的内容 "transforms.extractData.type": "org.apache.kafka.connect.transforms.ExtractField$Value", "transforms.extractData.field": "data", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "false" } }
如需处理删除操作(无主键场景)
由于源表无主键,无法直接通过JDBC Sink执行精准删除,可采用以下两种方案:
设置复合主键:将所有字段设置为复合主键,需调整以下配置项:
"pk.mode": "record_value", "pk.fields": "NAME,ADDRESS,PHONE", "delete.enabled": "true", "insert.mode": "upsert", // 修改filter条件,同时保留ins和del操作 "transforms.filter.condition": "$.meta.op == 'ins' || $.meta.op == 'del'"注意:若字段存在
null值(如消息中NAME:null的情况),复合主键会导致删除失败,需提前处理这类数据。软删除策略:将
del操作转换为插入或更新一条标记为删除的记录,需自定义转换逻辑或在数据库层面通过触发器等方式处理。
额外注意事项
- 确保Kafka Connect环境已安装
jsonscript依赖(用于Filter转换的条件判断),若未安装,可替换为RegexRouter或自定义Transform过滤消息。 - 测试时可先使用
kafka-console-consumer查看经过转换后的消息结构,确认是否符合JDBC Sink的写入要求。
内容的提问来源于stack exchange,提问作者호이호이
相关产品推荐
相关产品推荐

