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

Spring Boot Kafka集成Apache Avro反序列化异常求助

问题分析与解决方案

这个异常看起来特别诡异——明明是com.domain.User转com.domain.User,怎么会转换失败?其实这种情况几乎都是类加载器不一致或者序列化/反序列化实现不匹配导致的,结合你的代码和描述,核心问题出在这几个地方:

1. 自定义AvroSerializer是罪魁祸首

你提到没有用官方的Avro序列化依赖,而是复制了教程里的AvroSerializer类。从异常里的payload={"name": "a",...}就能看出猫腻:正常Avro序列化出来的是二进制数据,不会是明文JSON。这说明你的自定义序列化器根本没做Avro二进制序列化,而是把对象转成了JSON字符串。

消费者这边用官方AvroDeserializer<User>反序列化出来的User对象,和生产者那边自定义序列化得到的JSON字符串(或者Spring消息转换器解析后的临时对象),哪怕全类名一样,也会因为类加载器不同、实际对象结构不匹配,被Spring判定为不同类型,于是就抛出了这个看似矛盾的转换异常。

2. 序列化器与反序列化器不匹配

你用Maven插件生成的User类是Avro的SpecificRecord实现类,通用的AvroSerializer/AvroDeserializer是给GenericRecord用的,针对特定记录应该用专门的SpecificAvroSerializer/SpecificAvroDeserializer,否则也会出现序列化/反序列化不兼容的问题。


解决方案

步骤1:引入官方Avro Kafka依赖

先把你复制的自定义AvroSerializer删掉,然后在pom.xml里添加官方依赖(版本要和你的Kafka版本匹配):

<!-- Kafka Avro 序列化核心依赖 -->
<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-avro-serializer</artifactId>
    <version>${kafka.version}</version>
</dependency>
<!-- 若使用Confluent Schema Registry,添加这个依赖(可选) -->
<dependency>
    <groupId>io.confluent</groupId>
    <artifactId>kafka-avro-serializer</artifactId>
    <version>${confluent.version}</version>
</dependency>

步骤2:修正生产者配置

把生产者的序列化器换成SpecificAvroSerializer,并添加专属配置:

@Configuration
public class ProducerConfig {
    @Bean
    public Map<String, Object> producerConfigs() {
        Map<String, Object> props = new HashMap<>();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        // 替换为SpecificAvroSerializer,适配你的User类(SpecificRecord)
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, SpecificAvroSerializer.class);
        // 开启特定Avro读取器,确保序列化正确的类型
        props.put(AbstractKafkaAvroSerDeConfig.SPECIFIC_AVRO_READER_CONFIG, true);
        return props;
    }

    @Bean
    public ProducerFactory<String, User> producerFactory() {
        return new DefaultKafkaProducerFactory<>(producerConfigs());
    }

    // KafkaTemplate和sender Bean保持不变
}

步骤3:修正消费者配置

把消费者的反序列化器换成SpecificAvroDeserializer,同步配置:

@Configuration
public class ConsumerConfig {
    @Bean
    public Map<String, Object> consumerConfigs() {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        // 替换为SpecificAvroDeserializer
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, SpecificAvroDeserializer.class);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "avro");
        // 开启特定Avro读取器,确保反序列化为User类
        props.put(AbstractKafkaAvroSerDeConfig.SPECIFIC_AVRO_READER_CONFIG, true);
        return props;
    }

    @Bean
    public ConsumerFactory<String, User> consumerFactory() {
        return new DefaultKafkaConsumerFactory<>(consumerConfigs(), 
            new StringDeserializer(), 
            new SpecificAvroDeserializer<>(User.class));
    }

    // 监听器工厂和receiver Bean保持不变
}

步骤4:验证与重启

  1. 确认Maven生成的User类正确实现了SpecificRecord接口,编译无报错
  2. 重启应用,生产者发送的消息现在会是标准Avro二进制格式,消费者能正确反序列化为User对象,转换异常应该就消失了。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:42:22