Spring Kafka Streams Avro反序列化报错求助
问题:Spring Kafka Streams消费Avro消息时反序列化失败
使用Spring从Kafka的Queue_Garment Topic消费消息时,消费者抛出反序列化异常。已确认生产者与消费者使用相同的Avro Schema,此前运行正常但现在失效,怀疑是反序列化器配置问题。
错误栈信息
Exception in thread "Promotions_Handler-a2f86d00-554e-43c5-b984-faafb7524faf-StreamThread-1" org.apache.kafka.streams.errors.StreamsException: Deserialization exception handler is set to fail upon a deserialization error. If you would rather have the streaming pipeline continue after a deserialization error, please set the default.deserialization.exception.handler appropriately. at org.apache.kafka.streams.processor.internals.RecordDeserializer.deserialize(RecordDeserializer.java:82) at org.apache.kafka.streams.processor.internals.RecordQueue.updateHead(RecordQueue.java:176) at org.apache.kafka.streams.processor.internals.RecordQueue.addRawRecords(RecordQueue.java:112) at org.apache.kafka.streams.processor.internals.PartitionGroup.addRawRecords(PartitionGroup.java:185) at org.apache.kafka.streams.processor.internals.StreamTask.addRecords(StreamTask.java:895) at org.apache.kafka.streams.processor.internals.TaskManager.addRecordsToTasks(TaskManager.java:1008) at org.apache.kafka.streams.processor.internals.StreamThread.pollPhase(StreamThread.java:812) at org.apache.kafka.streams.processor.internals.StreamThread.runOnce(StreamThread.java:625) at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:564) at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:523) Caused by: org.apache.kafka.common.errors.SerializationException: Can't deserialize data [[0, 0, 0, 0, 1, 24, 77, 49, 50, 50, 57, 57, 115, 100, 97, 115, 102, 53]] from topic [Queue_Garment] Caused by: java.io.CharConversionException: Invalid UTF-32 character 0x1174d31 (above 0x0010ffff) at char #1, byte #7) at com.fasterxml.jackson.core.io.UTF32Reader.reportInvalid(UTF32Reader.java:195) at com.fasterxml.jackson.core.io.UTF32Reader.read(UTF32Reader.java:158) at com.fasterxml.jackson.core.json.ReaderBasedJsonParser._loadMore(ReaderBasedJsonParser.java:255) at com.fasterxml.jackson.core.json.ReaderBasedJsonParser._skipWSOrEnd(ReaderBasedJsonParser.java:2389) at com.fasterxml.jackson.core.json.ReaderBasedJsonParser.nextToken(ReaderBasedJsonParser.java:677) at com.fasterxml.jackson.databind.ObjectReader._initForReading(ObjectReader.java:355) at com.fasterxml.jackson.databind.ObjectReader._bindAndClose(ObjectReader.java:2023) at com.fasterxml.jackson.databind.ObjectReader.readValue(ObjectReader.java:1528) at org.springframework.kafka.support.serializer.JsonDeserializer.deserialize(JsonDeserializer.java:534) at org.apache.kafka.streams.processor.internals.SourceNode.deserializeKey(SourceNode.java:54) at org.apache.kafka.streams.processor.internals.RecordDeserializer.deserialize(RecordDeserializer.java:65) at org.apache.kafka.streams.processor.internals.RecordQueue.updateHead(RecordQueue.java:176) at org.apache.kafka.streams.processor.internals.RecordQueue.addRawRecords(RecordQueue.java:112) at org.apache.kafka.streams.processor.internals.PartitionGroup.addRawRecords(PartitionGroup.java:185) at org.apache.kafka.streams.processor.internals.StreamTask.addRecords(StreamTask.java:895) at org.apache.kafka.streams.processor.internals.TaskManager.addRecordsToTasks(TaskManager.java:1008) at org.apache.kafka.streams.processor.internals.StreamThread.pollPhase(StreamThread.java:812) at org.apache.kafka.streams.processor.internals.StreamThread.runOnce(StreamThread.java:625) at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:564) at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:523)
消费者代码
@StreamListener public KStream<GarmentKey, GarmentEvent> newGarment(@Input(BinderProcessor.Queue_Garment)KStream<GarmentKey,GarmentEvent> garment){ updateDatabase(garment); uploadGarmentToQueue(garment); return garment.peek(((GarmentKey,GarmentEvent) -> System.out.println("promotion body ="+GarmentKey.toString()))); }
生产者代码
public void postGarment(PostGarmentDto postGarmentDto) { if(goodReference(postGarmentDto.getReference())) { GarmentKey garmentKey = new GarmentKey(); garmentKey.setReference(postGarmentDto.getReference()); GarmentEvent garmentEvent = new GarmentEvent(); garmentEvent.setCategory(postGarmentDto.getCategory()); garmentEvent.setPrice(postGarmentDto.getPrice()); kafkaTemplate.send("Queue_Garment", garmentKey, garmentEvent); } }
消费者配置文件(Consumer.yml)
spring: config: activate: on-profile: default application: name: Promotions_Handler kafka: consumer: key-serializer: io.confluent.kafka.serializers.KafkaAvroSerializer value-serializer: io.confluent.kafka.serializers.KafkaAvroSerializer properties: streams: binder: auto-create-topics: true configuration: state.dir: /tmp commit.interval.ms: 100 topology.optimization: all session.timeout.ms: 10000 schema.registry.url: http://localhost:8081 auto.register.schemas: true default.key.serde: io.confluent.kafka.streams.serdes.avro.SpecificAvroSerde default.value.serde: io.confluent.kafka.streams.serdes.avro.SpecificAvroSerde spring.cloud.stream.bindings.Queue_Promotions_Applied: destination: Queue_Promotions_Applied producer: useNativeEncoding: true spring.cloud.stream.bindings.Queue_Garment: destination: Queue_Garment consumer: useNativeDecoding: true spring.cloud.stream.bindings.Queue_Promotions: destination: Queue_Promotions consumer: useNativeDecoding: true
生产者配置文件(Producer.yml)
spring: application: name: Garment_receiver kafka: properties: bootstrap.servers: localhost:9092 schema.registry.url: http://localhost:8081 producer: key-serializer: io.confluent.kafka.serializers.KafkaAvroSerializer value-serializer: io.confluent.kafka.serializers.KafkaAvroSerializer
解决方案
核心问题定位
错误栈显示消费者实际使用JsonDeserializer而非配置的Avro Serde,说明Avro反序列化配置未生效,导致用JSON解析逻辑处理Avro二进制数据,触发UTF-32解码错误。
配置修正点
修正消费者序列化器配置项:
消费者配置中spring.kafka.consumer下误用了生产者的key-serializer/value-serializer,需改为key-deserializer/value-deserializer,并指定Avro反序列化器,同时开启specific.avro.reader以适配生成的Avro实体类:spring: kafka: consumer: key-deserializer: io.confluent.kafka.serializers.KafkaAvroDeserializer value-deserializer: io.confluent.kafka.serializers.KafkaAvroDeserializer properties: specific.avro.reader: true确保Streams全局Serde配置生效:
将schema.registry.url移到spring.kafka.properties根节点,保证所有Serde都能读取到Schema Registry地址:spring: kafka: properties: schema.registry.url: http://localhost:8081 streams: binder: configuration: # 保留原有Streams配置为绑定明确指定Serde(可选):
在Queue_Garment绑定配置中添加Serde指定,确保优先级高于全局配置:spring.cloud.stream.bindings.Queue_Garment: destination: Queue_Garment consumer: useNativeDecoding: true configuration: key.serde: io.confluent.kafka.streams.serdes.avro.SpecificAvroSerde value.serde: io.confluent.kafka.streams.serdes.avro.SpecificAvroSerde
额外排查点
- 检查
Queue_GarmentTopic中是否存在非Avro格式的旧消息,若有需清理或配置异常处理器跳过错误消息; - 验证Schema Registry中
GarmentKey和GarmentEvent的Schema版本,确保生产者和消费者使用同一版本; - 确认生产者的
KafkaTemplate确实注入了Avro序列化器,未被默认JSON序列化器覆盖。
内容的提问来源于stack exchange,提问作者miguel30452
相关产品推荐
相关产品推荐

