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:验证与重启
- 确认Maven生成的
User类正确实现了SpecificRecord接口,编译无报错 - 重启应用,生产者发送的消息现在会是标准Avro二进制格式,消费者能正确反序列化为
User对象,转换异常应该就消失了。
内容的提问来源于stack exchange,提问作者Ouerghi Yassine

