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

@KafkaListener无法消费消息:反序列化问题求助

排查与解决@KafkaListener反序列化异常问题

看起来你踩了Spring Cloud Stream Kafka和Spring Kafka序列化机制不匹配的坑,我帮你梳理下核心问题和具体的排查修复步骤:

核心问题分析

你的问题本质是生产者和消费者的序列化/反序列化逻辑完全不匹配,再加上Spring Cloud Stream的自动序列化机制和你手动编码的冲突,以及消费者容器工厂配置的小失误,导致消息始终无法正常消费。


分步排查与修复

1. 修正生产者的序列化配置(最关键)

你第二个版本的生产者直接发送PageViewEvent对象,但配置里的default.value.serde是BytesSerde——这会让Spring Cloud Stream把POJO转换成Bytes对象而非JSON字节数组,消费者用JsonDeserializer自然无法解析。

修复方案:

修改生产者配置,统一使用JSON序列化:

# 全局配置JSON Serde(推荐)
spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde=org.springframework.kafka.support.serializer.JsonSerde
# 或者给指定绑定单独配置(更灵活)
spring.cloud.stream.bindings.pageViewsOut.content-type=application/json
spring.cloud.stream.kafka.bindings.pageViewsOut.producer.value-serde=org.springframework.kafka.support.serializer.JsonSerde

如果你还是想用第一个版本的手动序列化代码,记得要去掉Stream的自动序列化处理,或者确保手动序列化后的byte[]不会被再次处理——但这种方式不推荐,不如让Stream帮你统一处理序列化逻辑。

2. 修正消费者容器工厂的名称匹配

你更新后的消费者代码里,@KafkaListener指定的containerFactory是kafkaListenerContainerFactory,但你定义的容器工厂Bean名称是priceEventsKafkaListenerContainerFactory——名称完全不匹配,导致消费者根本没用到你配置的JsonDeserializer!

修复方案(二选一):

  • 要么修改@KafkaListener的容器工厂指定:
    @KafkaListener(topics = "test1" , groupId = "json", containerFactory = "priceEventsKafkaListenerContainerFactory")
    public void receive(@Payload PageViewEvent data,@Headers MessageHeaders headers) {
        // ... 你的代码
    }
    
  • 要么把容器工厂Bean名称改成默认的kafkaListenerContainerFactory,这样可以不用手动指定:
    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, PageViewEvent> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, PageViewEvent> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(priceEventConsumerFactory());
        return factory;
    }
    

3. 完善JsonDeserializer的配置

Spring的JsonDeserializer默认有安全限制,需要指定信任的包,否则会拒绝反序列化外部类型:

@Bean
public ConsumerFactory<String,PageViewEvent > priceEventConsumerFactory() {
    Map<String, Object> props = new HashMap<>();
    props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
    props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
    props.put(ConsumerConfig.GROUP_ID_CONFIG, "json");
    props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
    
    // 配置JsonDeserializer,信任你的POJO所在包
    JsonDeserializer<PageViewEvent> deserializer = new JsonDeserializer<>(PageViewEvent.class);
    deserializer.setTrustedPackages("com.your.package"); // 替换成你的PageViewEvent所在包名,或者用"*"允许所有(测试环境可用)
    return new DefaultKafkaConsumerFactory<>(props, new StringDeserializer(), deserializer);
}

4. 清理Topic中的旧消息

如果你的test1 Topic里已经存在之前生产者发送的错误格式消息(比如Bytes对象的序列化结果),即使修复了配置,消费者消费这些旧消息还是会报错。

处理方式(测试环境可用):

  • 重置消费者组的offset到最新位置:
    kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group json --reset-offsets --to-latest --topic test1 --execute
    
  • 或者直接删除并重建Topic:
    kafka-topics.sh --bootstrap-server localhost:9092 --delete --topic test1
    kafka-topics.sh --bootstrap-server localhost:9092 --create --topic test1 --partitions 1 --replication-factor 1
    

5. 验证消息格式(可选)

用Kafka命令行工具确认生产者发送的是正确的JSON格式:

kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic test1 --from-beginning --property print.value=true --property value.deserializer=org.apache.kafka.common.serialization.StringDeserializer

如果输出是类似{"userId":"priya","page":"blog","duration":10}的JSON字符串,说明生产者的序列化已经正常。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 04:45:20