Kafka Streams读取JDBC源Avro数据时Long反序列化兼容问题
我来帮你拆解这个问题——你遇到的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

