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

Spring Boot Kafka Avro反序列化异常:GenericRecord转CarDto失败

解决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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 11:25:58