Reactor Kafka消费消息Payload为Null及序列化异常问题求助
问题:Reactor Kafka消费消息时Value为null,反序列化报错
问题描述
使用Reactor Kafka时,能正常接收消息触发日志,但打印的消息value始终为null。尝试自定义序列化/反序列化逻辑后,出现java.io.StreamCorruptedException: invalid stream header: 4D657373错误。消息生产消费流程正常,每秒生成一条消息。
问题根源
- 初始序列化逻辑失效:原
Message类的serialize方法返回空字节数组,导致生产者发送的消息内容为空,消费者自然无法解析出有效value。 - 消息格式不匹配:修改序列化逻辑后,Kafka Topic中残留的旧消息(空字节或错误格式)与新的反序列化逻辑不兼容,引发流格式异常。
- 自定义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 1auto.offset.reset为latest,仅消费新消息。
步骤4:重启服务测试
启动生产者和消费者服务,调用/kafka/generate-messages接口生成消息,查看消费者日志即可看到正常解析的Message对象。
内容的提问来源于stack exchange,提问作者user1354825
相关产品推荐
相关产品推荐

