基于Strimzi的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
咨询问题
- 适配该纯JSON的Schema结构应该是什么样的?
- 应该选用Avro Converter还是JsonSchema Converter?
- 问题根源是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

