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

无主键场景下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"
    }
}
正确配置方案

核心问题分析

  1. 错误使用Debezium的ExtractNewRecordState转换:该转换仅适配Debezium格式的消息,而SharePlex的消息结构为meta+data的自定义格式,完全不兼容。
  2. delete.enabled参数值错误:合法值为true/false,而非none。
  3. 未过滤非数据操作消息:topic中包含schema、commit等非业务数据消息,会干扰Sink的正常写入逻辑。
  4. 无主键场景下的删除操作限制: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执行精准删除,可采用以下两种方案:

  1. 设置复合主键:将所有字段设置为复合主键,需调整以下配置项:

    "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的情况),复合主键会导致删除失败,需提前处理这类数据。

  2. 软删除策略:将del操作转换为插入或更新一条标记为删除的记录,需自定义转换逻辑或在数据库层面通过触发器等方式处理。

额外注意事项

  • 确保Kafka Connect环境已安装jsonscript依赖(用于Filter转换的条件判断),若未安装,可替换为RegexRouter或自定义Transform过滤消息。
  • 测试时可先使用kafka-console-consumer查看经过转换后的消息结构,确认是否符合JDBC Sink的写入要求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 11:01:02