Spring Kafka消费自定义对象报Can't deserialize data异常
问题背景
- 项目集成Kafka组件初期传输字符串类型消息时运行正常,参照教程调整配置改为传输自定义对象后,配置代码可正常编译构建,但添加监听器代码后消费功能无法正常运行。
- 异常日志出现频次远高于实际发送的业务事件数量,初步判断是Kafka主题中存在无有效消息体的空记录触发报错。
现有代码实现
Kafka配置类代码
import org.apache.kafka.clients.admin.AdminClientConfig; import org.apache.kafka.clients.admin.NewTopic; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.common.serialization.ByteArrayDeserializer; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.serialization.StringSerializer; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.core.env.Environment; import org.springframework.kafka.annotation.EnableKafka; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.*; import org.springframework.kafka.support.serializer.JsonDeserializer; import org.springframework.kafka.support.serializer.JsonSerializer; import paysys.persist.event.StatusUpdatedEvent; import java.util.HashMap; import java.util.Map; @EnableKafka @Configuration @ConditionalOnProperty(name = "status.topic.enabled") public class KafkaConfig { private final Environment environment; private final String topicName = "payOperationStatusChanges"; private final String kafkaGroupId = "status"; public KafkaConfig(Environment environment) { this.environment = environment; } //TopicConfig @Bean public KafkaAdmin kafkaAdmin() { Map<String, Object> configs = new HashMap<>(); configs.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, environment.getProperty("sender.kafka.bootstrap-servers")); return new KafkaAdmin(configs); } @Bean @ConditionalOnProperty(name = "status.topic.enabled") public NewTopic eventTopic() { return new NewTopic(topicName, 1, (short) 1); } //ProducerConfig @Bean public ProducerFactory<String, StatusUpdatedEvent> producerFactory() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, environment.getProperty("sender.kafka.bootstrap-servers")); configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class); configProps.put(ProducerConfig.ACKS_CONFIG, "all"); return new DefaultKafkaProducerFactory<>(configProps); } @Bean public KafkaTemplate<String, StatusUpdatedEvent> statusKafkaTemplate() { return new KafkaTemplate<>(producerFactory()); } //ConsumerConfig @Bean public ConsumerFactory<String, StatusUpdatedEvent> consumerFactory() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, environment.getProperty("sender.kafka.bootstrap-servers")); configProps.put(ConsumerConfig.GROUP_ID_CONFIG, kafkaGroupId); configProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); configProps.put(JsonDeserializer.TRUSTED_PACKAGES, "*"); return new DefaultKafkaConsumerFactory<>(configProps, new StringDeserializer(), new JsonDeserializer<>(StatusUpdatedEvent.class)); } @Bean public ConcurrentKafkaListenerContainerFactory<String, StatusUpdatedEvent> kafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, StatusUpdatedEvent> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); return factory; } }
消息监听器代码
import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Component; import paysys.persist.event.StatusUpdatedEvent; @Component @ConditionalOnProperty(name = "status.topic.enabled") public class StatusEventListener { private static final Logger log = LoggerFactory.getLogger(StatusEventListener.class); @KafkaListener(topics = "operationStatusChanges", groupId = "status", containerFactory = "kafkaListenerContainerFactory") public void listenGroupFoo(StatusUpdatedEvent message) { log.info("Received status update event : {}", message); } }
自定义事件实体代码
事件主类
public class StatusUpdatedEvent extends PayOperationEvent { public static final String STATUS_UPDATED_EVENT_TYPE = "statusUpdated"; private final PayOperation.PayOperationStatus oldStatus; private final PayOperation.PayOperationStatus newStatus; public StatusUpdatedEvent(PayOperation.PayOperationStatus oldStatus, PayOperation.PayOperationStatus newStatus, PayOperation payOperation) { super(STATUS_UPDATED_EVENT_TYPE, payOperation); this.oldStatus = oldStatus; this.newStatus = newStatus; } public PayOperation.PayOperationStatus getOldStatus() { return oldStatus; } public PayOperation.PayOperationStatus getNewStatus() { return newStatus; } @Override public String toString() { return "StatusUpdatedEvent{" + "oldStatus=" + oldStatus + ", newStatus=" + newStatus + ", payOperation=" + this.getPayOperation() + '}'; } }
事件抽象父类
public abstract class PayOperationEvent { private final String type; private final PayOperation payOperation; protected PayOperationEvent(String type, PayOperation payOperation) { this.type = type; this.payOperation = payOperation; } public String getType() { return type; } public PayOperation getPayOperation() { return payOperation; } @Override public String toString() { return "PayOperationEvent{" + "type='" + type + '\'' + ", payOperation=" + payOperation.toString() + '}'; } }
异常信息
org.apache.kafka.common.errors.SerializationException: Error deserializing key/value for partition operationStatusChanges-0 at offset 1987. If needed, please seek past the record to continue consumption. Caused by: org.apache.kafka.common.errors.SerializationException: Can't deserialize data [[]] from topic [operationStatusChanges] Caused by: com.fasterxml.jackson.databind.exc.MismatchedInputException: No content to map due to end-of-input at [Source: (byte[])""; line: 1, column: 0]
问题根因与修复方案
核心问题及对应修复方式如下:
- 主题名配置不一致:配置类中定义的自动创建主题名为
payOperationStatusChanges,但监听器注解中配置的监听主题为operationStatusChanges,消费者实际在消费另一个存量主题,该主题中残留了之前测试产生的空消息/非JSON格式的历史消息,触发反序列化报错。
修复方式:将监听器@KafkaListener注解的topics属性值改为和配置类一致的payOperationStatusChanges,测试环境可直接删除错连的operationStatusChanges主题,避免消费到历史脏数据。 - Json反序列化器未配置空值容错:默认的
JsonDeserializer遇到空字节数组、null值(Kafka墓碑消息)时会直接抛出反序列化异常,即使是正确的主题,后续如果产生空消息也会导致消费阻塞。
修复方式:在消费者配置中添加空值容错配置,给JsonDeserializer开启忽略空内容的参数,修改consumerFactory方法代码如下:
同时在监听器容器工厂中配置错误处理逻辑,跳过无法反序列化的脏消息,避免消费线程永久阻塞:@Bean public ConsumerFactory<String, StatusUpdatedEvent> consumerFactory() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, environment.getProperty("sender.kafka.bootstrap-servers")); configProps.put(ConsumerConfig.GROUP_ID_CONFIG, kafkaGroupId); configProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); configProps.put(JsonDeserializer.TRUSTED_PACKAGES, "*"); // 反序列化时指定默认目标类型 configProps.put(JsonDeserializer.VALUE_DEFAULT_TYPE, StatusUpdatedEvent.class); configProps.put(JsonDeserializer.USE_TYPE_INFO_HEADERS, false); JsonDeserializer<StatusUpdatedEvent> valueDeserializer = new JsonDeserializer<>(StatusUpdatedEvent.class); // 配置反序列化器遇到空内容时返回null而非抛出异常 valueDeserializer.ignoreTypeHeaders(); return new DefaultKafkaConsumerFactory<>(configProps, new StringDeserializer(), valueDeserializer); }@Bean public ConcurrentKafkaListenerContainerFactory<String, StatusUpdatedEvent> kafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, StatusUpdatedEvent> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); // 遇到反序列化错误时直接跳过当前消息,记录日志后提交offset继续消费后续消息 factory.setErrorHandler((e, consumerRecord) -> { log.warn("Skip invalid kafka record, topic:{}, partition:{}, offset:{}, reason:{}", consumerRecord.topic(), consumerRecord.partition(), consumerRecord.offset(), e.getMessage()); }); return factory; } - 自定义实体缺少Jackson反序列化必要配置:
StatusUpdatedEvent和父类PayOperationEvent都只提供了全参构造,没有无参构造方法,Jackson默认反序列化需要无参构造,否则正常的业务消息也可能出现反序列化失败。
修复方式:给两个事件类补充无参构造,或者使用Lombok的@NoArgsConstructor、@AllArgsConstructor注解标注类,也可以手动配置Jackson支持带参构造的反序列化。
内容的提问来源于stack exchange,提问作者agingcabbage32
相关产品推荐
相关产品推荐

