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

Spring-Kafka同时作为生产者和消费者时出现反序列化异常

Spring-Kafka 生产者发送消息触发消费者反序列化错误排查

可能的原因及排查步骤

1. 检查Topic配置是否混淆

  • 确认生产者代码中指定的发送Topic确实是topic-A,消费者@KafkaListener注解中监听的是topic-B,避免代码中写错Topic名称,导致生产者发送的消息被自身消费者监听并尝试反序列化(比如误写为同一个Topic)。
  • 用kafka-topics.sh工具查看集群中Topic的实际名称,注意Kafka Topic名称默认区分大小写,避免大小写或拼写错误。

2. 排查生产者与消费者的配置隔离情况

Spring Boot自动配置会默认创建全局的ProducerFactory、ConsumerFactory和KafkaTemplate,如果生产者和消费者需要不同的序列化/反序列化配置,必须明确隔离:

  • 确保自定义了独立的ProducerFactory和KafkaTemplate,指定生产者专属的value.serializer(对应生产者POJO);同时为消费者创建独立的ConsumerFactory和容器工厂,指定对应消费者POJO的value.deserializer。
  • 示例隔离配置:
    // 生产者配置
    @Bean
    public ProducerFactory<String, ProducerPojo> producerFactory() {
        Map<String, Object> configProps = new HashMap<>();
        configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "your-kafka-server:9092");
        configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
        return new DefaultKafkaProducerFactory<>(configProps);
    }
    
    @Bean
    public KafkaTemplate<String, ProducerPojo> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactory());
    }
    
    // 消费者配置
    @Bean
    public ConsumerFactory<String, ConsumerPojo> consumerFactory() {
        Map<String, Object> configProps = new HashMap<>();
        configProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "your-kafka-server:9092");
        configProps.put(ConsumerConfig.GROUP_ID_CONFIG, "consumer-group-b");
        configProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        configProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
        configProps.put(JsonDeserializer.TRUSTED_PACKAGES, "com.yourpackage.consumer.pojo");
        return new DefaultKafkaConsumerFactory<>(configProps);
    }
    
    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, ConsumerPojo> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, ConsumerPojo> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        return factory;
    }
    
  • 注意@KafkaListener要指定对应的容器工厂,确保用消费者专属配置:
    @KafkaListener(topics = "topic-B", containerFactory = "kafkaListenerContainerFactory")
    public void listen(ConsumerPojo message) {
        // 消费逻辑
    }
    

3. 检查全局配置是否干扰

  • 查看application.properties/application.yml中的全局配置,确认spring.kafka.producer.value.serializer和spring.kafka.consumer.value.deserializer是否正确设置,避免生产者配置遗漏value.serializer,导致Spring Boot自动误用消费者的反序列化器配置。
  • 示例正确的全局配置:
    spring:
      kafka:
        bootstrap-servers: your-kafka-server:9092
        producer:
          key-serializer: org.apache.kafka.common.serialization.StringSerializer
          value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
        consumer:
          group-id: consumer-group-b
          key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
          value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
          properties:
            spring.json.trusted.packages: com.yourpackage.consumer.pojo
            spring.json.value.default.type: com.yourpackage.consumer.pojo.ConsumerPojo
    

4. 排查序列化/反序列化的类型信息冲突

如果使用JsonSerializer和JsonDeserializer,默认会在消息头中添加__TypeId__字段:

  • 检查生产者的JsonSerializer是否配置了addTypeHeaders=false,避免消费者反序列化时被生产者的类型头干扰。
  • 确保消费者的JsonDeserializer指定了trusted-packages和value.default.type,避免因类型信息不匹配抛出错误。

5. 检查重试/死信队列的干扰

  • 如果消费者配置了重试或死信队列,可能在启动时处理旧消息,与生产者发送动作重叠导致错误混淆。可以暂时关闭消费者重试机制,或清空topic-B的消息后重新测试。

快速验证方法

  1. 暂时注释掉消费者的@KafkaListener代码,启动应用后发送消息到topic-A,确认是否仍出现反序列化错误。如果无错误,说明问题源于消费者配置与生产者的干扰。
  2. 打印生产者和消费者的配置信息,确认value.serializer和value.deserializer是否对应各自的POJO:
    // 生产者配置中打印
    System.out.println("Producer value serializer: " + producerFactory().getConfigurationProperties().get(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG));
    // 消费者配置中打印
    System.out.println("Consumer value deserializer: " + consumerFactory().getConfigurationProperties().get(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG));
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 13:40:47