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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 00:36:28