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

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解码错误。

配置修正点

  1. 修正消费者序列化器配置项:
    消费者配置中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
    
  2. 确保Streams全局Serde配置生效:
    将schema.registry.url移到spring.kafka.properties根节点,保证所有Serde都能读取到Schema Registry地址:

    spring:
      kafka:
        properties:
          schema.registry.url: http://localhost:8081
          streams:
            binder:
              configuration:
                # 保留原有Streams配置
    
  3. 为绑定明确指定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_Garment Topic中是否存在非Avro格式的旧消息,若有需清理或配置异常处理器跳过错误消息;
  • 验证Schema Registry中GarmentKey和GarmentEvent的Schema版本,确保生产者和消费者使用同一版本;
  • 确认生产者的KafkaTemplate确实注入了Avro序列化器,未被默认JSON序列化器覆盖。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 23:55:26