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

Kafka Connect序列化错误排查:Debezium源到JDBC Sink报错

Debezium Oracle源连接器对接JDBC Sink序列化错误排查

问题描述

使用Debezium Oracle源连接器同步数据到PostgreSQL时,出现「Converting byte[] to Kafka Connect data failed due to serialization error」错误,相关配置及错误栈如下:

1. Connect全局配置(connect-standalone-file.properties)

bootstrap.servers=localhost:9092

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=true

offset.storage.file.filename=/tmp/connect.offsets

plugin.path=/dir/kafka_2.12-3.2.3/libs

2. Debezium Oracle源连接器配置

name=inventory-connector-test31
connector.class=io.debezium.connector.oracle.OracleConnector
tasks.max=1
database.server.name=server1
database.hostname=***.**.*.*
database.port=1521
database.user=username
database.password=password
database.dbname=dbname
database.pdb.name=ORCLPDB1
database.connection.adapter=logminer
database.history.kafka.bootstrap.servers=localhost:9092
database.history.kafka.topic=schema-changes.inventory
#snapshot.mode=SCHEMA_ONLY_RECOVERY
transforms=filter,route
transforms.filter.type=io.debezium.transforms.Filter
transforms.filter.language=jsr223.groovy
transforms.filter.condition=value.source.table == 'CUSTOMERS' && ((value.before != null && value.before.ID >= 1000 && value.before.ID <= 2000) || (value.after != null && value.after.ID >= 1000 && value.after.ID <= 2000))
transforms.filter.topic.regex=server1.DEBEZIUM.*
transforms.route.type=org.apache.kafka.connect.transforms.RegexRouter
transforms.route.regex=([^.]+)\.([^.]+)\.([^.]+)
transforms.route.replacement=$3

3. JDBC Sink连接器配置

name=students-sink14
connector.class=io.confluent.connect.jdbc.JdbcSinkConnector
tasks.max=1
topics=CUSTOMERS
connection.url=jdbc:postgresql://ipaddr:5432/productorcl?user=username&password=pass
transforms=unwrap
transforms.unwrap.type=io.debezium.transforms.ExtractNewRecordState
transforms.unwrap.drop.tombstones=false
auto.create=true
insert.mode=upsert
delete.enabled=true
pk.mode=record_key

完整错误栈

org.apache.kafka.connect.errors.ConnectException: Tolerance exceeded in error handler
    at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:223)
    at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execute(RetryWithToleranceOperator.java:149)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.convertAndTransformRecord(WorkerSinkTask.java:513)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.convertMessages(WorkerSinkTask.java:493)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:332)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:234)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:203)
    at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:188)
    at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:243)
    at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
    at java.util.concurrent.FutureTask.run(FutureTask.java:266)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:750)
Caused by: org.apache.kafka.connect.errors.DataException: Converting byte[] to Kafka Connect data failed due to serialization error: 
    at org.apache.kafka.connect.json.JsonConverter.toConnectData(JsonConverter.java:324)
    at org.apache.kafka.connect.storage.Converter.toConnectData(Converter.java:88)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.lambda$convertAndTransformRecord$3(WorkerSinkTask.java:513)
    at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndRetry(RetryWithToleranceOperator.java:173)
    at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:207)
    ... 13 more
Caused by: org.apache.kafka.common.errors.SerializationException: com.fasterxml.jackson.core.JsonParseException: Unrecognized token 'Struct': was expecting (JSON String, Number, Array, Object or token 'null', 'true' or 'false')
 at [Source: (byte[])"Struct{ID=1001}"; line: 1, column: 8]
    at org.apache.kafka.connect.json.JsonDeserializer.deserialize(JsonDeserializer.java:66)
    at org.apache.kafka.connect.json.JsonConverter.toConnectData(JsonConverter.java:322)
    ... 17 more
Caused by: com.fasterxml.jackson.core.JsonParseException: Unrecognized token 'Struct': was expecting (JSON String, Number, Array, Object or token 'null', 'true' or 'false')
 at [Source: (byte[])"Struct{ID=1001}"; line: 1, column: 8]
    at com.fasterxml.jackson.core.JsonParser._constructError(JsonParser.java:2391)
    at com.fasterxml.jackson.core.base.ParserMinimalBase._reportError(ParserMinimalBase.java:745)
    at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._reportInvalidToken(UTF8StreamJsonParser.java:3635)
    at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._handleUnexpectedValue(UTF8StreamJsonParser.java:2734)
    at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._nextTokenNotInObject(UTF8StreamJsonParser.java:902)
    at com.fasterxml.jackson.core.json.UTF8StreamJsonParser.nextToken(UTF8StreamJsonParser.java:794)
    at com.fasterxml.jackson.databind.ObjectMapper._readTreeAndClose(ObjectMapper.java:4703)
    at com.fasterxml.jackson.databind.ObjectMapper.readTree(ObjectMapper.java:3090)
    at org.apache.kafka.connect.json.JsonDeserializer.deserialize(JsonDeserializer.java:64)
    ... 18 more

错误原因分析

从错误栈里的Source: (byte[])"Struct{ID=1001}"可以明确看出:

  • Debezium生成的消息Key是Struct格式(包含表的主键字段)
  • 全局配置中key.converter.schemas.enable=false,导致JsonConverter将Struct直接序列化为其toString形式(即Struct{ID=1001}),而非标准JSON格式
  • Sink端使用JsonConverter反序列化Key时,无法识别这种非JSON字符串,因此抛出解析错误

另外,Sink配置中pk.mode=record_key意味着需要从消息Key中提取主键,这进一步依赖Key的正确序列化/反序列化。

解决方案

方案1:修改全局配置,开启Key的Schema支持

修改connect-standalone-file.properties中的Key转换器配置:

key.converter.schemas.enable=true

这样Key会被序列化为带Schema的JSON格式,Sink端可以正确解析。

方案2:仅在Sink连接器中单独配置Key转换器

如果不想修改全局配置,可以在Sink连接器配置中添加单独的Key转换器设置:

key.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=true

确保Sink端使用的Key转换器与源端的序列化格式匹配。

方案3:将Key转换为基本类型(可选)

如果主键是单个字段(比如ID),可以在源端添加ExtractField Transform,将Key从Struct提取为单个值:
在源连接器配置中添加:

transforms=filter,route,extractKey
transforms.extractKey.type=org.apache.kafka.connect.transforms.ExtractField$Key
transforms.extractKey.field=ID

此时Key会变成单个数值,即使key.converter.schemas.enable=false也能序列化为JSON格式的数值,Sink端可以正常解析。

内容的提问来源于stack exchange,提问作者Muhammad Affan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 10:34:53