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

Kafka Sink Connector Avro数据反序列化失败问题求助

问题排查与解决方案

核心原因

你的Sink连接器配置了io.confluent.connect.avro.AvroConverter作为value转换器,但写入Kafka Topic的是纯JSON字符串,而AvroConverter仅支持解析Confluent Avro二进制格式的消息(消息前5字节包含魔术位和Schema Registry中的schema ID,后续为Avro序列化的二进制数据)。纯JSON消息没有这些标识,因此反序列化失败。

解决方案

方案一:发送Avro格式消息(推荐)

使用Confluent官方的Avro生产者工具发送消息,自动将JSON转换为符合要求的Avro二进制格式,并关联已注册的schema:

kafka-avro-console-producer --broker-list <你的Kafka Broker地址>:9092 --topic testtopic --property schema.registry.url=http://localhost:8081 --property value.schema='{"type":"record","name":"sampleTest","fields":[{"name":"testid","type":"string"},{"name":"testname","type":"string"}]}'

启动工具后,输入消息内容:

{"testid": "226", "testname": "testSample21"}

方案二:改用JSON转换器

若需保留纯JSON消息,修改Sink连接器的value转换器配置:

"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable": false,  // 你的JSON消息不包含schema信息
// 移除原有的value.converter.schema.registry.url配置

注意:使用JSONConverter时,auto.create和auto.evolve依赖消息字段结构,需确保消息字段与目标表结构匹配,或提前创建表。

验证步骤

可通过以下命令确认Topic中的消息格式:

kafka-console-consumer --bootstrap-server <你的Kafka Broker地址>:9092 --topic testtopic --from-beginning --property print.value=true

若输出为可读的JSON字符串,说明是纯JSON格式;若为乱码/不可读二进制,则是Avro格式。

内容的提问来源于stack exchange,提问作者Sivapandian Jayakumar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 23:17:30