能否使用spring-kafka的RecoveringBatchErrorHandler处理反序列化异常?
解决方案
方案一(最优,无业务侵入)
自定义批量消息转换器,在转换阶段捕获单条消息的转换异常,携带失败索引抛出BatchListenerFailedException,让RecoveringBatchErrorHandler可以识别并按预期处理。
- 实现自定义批量消息转换器
import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.converter.BatchMessagingMessageConverter; import org.springframework.kafka.support.converter.ConversionException; import org.springframework.kafka.support.converter.RecordMessageConverter; import org.springframework.messaging.Message; import org.springframework.kafka.listener.BatchListenerFailedException; import java.lang.reflect.Type; import java.util.ArrayList; import java.util.List; public class FaultTolerantBatchMessagingMessageConverter extends BatchMessagingMessageConverter { public FaultTolerantBatchMessagingMessageConverter(RecordMessageConverter recordConverter) { super(recordConverter); } @Override public List<Message<?>> toMessage(List<ConsumerRecord<?, ?>> records, Acknowledgment acknowledgment, Consumer<?, ?> consumer, Type payloadType) { List<Message<?>> convertedMessages = new ArrayList<>(); int failIndex = 0; try { for (ConsumerRecord<?, ?> record : records) { Message<?> msg = getRecordMessageConverter().toMessage(record, acknowledgment, consumer, payloadType); convertedMessages.add(msg); failIndex++; } } catch (ConversionException e) { throw new BatchListenerFailedException("消息反序列化失败", e, failIndex); } return convertedMessages; } }
- 替换容器工厂中的消息转换器
修改你原有的kafkaListenerContainerFactory配置中的转换器声明:
// 替换原有的BatchMessagingMessageConverter BatchMessagingMessageConverter messageConverter = new FaultTolerantBatchMessagingMessageConverter(new BytesJsonMessageConverter()); factory.setMessageConverter(messageConverter);
- (可选)配置失败消息恢复逻辑
可以在RecoveringBatchErrorHandler中自定义重试耗尽后的处理逻辑,比如将失败消息存入死信队列或错误日志表,避免消息丢失:
RecoveringBatchErrorHandler errorHandler = new RecoveringBatchErrorHandler( (failedRecord, e) -> { // 自定义恢复逻辑示例 log.error("消息处理失败,topic: {}, partition: {}, offset: {}, 异常: {}", failedRecord.topic(), failedRecord.partition(), failedRecord.offset(), e.getMessage(), e); // 可扩展为写入死信队列、存储到错误库等操作 }, new FixedBackOff(FixedBackOff.DEFAULT_INTERVAL, 2) );
方案二(轻量实现,有业务侵入)
如果不想自定义转换器,可以直接调整监听器入参类型,在业务代码中自行处理反序列化:
- 修改监听器入参为原始字节类型的
ConsumerRecord列表:
@KafkaListener(id = "myKafkaListener", idIsGroup = false, autoStartup = "true", topics = {"pojo-topic"}, containerFactory = "kafkaListenerContainerFactory") public void receive(List<ConsumerRecord<byte[], byte[]>> records) { ObjectMapper objectMapper = new ObjectMapper(); int index = 0; try { for (ConsumerRecord<byte[], byte[]> record : records) { Pojo pojo = objectMapper.readValue(record.value(), Pojo.class); // 原有的业务处理逻辑 System.out.println("received: " + pojo + " at offset " + record.offset()); index++; } } catch (Exception e) { throw new BatchListenerFailedException("处理失败", e, index); } }
- 移除原来配置的
BatchMessagingMessageConverter即可。
原理解释
原生BatchMessagingMessageConverter在批量转换时,只要任意一条消息转换失败就会直接抛出不带索引的ConversionException,RecoveringBatchErrorHandler无法定位失败位置,只能整批回溯重试。自定义转换器后会逐个转换消息,捕获到异常时携带失败位置抛出BatchListenerFailedException,错误处理器即可识别并提交失败位置之前的所有偏移量,重试/恢复失败消息后,从下一个偏移量开始拉取新批次。
内容的提问来源于stack exchange,提问作者Chad Showalter
相关产品推荐
相关产品推荐

