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

Spring Boot Kafka中ErrorHandlingDeserializer后获取ConsumerRecord值的问题

解决方案

1. 利用ErrorHandlingDeserializer内置异常获取原始数据

ErrorHandlingDeserializer在反序列化失败时,会抛出FailedDeserializationException,该异常包含FailedDeserializationInfo对象,直接存储了未反序列化的原始字节数组,完全不需要调用KafkaTemplate.receive()(该方法会再次触发反序列化,必然重复报错)。

修改你的KafkaErrorHandler代码:

public class KafkaErrorHandler implements ErrorHandler {

    // 注入数据库操作DAO(根据你的实际实现调整)
    @Autowired
    private FailedRecordDao failedRecordDao;

    @Override
    public void handle(Exception thrownException, ConsumerRecord<?, ?> data) {
        // 精准捕获反序列化失败异常
        if (thrownException instanceof FailedDeserializationException) {
            FailedDeserializationException deserializationEx = (FailedDeserializationException) thrownException;
            FailedDeserializationInfo errorInfo = deserializationEx.getFailedDeserializationInfo();
            
            // 提取原始字节数组和消息元数据
            byte[] rawMessage = errorInfo.getRawValue();
            String topic = data.topic();
            int partition = data.partition();
            long offset = data.offset();
            String messageKey = (String) data.key();
            
            // 组装数据写入数据库
            FailedRecord failedRecord = new FailedRecord();
            failedRecord.setTopic(topic);
            failedRecord.setPartition(partition);
            failedRecord.setOffset(offset);
            failedRecord.setMessageKey(messageKey);
            failedRecord.setRawContent(rawMessage);
            failedRecord.setErrorDetail(thrownException.getMessage());
            failedRecordDao.save(failedRecord);
        } else {
            // 处理非反序列化类异常
            log.error("Kafka消息处理非反序列化错误: {}", thrownException.getMessage(), thrownException);
        }
    }
}

2. 将ErrorHandler注册到Kafka监听容器工厂

需要把自定义的错误处理器配置到Spring Kafka的监听容器工厂,确保它能生效:

@Configuration
public class KafkaConfig {

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory(
            ConsumerFactory<String, Object> consumerFactory,
            KafkaErrorHandler kafkaErrorHandler) {
        ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory);
        // 绑定自定义错误处理器
        factory.setErrorHandler(kafkaErrorHandler);
        return factory;
    }

    @Bean
    public KafkaErrorHandler kafkaErrorHandler() {
        return new KafkaErrorHandler();
    }
}

3. 确认application.properties配置有效性

你的现有配置已经正确使用了ErrorHandlingDeserializer作为代理,无需修改,可添加日志配置辅助排查:

spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer
spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
spring.kafka.properties.spring.deserializer.value.delegate.class=org.springframework.kafka.support.serializer.JsonDeserializer
# 可选:开启反序列化组件日志,方便排查问题
logging.level.org.springframework.kafka.support.serializer=DEBUG

关键注意点

  • 禁止使用KafkaTemplate.receive():该方法会复用当前消费者的反序列化器再次尝试解析,必然触发相同异常,无法获取原始字节。
  • FailedDeserializationInfo还提供了反序列化目标类型、委托反序列化器等信息,可用于后续错误分析。
  • 如果需要处理key的反序列化错误,只需给key-deserializer也配置ErrorHandlingDeserializer,逻辑与value处理一致。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 17:47:39