spring-kafka 2.2.7中ErrorHandlingDeserializer2配置及异常与记录获取
嘿,针对你在spring-kafka 2.2.7.RELEASE里配置ErrorHandlingDeserializer2处理反序列化异常的需求,我来一步步给你讲清楚怎么拿到异常信息和对应的消费者记录,同时实现日志记录或者发送到DLT的功能:
1. 先把基础配置搭好
首先,你需要在消费者配置里指定ErrorHandlingDeserializer2作为key/value的反序列化器,同时配置它的委托反序列化器(也就是你原本想用的那个反序列化器,比如JsonDeserializer)。
举个Java配置的例子:
@Bean public ConsumerFactory<String, MyBusinessObject> consumerFactory() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "your-kafka-servers"); configProps.put(ConsumerConfig.GROUP_ID_CONFIG, "your-consumer-group"); // 配置value的反序列化器为ErrorHandlingDeserializer2 configProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer2.class); // 指定委托的实际反序列化器(比如JSON反序列化) configProps.put(ErrorHandlingDeserializer2.VALUE_DESERIALIZER_CLASS, JsonDeserializer.class.getName()); // 如果需要处理key的反序列化异常,同理配置key的部分 // configProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer2.class); // configProps.put(ErrorHandlingDeserializer2.KEY_DESERIALIZER_CLASS, StringDeserializer.class.getName()); return new DefaultKafkaConsumerFactory<>(configProps); }
2. 捕获反序列化异常与消费者记录
有两种常用方式可以拿到异常信息和原始的ConsumerRecord:
方式一:自定义ErrorReporter(推荐直接捕获反序列化环节的异常)
ErrorHandlingDeserializer2允许你配置一个ErrorReporter,当反序列化失败时,会主动回调这个类,把异常和原始记录传进来。
第一步:写自定义的ErrorReporter
@Component public class CustomDeserializationErrorReporter implements ErrorReporter { private static final Logger logger = LoggerFactory.getLogger(CustomDeserializationErrorReporter.class); @Autowired private KafkaTemplate<String, byte[]> kafkaTemplate; // 用于发送到DLT @Override public void report(Exception exception, ConsumerRecord<?, ?> record) { // 这里能直接拿到反序列化异常和原始的消费者记录 logger.error("反序列化失败!Topic: {}, Partition: {}, Offset: {}, 原始Key: {}, 原始Value: {}", record.topic(), record.partition(), record.offset(), record.key(), record.value(), exception); // 可选:把失败的记录发送到死信队列(DLT) sendToDlt(record); } private void sendToDlt(ConsumerRecord<?, ?> record) { String dltTopic = record.topic() + ".DLT"; // 约定DLT主题名 kafkaTemplate.send(dltTopic, record.key(), record.value()); logger.info("已将失败记录发送到DLT主题: {}", dltTopic); } }
第二步:在消费者配置里指定这个ErrorReporter
在之前的consumerFactory配置里加上一行,指定ErrorReporter的bean名称:
configProps.put(ErrorHandlingDeserializer2.ERROR_REPORTER_BEAN_NAME, "customDeserializationErrorReporter");
方式二:通过ListenerErrorHandler捕获(适合在消费环节统一处理错误)
如果你想在@KafkaListener的消费逻辑里统一处理所有错误(包括反序列化异常),可以自定义一个ListenerErrorHandler:
第一步:写自定义的ListenerErrorHandler
@Bean public ListenerErrorHandler customKafkaErrorHandler() { return (message, exception) -> { // 反序列化异常会被包装在DeserializationException里 if (exception.getCause() instanceof DeserializationException) { DeserializationException de = (DeserializationException) exception.getCause(); ConsumerRecord<?, ?> originalRecord = de.getOriginalMessage(); logger.error("消费时捕获反序列化异常!Topic: {}, Offset: {}", originalRecord.topic(), originalRecord.offset(), de); // 同样可以在这里发送到DLT sendToDlt(originalRecord); } // 返回null表示跳过这条消息,或者抛出异常让容器处理(比如重试) return null; }; } // 这里可以复用上面的sendToDlt方法,记得注入KafkaTemplate
第二步:在@KafkaListener里指定这个错误处理器
@KafkaListener( topics = "your-topic", groupId = "your-consumer-group", errorHandler = "customKafkaErrorHandler" ) public void handleMessage(MyBusinessObject payload) { // 正常的消费业务逻辑 }
3. 注意事项
- 确保你的
KafkaTemplate配置正确,能正常连接到Kafka集群,否则发送DLT会失败。 - 如果你用的是yaml/properties配置文件,也可以把上述配置项写到配置文件里,比如:
spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.ErrorHandlingDeserializer2 spring.kafka.consumer.properties.spring.deserializer.value.delegate.class=org.springframework.kafka.support.serializer.JsonDeserializer spring.kafka.consumer.properties.spring.deserializer.error.reporter.bean.name=customDeserializationErrorReporter
- 在2.2.7.RELEASE版本中,
ErrorHandlingDeserializer2已经能很好地封装反序列化异常,不用额外的复杂配置。
内容的提问来源于stack exchange,提问作者Raj
相关产品推荐
相关产品推荐

