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
相关产品推荐
相关产品推荐

