Spring Boot Kafka Avro反序列化异常:GenericRecord转CarDto失败
问题背景
给Avro Schema新增了一个带默认值的int类型属性:
{ "name": "minimum", "type": [ "null", "int" ], "default": null }
Schema Registry已设置为向前兼容模式,生产者与消费者在同一项目中,使用相同的生成Avro类,但测试时消费者抛出反序列化异常:
org.springframework.messaging.converter.MessageConversionException: Cannot handle message; nested exception is org.springframework.messaging.converter.MessageConversionException: Cannot convert from [org.apache.avro.generic.GenericData$Record] to [com.example.dto.CarDto] for GenericMessage [payload={....}
项目基于Spring Boot 2.5,未使用DevTools,已配置:
spring.kafka.consumer.value-deserializer=io.confluent.kafka.serializers.KafkaAvroDeserializer spring.kafka.consumer.properties.specific.avro.reader=true
更新说明:问题与Avro文件更新无关,此前异常被抑制,修改后才打印出日志,论坛线索指向可能是Avro类包路径问题。
项目结构
- src->main->avro->Car.avsc
- src->main->java->com->example->dto->CarDto.java
- src->main->java->com->example->configs->producer/consumer
- src->main->java->com->example->service->ConsumerProcessor.java
原消费者配置
@Slf4j @Configuration public class ConsumerConfig { @Value("${kafka.bootstrapAddress}") private String bootstrapAddress; @Value("${kafka.schemaRegistryUrl}") private String schemaRegistryUrl; // Consumers Configs private Map<String, Object> getConsumerProps() { final var props = new HashMap<String, Object>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress); props.put(ConsumerConfig.GROUP_ID_CONFIG, "topic-a"); props.put(AbstractKafkaAvroSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, schemaRegistryUrl); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); return props; } private <T, K> ConcurrentKafkaListenerContainerFactory<K, T> getConcurrentKafkaListenerContainerFactory(final ConsumerFactory<K, T> kafkaConsumerFactory, final KafkaTemplate<K, T> template) { final ConcurrentKafkaListenerContainerFactory<K, T> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(kafkaConsumerFactory); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); return factory; } private KafkaAvroDeserializer getAvroDeserializer() { final var deserializerProps = new HashMap<String, Object>(); deserializerProps.put(AbstractKafkaAvroSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, schemaRegistryUrl); deserializerProps.put(KafkaAvroDeserializerConfig.SPECIFIC_AVRO_READER_CONFIG, "true"); final var deserializer = new KafkaAvroDeserializer(); deserializer.configure(deserializerProps, false); return deserializer; } @Bean public ConcurrentKafkaListenerContainerFactory<Integer, com.exmaple.dto.CarDto> dataKafkaListenerContainerFactory( final ConsumerFactory<Integer, com.example.dto.CarDto> kafkaConsumerFactory, final KafkaTemplate<Integer, com.example.dto.CarDto> template) { return getConcurrentKafkaListenerContainerFactory(kafkaConsumerFactory, template); } @Bean("kafkaListenerContainerFactory") public ConsumerFactory<Integer, com.example.dto.CarDto> dataConsumerFactory() { final Map<String, Object> props = getConsumerProps(); final KafkaAvroDeserializer deserializer = getAvroDeserializer(); return new DefaultKafkaConsumerFactory<>(props, new IntegerDeserializer(), new ErrorHandlingDeserializer(deserializer)); } }
Avro Schema定义
{ "type": "record", "name": "CarDto", "namespace": "com.example.dto", "fields": [ { "name": "id", "type": [ "null", "int" ], "default": null }, { "name": "minimum", "type": ["null", {"type": "int", "logicalType": "date"}], "default": null } ] }
消费者代码
@Slf4j @Service public class dataConsumer implements AvroKafkaConsumer<CarDto> { Integer counter = 0; @Override @KafkaListener(topics = "topic-a", containerFactory = "dataKafkaListenerContainerFactory") public void listen(final CarDto carDto, final Acknowledgment acknowledgment) { acknowledgment.acknowledge(); } @Override public String getName() { return "data"; } }
发现的配置疑点
消费者@KafkaListener指定的containerFactory名称为dataKafkaListenerContainerFactory,但配置中@Bean声明的ConsumerFactory名称是kafkaListenerContainerFactory,两者是否存在匹配问题?
修改后的配置
private Map<String, Object> getProps() { final var props = new HashMap<String, Object>(); // ... 原有配置 props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, Integer.class); props.put(KafkaAvroDeserializerConfig.SPECIFIC_AVRO_READER_CONFIG, true); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class); props.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, KafkaAvroDeserializer.class); return props; } @Bean public ConcurrentKafkaListenerContainerFactory<Integer, CarDto> myKafkaContainerFactory() { final ConcurrentKafkaListenerContainerFactory<Integer, CarDto> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(kafkaConsumerFactory()); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); return factory; } public ConsumerFactory<Integer, CarDto> kafkaConsumerFactory() { final Map<String, Object> props = getProps(); final KafkaAvroDeserializer deserializer = new KafkaAvroDeserializer(); return new DefaultKafkaConsumerFactory<>(props, new IntegerDeserializer(), new ErrorHandlingDeserializer(deserializer)); }
解决方案
1. 修正容器工厂名称匹配
修改消费者代码中@KafkaListener的containerFactory参数为myKafkaContainerFactory,确保与配置类中声明的Bean名称完全一致,避免Spring注入错误的容器工厂。
2. 修复类路径拼写错误
原配置中存在com.exmaple.dto.CarDto的拼写错误(exmaple应为example),这个错误会导致无法找到目标类,强制使用GenericRecord反序列化。检查所有配置和代码中引用的类路径,确保拼写完全正确。
3. 确保SPECIFIC_AVRO_READER配置生效
移除application.properties中的重复配置,仅保留配置类中的KafkaAvroDeserializerConfig.SPECIFIC_AVRO_READER_CONFIG=true,同时避免手动创建KafkaAvroDeserializer时的重复配置。可简化ConsumerFactory的创建:
public ConsumerFactory<Integer, CarDto> kafkaConsumerFactory() { final Map<String, Object> props = getProps(); return new DefaultKafkaConsumerFactory<>(props); }
让Spring根据props自动创建反序列化器,保证配置参数能正确传递到底层的KafkaAvroDeserializer。
4. 验证Avro类与Schema命名空间一致性
确认生成的CarDto.java在com.example.dto包下,且Schema中的namespace与类包路径完全匹配。同时执行clean和compile命令,确保编译后的Avro类是最新的。
5. 清理冲突配置
删除application.properties中的以下配置,避免与自定义配置类冲突:
spring.kafka.consumer.value-deserializer=io.confluent.kafka.serializers.KafkaAvroDeserializer spring.kafka.consumer.properties.specific.avro.reader=true
内容的提问来源于stack exchange,提问作者user1555190

