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

Kafka Streams读取JDBC源Avro数据时Long反序列化兼容问题

解决Kafka Streams读取JDBC Avro数据时的LongDeserializer序列化异常

我来帮你拆解这个问题——你遇到的org.apache.kafka.common.errors.SerializationException: Size of data received by LongDeserializer is not 8异常,本质是你的Kafka Streams应用期望消息Key是8字节的Long类型,但实际从JDBC源连接器过来的Key格式完全不匹配。

先搞懂为什么会出这个问题

默认情况下,如果你不给JDBC源连接器指定Key相关的配置,也没通过SMT(比如ValueToKey)手动设置Key,连接器会把消息的Key设为null。而Kafka中null对应的序列化数据长度是0,显然和Long类型要求的8字节对不上,自然会抛异常。

你说尝试用SMT设置Key还是报错,大概率是SMT配置不到位,或者你的Streams应用的Key反序列化器和连接器的Key序列化器不匹配——毕竟你用了Avro序列化器处理Value,但Key的序列化逻辑可能没同步调整。

给你两种可行的解决思路

思路一:让Key也用Avro序列化(推荐,保持格式统一)

既然Value用了Avro,Key也用Avro序列化是最稳妥的做法,只需要调整JDBC源连接器的配置:

# 用SMT把数据库主键字段提取为Key
transforms=ValueToKey
transforms.ValueToKey.fields=your_primary_key_field # 替换成你表的主键,比如id

# 配置Key的Avro序列化器
key.converter=io.confluent.connect.avro.AvroConverter
key.converter.schema.registry.url=http://你的Schema Registry地址:8081

# 保持Value的Avro配置不变
value.converter=io.confluent.connect.avro.AvroConverter
value.converter.schema.registry.url=http://你的Schema Registry地址:8081

然后你的Kafka Streams应用里,Key的反序列化器要换成Avro对应的Serde,比如GenericAvroSerde或者SpecificAvroSerde,和Value保持一致。

思路二:强制让Key为Long类型(如果业务确实需要)

如果你确定必须用Long类型的Key,那得确保连接器生产的Key就是纯Long格式:

# 先把主键字段转成Key,再提取出纯字段值
transforms=ValueToKey,ExtractField
transforms.ValueToKey.fields=your_primary_key_field
transforms.ExtractField.type=org.apache.kafka.connect.transforms.ExtractField$Key
transforms.ExtractField.field=your_primary_key_field

# Key用Long序列化器
key.converter=org.apache.kafka.connect.converters.LongConverter

# Value依旧用Avro
value.converter=io.confluent.connect.avro.AvroConverter
value.converter.schema.registry.url=http://你的Schema Registry地址:8081

这时候你的Streams应用用LongDeserializer就没问题了,因为连接器生产的Key就是标准的8字节Long数据。

验证配置是否正确

你可以用Kafka的命令行工具先验证消息格式:

# 验证Avro Key的情况
kafka-console-consumer.sh --bootstrap-server 你的Kafka Broker地址:9092 --topic 你的主题名 \
  --property print.key=true \
  --property key.deserializer=io.confluent.kafka.serializers.KafkaAvroDeserializer \
  --property value.deserializer=io.confluent.kafka.serializers.KafkaAvroDeserializer \
  --property schema.registry.url=http://你的Schema Registry地址:8081

# 验证Long Key的情况
kafka-console-consumer.sh --bootstrap-server 你的Kafka Broker地址:9092 --topic 你的主题名 \
  --property print.key=true \
  --property key.deserializer=org.apache.kafka.common.serialization.LongDeserializer \
  --property value.deserializer=io.confluent.kafka.serializers.KafkaAvroDeserializer \
  --property schema.registry.url=http://你的Schema Registry地址:8081

如果能正常打印Key和Value,说明配置没问题,再启动Streams应用就不会报错了。

最后总结一下

问题的核心就是Streams的Key反序列化器和连接器生产的Key格式不匹配——要么调整连接器配置,让Key符合Streams的预期;要么调整Streams的反序列化器,匹配连接器的Avro Key格式。

内容的提问来源于stack exchange,提问作者Allen Underwood

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:52:48