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

如何配置Kafka自定义反序列化器通过@KafkaListener获取消费的JSON数据

核心错误修正与实现步骤

你现有代码存在3个核心问题,按以下步骤修改即可正常运行:

1. 修复无效的自定义反序列化器

你写的反序列化器调用了SerializationUtils.deserialize,这是处理Java原生序列化字节流的方法,完全无法解析JSON格式消息,且没有处理泛型目标类型的读取,修正后代码如下:

public class CustomDeserializer<T extends Serializable> implements Deserializer<T> {
    private final ObjectMapper objectMapper = new ObjectMapper();
    public static final String VALUE_CLASS_NAME_CONFIG = "value.class.name";
    private Class<T> targetClass;

    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {
        // 从配置中读取要反序列化的目标实体类型
        targetClass = (Class<T>) configs.get(VALUE_CLASS_NAME_CONFIG);
    }

    @Override
    public T deserialize(String topic, byte[] objectData) {
        if (objectData == null) {
            return null;
        }
        try {
            return objectMapper.readValue(objectData, targetClass);
        } catch (IOException e) {
            throw new RuntimeException("JSON消息反序列化失败", e);
        }
    }

    @Override
    public void close() {}
}

注意:提前给KafkaPayload、EventHeader、NewFields、OldFields四个实体类加无参构造、getter/setter,或者直接加Lombok的@Data注解,ObjectMapper才能正常反射赋值;如果JSON字段名和实体类属性名不匹配(比如JSON里是大驼峰EventHeader,实体属性是小驼峰eventHeader),要在属性上加@JsonProperty("EventHeader")做映射。


2. 修正Kafka配置类

你原来的consumerFactory方法返回类型错误:Spring Kafka的监听容器需要的是ConsumerFactory(消费者实例工厂),而非直接new出来的原生KafkaConsumer实例——原生KafkaConsumer是非线程安全的,容器需要根据并发配置自行创建多个消费者实例,不能直接注入单例的消费者对象。修正后配置类如下:

@Configuration
// 启动类记得加@EnableKafka注解,否则@KafkaListener不会生效
public class KafkaConfig {

    @Bean
    public ConsumerFactory<String, KafkaPayload> consumerFactory(){
        Map<String, Object> config = new HashMap<>();
        config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, CustomDeserializer.class);
        // 给自定义反序列化器传入目标类型
        config.put(CustomDeserializer.VALUE_CLASS_NAME_CONFIG, KafkaPayload.class);
        return new DefaultKafkaConsumerFactory<>(config);
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, KafkaPayload> kafkaListener(){
        ConcurrentKafkaListenerContainerFactory<String, KafkaPayload> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        return factory;
    }
}

提效提示:如果没有特殊自定义需求,完全可以不用自己写反序列化器,直接用Spring Kafka自带的JsonDeserializer,只需要把反序列化类配置替换为JsonDeserializer.class,再加一行配置config.put(JsonDeserializer.VALUE_DEFAULT_TYPE, KafkaPayload.class)即可,稳定性更高。


3. 修正消费监听类

给消费类加@Service注解让Spring上下文托管,否则监听器不会被扫描加载:

@Service
public class ConsumerService {

    @KafkaListener(topics = "Topic", groupId = "groupId", containerFactory = "kafkaListener")
    public void consume(KafkaPayload kafkaPayload){
        System.out.println("Consumed Message :" + kafkaPayload);
    }
}

启动前检查项

  • 确认Kafka服务端正常运行,监听的Topic已提前创建
  • 确认生产端发送的消息是标准JSON格式,无多余转义、不可见字符
  • 确认Spring Boot启动类上添加了@EnableKafka注解

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.01 17:01:27