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

Kafka报SerializationException异常如何配置忽略反序列化失败消息

问题根因

Kafka消费流程中,消息反序列化动作执行在消费者拉取消息之后、@KafkaListener标注的监听方法被调用之前,因此写在监听方法内部的try-catch无法捕获SerializationException。未做特殊配置时,该异常会导致消费线程反复重试拉取同一条格式错误的消息,直接阻断整个消费流程。

实现方案

1. 修改Kafka配置类,包装反序列化器

使用Spring Kafka内置的ErrorHandlingDeserializer分别包装原有的key、value反序列化器。反序列化失败时该包装类不会直接向外抛出阻断异常,而是将异常信息写入消息Header,将消息值置为null传递给后续流程。
修改后的KafkaConfig代码如下:

@Configuration
@Slf4j
public class KafkaConfig {
    private final KafkaProperties kafkaProperties;

    public KafkaConfig(KafkaProperties kafkaProperties) {
        this.kafkaProperties = kafkaProperties;
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, ExternalCardServiceData> kafkaContainerFactory() {
        // 构造目标类型JSON反序列化器,配置信任所有包避免类型校验异常
        JsonDeserializer<ExternalCardServiceData> jsonDeserializer = new JsonDeserializer<>(ExternalCardServiceData.class);
        jsonDeserializer.addTrustedPackages("*");
        // 用ErrorHandlingDeserializer包装key、value的实际反序列化器
        ErrorHandlingDeserializer<String> keyErrorDeserializer = new ErrorHandlingDeserializer<>(new StringDeserializer());
        ErrorHandlingDeserializer<ExternalCardServiceData> valueErrorDeserializer = new ErrorHandlingDeserializer<>(jsonDeserializer);

        DefaultKafkaConsumerFactory<String, ExternalCardServiceData> consumerFactory =
                new DefaultKafkaConsumerFactory<>(kafkaProperties.buildConsumerProperties(),
                        keyErrorDeserializer, valueErrorDeserializer);
        
        ConcurrentKafkaListenerContainerFactory<String, ExternalCardServiceData> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory);
        // 配置错误处理器:异常触发时记录日志,不重试直接提交偏移量跳过坏消息
        factory.setCommonErrorHandler(new DefaultErrorHandler((consumerRecord, e) -> {
            log.error("消费消息失败,跳过当前消息,topic:{}, partition:{}, offset:{}, 异常信息:{}",
                    consumerRecord.topic(), consumerRecord.partition(), consumerRecord.offset(), e.getMessage());
        }, new FixedBackOff(0L, 0L)));
        // 配置消费完成单条消息后自动提交偏移量
        factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.RECORD);
        return factory;
    }
}

注:如果使用的Spring Kafka版本低于2.8,将上述代码中的DefaultErrorHandler替换为SeekToCurrentErrorHandler,构造参数逻辑完全一致。

2. 调整监听方法逻辑,兼容反序列化失败场景

反序列化失败时消息value为null,需要提前判空避免空指针,可按需记录错误信息:

@KafkaListener(topics = "#{'${topic}'}",
        groupId = "#{'${groupid}'}",
        autoStartup = "#{'${enabled}'}",
        containerFactory = "kafkaContainerFactory")
public void updateExternalCardToken(ConsumerRecord<String, ExternalCardServiceData> record) {
    String key = record.key();
    ExternalCardServiceData externalCardServiceData = record.value();
    // 反序列化失败时value为null,直接记录错误后返回
    if (externalCardServiceData == null) {
        log.error("消息反序列化失败,跳过处理,offset:{}", record.offset());
        externalCardService.saveTokensErrors(key, "消息格式不合法,反序列化失败", "SerializationException");
        return;
    }
    try {
        log.info("ExternalCardsListener. Received message: {}, offset={}", externalCardServiceData, record.offset());
        externalCardService.updateToken(key, externalCardServiceData);
    } catch (Exception e) {
        log.error("Error processing message received from kafka. [Message={}]", externalCardServiceData);
        externalCardService.saveTokensErrors(key, externalCardServiceData.toString(),
                Arrays.toString(e.getStackTrace()));
    }
}
配置说明
  • ErrorHandlingDeserializer为委托型反序列化器,本身不执行实际反序列化逻辑,仅捕获下层反序列化器抛出的所有异常,避免异常直接穿透到消费容器层导致消费线程卡死
  • 上述配置中FixedBackOff(0L, 0L)代表异常触发后重试间隔0ms、最大重试次数0次,即失败后直接跳过,不会因为单条坏消息阻塞整个消费进度
  • 反序列化失败的完整异常栈会写入消息Header,对应key为ErrorHandlingDeserializer.VALUE_DESERIALIZER_EXCEPTION_HEADER,如果需要持久化完整异常信息,可以从该Header中读取序列化后的异常对象

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 23:09:23