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的消息后重新测试。
快速验证方法
- 暂时注释掉消费者的
@KafkaListener代码,启动应用后发送消息到topic-A,确认是否仍出现反序列化错误。如果无错误,说明问题源于消费者配置与生产者的干扰。 - 打印生产者和消费者的配置信息,确认
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
相关产品推荐
相关产品推荐

