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

基于Strimzi的Kafka JDBC Sink连接器对接Oracle技术咨询

Kafka JDBC Sink 同步 Oracle 问题排查

背景

在Strimzi上搭建Kafka集群,需通过JDBC Sink连接器将test1主题的纯JSON数据同步至Oracle数据库。主题数据格式固定为:

{"ID":"348010961","TRANSACTION_REQUEST_TYPE_ID":"111"}

无法控制生产者逻辑,只能处理现有数据。已知JDBC Sink连接器写入关系型数据库需要Schema支持,因此部署了Schema Registry并完成与Broker的连接,但不确定适配该JSON的Schema结构,也纠结于选择Avro Converter还是JsonSchema Converter。

当前连接器配置

spec:
  class: io.confluent.connect.jdbc.JdbcSinkConnector
  config:
    value.converter.schema.registry.url: 'http://schema-registry-test:8081'
    value.converter: io.confluent.connect.json.JsonSchemaConverter
    key.converter: io.confluent.connect.json.JsonSchemaConverter
    topics: test1
    value.converter.schema.registry.version: 1
    value.converter.schema.registry.id: 1
    value.converter.schemas.enable: true
    key.converter.schema.registry.subject: my-schema-value
    connection.password: 'xxxx'
    pk.fields: ID
    key.converter.schema.registry.url: 'http://schema-registry-test:8081'
    pk.mode: record_value
    tasksMax: 1
    insert.mode: insert
    connection.user: xxx
    auto.create: true
    value.converter.schema.registry.subject: my-schema-value
    connection.url: 'xxxx'

日志报错

1. 连接器日志

yourorg.apache.kafka.connect.errors.DataException: Converting byte[] to Kafka Connect data failed due to serialization error of topic test1:  at io.confluent.connect.json.JsonSchemaConverter.toConnectData(JsonSchemaConverter.java:144) at org.apache.kafka.connect.runtime.WorkerSinkTask.convertValue(WorkerSinkTask.java:540) at org.apache.kafka.connect.runtime.WorkerSinkTask.lambda$convertAndTransformRecord$2(WorkerSinkTask.java:496) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndRetry(RetryWithToleranceOperator.java:156) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:190) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execute(RetryWithToleranceOperator.java:132) at org.apache.kafka.connect.runtime.WorkerSinkTask.convertAndTransformRecord(WorkerSinkTask.java:496) at org.apache.kafka.connect.runtime.WorkerSinkTask.convertMessages(WorkerSinkTask.java:473) at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:328) at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:232) at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:201) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:186) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:241) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829) Caused by: org.apache.kafka.common.errors.SerializationException: Error deserializing JSON message for id -1 at io.confluent.kafka.serializers.json.AbstractKafkaJsonSchemaDeserializer.deserialize(AbstractKafkaJsonSchemaDeserializer.java:236) at io.confluent.kafka.serializers.json.AbstractKafkaJsonSchemaDeserializer.deserializeWithSchemaAndVersion(AbstractKafkaJsonSchemaDeserializer.java:313) at io.confluent.connect.json.JsonSchemaConverter$Deserializer.deserialize(JsonSchemaConverter.java:193) at io.confluent.connect.json.JsonSchemaConverter.toConnectData(JsonSchemaConverter.java:127) ... 17 more text Caused by: org.apache.kafka.common.errors.SerializationException: Unknown magic byte!

无论是否注册Schema,该错误都会出现。

2. Schema Registry日志

[2023-07-17 17:41:17,794] INFO [Consumer clientId=KafkaStore-reader-_schemas, groupId=schema-registry-schema-registry-8081] Node -1 disconnected. (org.apache.kafka.clients.NetworkClient) [2023-07-17 17:41:18,516] INFO [Schema registry clientId=sr-1, groupId=schema-registry] Node -1 disconnected. (org.apache.kafka.clients.NetworkClient) [2023-07-17 17:41:40,576] INFO [Producer clientId=producer-1] Node -1 disconnected. (org.apache.kafka.clients.NetworkClient)
Dynamic member with unknown member id joins group schema-registry in Empty state. Created a new member id sr-1-8343ba38-46a2-4b36-81ec-ddb874b30ae5 and request the member to rejoin with this id. (kafka.coordinator.group.GroupCoordinator) [data-plane-kafka-request-handler-3] 2023-07-17 17:32:18,443 INFO [GroupCoordinator 0]: Preparing to rebalance group schema-registry in state PreparingRebalance with old generation 27 (__consumer_offsets-29) (reason: Adding new member sr-1-8343ba38-46a2-4b36-81ec-ddb874b30ae5 with group instance id None; client reason: rebalance failed due to MemberIdRequiredException

咨询问题

  1. 适配该纯JSON的Schema结构应该是什么样的?
  2. 应该选用Avro Converter还是JsonSchema Converter?
  3. 问题根源是Schema本身,还是连接器访问Schema Registry时存在其他故障?

解答

1. 适配的Schema结构

如果一定要用Schema Registry,针对你的JSON数据,对应的JSON Schema结构如下(注册到Schema Registry时使用):

{
  "$schema": "http://json-schema.org/draft-07/schema#",
  "type": "object",
  "title": "TransactionRecord",
  "properties": {
    "ID": {
      "type": "string"
    },
    "TRANSACTION_REQUEST_TYPE_ID": {
      "type": "string"
    }
  },
  "required": ["ID"]
}

若后续需要将ID转为数值类型,可将type改为integer,但需确保生产者发送的数据类型匹配。

2. Converter选择

你不需要使用依赖Schema Registry的Converter(Avro或JsonSchema Converter),因为生产者发送的是无Schema的原生JSON,而这类Converter要求消息是经过Schema Registry序列化的格式(包含magic byte和schema ID),这正是你报错Unknown magic byte!的核心原因。

正确的选择是使用原生JSON Converter:org.apache.kafka.connect.json.JsonConverter,并关闭Schema自动启用,配合JDBC Sink的自身能力即可完成同步。修改后的核心配置如下:

key.converter: org.apache.kafka.connect.json.JsonConverter
value.converter: org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable: false
value.converter.schemas.enable: false

结合你已配置的auto.create: true和pk.mode: record_value,只要JSON字段名与Oracle表字段名匹配,即可自动创建表并写入数据。

3. 问题根源分析

  • 连接器的Unknown magic byte!错误:完全是Converter配置错误导致,和Schema本身无关。你用了依赖Schema Registry的Converter,但消息是原生JSON,格式不匹配导致解析失败。
  • Schema Registry的日志问题:显示Node -1 disconnected和组重平衡错误,说明Schema Registry与Kafka Broker之间存在连接故障,可能是Broker地址配置错误、网络不通或Broker集群状态异常。这是独立于连接器序列化错误的另一个问题,需要先修复Schema Registry的Broker连接问题,再调整Converter配置。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 04:35:05