Spring Boot中@KafkaHandler无法将Kafka消息反序列化为指定对象
问题描述
我有一个Java Spring Boot应用,原本应该消费Kafka主题product-created-events-topic的消息并转换为ProductCreatedEvent对象,实际却只调用处理字符串的handleDefault方法,而非对应对象的handle方法,请问这是什么原因?
相关代码
消息处理器类
@Component @KafkaListener(topics="product-created-events-topic", groupId = "product-created-events") public class ProductCreatedEventHandler { private final Logger LOGGER = LoggerFactory.getLogger(this.getClass()); @KafkaHandler(isDefault = true) public void handle(ProductCreatedEvent productCreatedEvent){ LOGGER.info("get msg from kafka:"+productCreatedEvent.title()); } @KafkaHandler(isDefault = true) public void handleDefault(String message) { LOGGER.warn("Received unknown message type from Kafka: " + message); } }
ProductCreatedEvent类
package com.test.ws.core; import java.math.BigDecimal; public class ProductCreatedEvent { private String productId; private String title; private BigDecimal price; private Integer quantity; public ProductCreatedEvent(String productId, String title, BigDecimal price, Integer quantity) { this.productId = productId; this.title = title; this.price = price; this.quantity = quantity; } public String productId() { return this.productId; } public ProductCreatedEvent setProductId(String productId) { this.productId = productId; return this; } public String title() { return this.title; } public ProductCreatedEvent setTitle(String title) { this.title = title; return this; } public BigDecimal price() { return this.price; } public ProductCreatedEvent setPrice(BigDecimal price) { this.price = price; return this; } public Integer quantity() { return this.quantity; } public ProductCreatedEvent setQuantity(Integer quantity) { this.quantity = quantity; return this; } }
启动类
@SpringBootApplication public class ConsumersApplication { public static void main(String[] args) { SpringApplication.run(ConsumersApplication.class, args); } }
应用配置
spring.application.name=consumers server.port=8083 spring.kafka.consumer.bootstrap-servers=localhost:9092 spring.kafka.consumer.key-serializer=org.apache.kafka.common.serialization.StringSerializer spring.kafka.consumer.value-serializer=org.springframework.kafka.support.serializer.JsonSerializer spring.kafka.consumer.group-id=product-created-events spring.kafka.consumer.properties.spring.json.trusted.packages=* spring.kafka.consumer.auto-offset-reset=earliest spring.kafka.consumer.enable-auto-commit=false
Kafka主题消息内容
~/kafka/kafka_2.13-3.7.0/bin$ ./kafka-console-consumer.sh --topic product-created-events-topic --from-beginning --bootstrap-server localhost:9092 --property print.key=true --property print.value=true ea59cd46-b205-4769-84d7-69cf37b3ba79 {"productId":"ea59cd46-b205-4769-84d7-69cf37b3ba79","title":"iphone1","price":222,"quantity":19} ea59cd46-b205-4769-84d7-69cf37b3ba79 {"productId":"ea59cd46-b205-4769-84d7-69cf37b3ba79","title":"iphone1","price":222,"quantity":19} b1f9dc43-58b4-44e5-9184-afe9159cd757 {"productId":"b1f9dc43-58b4-44e5-9184-afe9159cd757","title":"iphone2","price":212,"quantity":19}
原因分析及解决方案
1. @KafkaHandler注解配置错误
你给两个方法都加了@KafkaHandler(isDefault = true),但Spring Kafka中一个@KafkaListener下只能有一个默认处理器(isDefault=true)。这个错误会导致处理ProductCreatedEvent的方法没有被识别为对应类型的处理器,所有消息都会走默认的handleDefault方法。
解决: 只给handleDefault方法保留isDefault=true,处理对象的方法去掉该属性:
@KafkaHandler public void handle(ProductCreatedEvent productCreatedEvent){ LOGGER.info("get msg from kafka:"+productCreatedEvent.title()); } @KafkaHandler(isDefault = true) public void handleDefault(String message) { LOGGER.warn("Received unknown message type from Kafka: " + message); }
2. 消费者配置用了序列化器而非反序列化器
配置里的key-serializer和value-serializer是生产者的配置项,消费者应该用key-deserializer和value-deserializer,否则无法正确反序列化Kafka中的JSON消息为Java对象。
解决: 修改配置项:
spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.JsonDeserializer
3. ProductCreatedEvent缺少无参构造函数
JSON反序列化(如Jackson)需要类有无参构造函数来实例化对象,你的ProductCreatedEvent只有全参构造,导致无法创建对象,只能回退到字符串处理。
解决: 添加无参构造函数:
public ProductCreatedEvent() { }
4. 消息缺少类型信息(可选补充)
如果生产者发送消息时没有通过JsonSerializer添加类型头信息,消费者无法自动识别消息对应的Java类。这种情况下可以指定默认类型:
解决: 在消费者配置中添加:
spring.kafka.consumer.properties.spring.json.value.default.type=com.test.ws.core.ProductCreatedEvent
或者确保生产者发送时使用JsonSerializer并开启类型信息(比如配置spring.kafka.producer.value-serializer=org.springframework.kafka.support.serializer.JsonSerializer,并添加spring.kafka.producer.properties.spring.json.add.type.headers=true)。
内容的提问来源于stack exchange,提问作者user63898

