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

Spring Boot中@KafkaHandler无法将Kafka消息反序列化为指定对象

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 05:35:13