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

Reactor Kafka消费消息Payload为Null及序列化异常问题求助

问题:Reactor Kafka消费消息时Value为null,反序列化报错

问题描述

使用Reactor Kafka时,能正常接收消息触发日志,但打印的消息value始终为null。尝试自定义序列化/反序列化逻辑后,出现java.io.StreamCorruptedException: invalid stream header: 4D657373错误。消息生产消费流程正常,每秒生成一条消息。

问题根源

  1. 初始序列化逻辑失效:原Message类的serialize方法返回空字节数组,导致生产者发送的消息内容为空,消费者自然无法解析出有效value。
  2. 消息格式不匹配:修改序列化逻辑后,Kafka Topic中残留的旧消息(空字节或错误格式)与新的反序列化逻辑不兼容,引发流格式异常。
  3. 自定义Java序列化的局限性:手动实现Java序列化不仅繁琐,还容易出现版本兼容、格式不匹配问题,不如Spring Kafka提供的JSON序列化可靠。

解决方案

步骤1:简化Message类,移除自定义序列化逻辑

将Message改为普通POJO,依赖Lombok注解实现getter/setter,无需手动实现Serializer/Deserializer:

package com.example.demo.model;

import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;

@Data
@AllArgsConstructor
@NoArgsConstructor
public class Message {
    private int id;
    private String content;
    private String timestamp;
}

步骤2:修改Kafka配置,使用Spring Kafka的JSON序列化器

在配置类中替换自定义序列化器为JsonSerializer和JsonDeserializer,并配置信任包:

package com.example.demo.config;

import com.example.demo.model.Message;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.support.serializer.JsonDeserializer;
import org.springframework.kafka.support.serializer.JsonSerializer;
import reactor.kafka.receiver.KafkaReceiver;
import reactor.kafka.receiver.ReceiverOptions;
import reactor.kafka.receiver.internals.ConsumerFactory;
import reactor.kafka.receiver.internals.DefaultKafkaReceiver;
import reactor.kafka.sender.KafkaSender;
import reactor.kafka.sender.SenderOptions;

import java.util.Collections;
import java.util.HashMap;
import java.util.Map;

@Configuration
public class KafkaConfiguration {

    private static final String TOPIC = "data-store";
    private static final String BOOTSTRAP_SERVERS = "localhost:9092";
    private static final String CLIENT_ID_CONFIG = "webflux-client";
    private static final String GROUP_ID_CONFIG = "webflux-group";

    @Bean
    public KafkaReceiver<String, Message> kafkaReceiver(){

        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVERS);
        props.put(ConsumerConfig.CLIENT_ID_CONFIG, CLIENT_ID_CONFIG);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, GROUP_ID_CONFIG);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        // 配置JSON反序列化器
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class.getName());
        // 允许反序列化指定包下的类
        props.put(JsonDeserializer.TRUSTED_PACKAGES, "com.example.demo.model");
        // 指定反序列化的目标类
        props.put(JsonDeserializer.VALUE_DEFAULT_TYPE, Message.class.getName());
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true);

        return new DefaultKafkaReceiver<>(ConsumerFactory.INSTANCE, ReceiverOptions.create(props).subscription(Collections.singleton(TOPIC)));
    }

    @Bean
    public KafkaSender<String, Message> kafkaSender(){
        Map<String, Object> props = new HashMap<>();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVERS);
        props.put(ProducerConfig.CLIENT_ID_CONFIG, CLIENT_ID_CONFIG);
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        // 配置JSON序列化器
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class.getName());

        SenderOptions<String, Message> senderOptions = SenderOptions.create(props);
        return KafkaSender.create(senderOptions);
    }
}

步骤3:清理旧消息,避免格式冲突

  • 停止所有相关服务
  • 清理Kafka Topic(使用命令行工具):
    kafka-topics.sh --delete --topic data-store --bootstrap-server localhost:9092
    kafka-topics.sh --create --topic data-store --bootstrap-server localhost:9092 --partitions 1 --replication-factor 1
    
    或者修改消费者配置auto.offset.reset为latest,仅消费新消息。

步骤4:重启服务测试

启动生产者和消费者服务,调用/kafka/generate-messages接口生成消息,查看消费者日志即可看到正常解析的Message对象。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 06:54:56