Spring Boot Kafka批量消费中DefaultErrorHandler.addNotRetryableExceptions失效问题
问题描述
在Spring Boot Kafka批量消费场景下,已将NullPointerException添加到DefaultErrorHandler.addNotRetryableExceptions()中,预期触发该异常时直接进入ConsumerRecordRecoverer且不重试,但测试发现会先重试3次才进入恢复器。使用版本为spring-kafka:2.9.2,测试代码如下:
@SpringBootApplication public class DemoKafkaApplication { public static void main(String[] args) { SpringApplication.run(DemoKafkaApplication.class, args); } @Bean CommandLineRunner commandLineRunner(){ return new CommandLineRunner() { @Autowired KafkaTemplate<Integer, String> kafkaTemplate; @Override public void run(String... args) throws Exception { kafkaTemplate.send("topic1", "hello"); kafkaTemplate.send("topic1", "foo"); kafkaTemplate.send("topic1", "derp"); kafkaTemplate.send("topic1", "cheese"); kafkaTemplate.send("topic1", "bar"); } }; } @KafkaListener(id = "myId", topics = "topic1") public void listen(List<String> in) { System.out.println("----------------"); in.forEach(str -> { System.out.println(str); if(str.equalsIgnoreCase("cheese")) throw new NullPointerException("cheese not allowed"); }); System.out.println("----------------"); } class MyConsumerRecordRecoverer implements ConsumerRecordRecoverer{ @Override public void accept(ConsumerRecord<?, ?> consumerRecord, Exception e) { System.out.println(consumerRecord.toString()); } } @Configuration @EnableKafka class KafkaConfig{ @Bean NewTopic topic(){ return TopicBuilder.name("topic1") .build(); } @Bean ConcurrentKafkaListenerContainerFactory<Integer, String> kafkaListenerContainerFactory(ConsumerFactory<Integer, String> consumerFactory) { ConcurrentKafkaListenerContainerFactory<Integer, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); factory.setBatchListener(true); factory.setCommonErrorHandler(commonErrorHandler()); return factory; } @Bean public ConsumerFactory<Integer, String> consumerFactory() { return new DefaultKafkaConsumerFactory<>(consumerProps()); } private Map<String, Object> consumerProps() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ConsumerConfig.GROUP_ID_CONFIG, "group"); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, IntegerDeserializer.class); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); props.put(ConsumerConfig.FETCH_MAX_WAIT_MS_CONFIG, "5000"); props.put(ConsumerConfig.FETCH_MIN_BYTES_CONFIG, "1000"); props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, "2"); return props; } @Bean public CommonErrorHandler commonErrorHandler() { DefaultErrorHandler defaultErrorHandler = new DefaultErrorHandler(myConsumerRecordRecoverer(), new FixedBackOff(2000,3)); defaultErrorHandler.addNotRetryableExceptions(NullPointerException.class); return defaultErrorHandler; } @Bean public MyConsumerRecordRecoverer myConsumerRecordRecoverer(){ return new MyConsumerRecordRecoverer(); } } }
原因分析
DefaultErrorHandler在批量消费模式下,默认将整个批次作为一个处理单元。当遍历批次记录时抛出异常,整个批次的处理会被判定为失败,错误处理器会按照配置的重试次数(此处为3次)重试整个批次,而非针对单个失败记录跳过重试。
addNotRetryableExceptions()的生效逻辑是针对单个记录的异常,但当前代码中异常是从整个批次的处理方法中抛出的,错误处理器无法识别具体是哪个记录导致的异常,因此会执行整个批次的重试流程。
解决方案
要实现单个记录触发不可重试异常时直接进入恢复器,需配置DefaultErrorHandler开启逐记录处理模式,让错误处理器能够针对单个记录判断是否重试:
修改commonErrorHandler()方法,添加defaultErrorHandler.setRetryOnEachRecord(true);:
@Bean public CommonErrorHandler commonErrorHandler() { DefaultErrorHandler defaultErrorHandler = new DefaultErrorHandler(myConsumerRecordRecoverer(), new FixedBackOff(2000,3)); defaultErrorHandler.addNotRetryableExceptions(NullPointerException.class); // 开启逐记录处理模式 defaultErrorHandler.setRetryOnEachRecord(true); return defaultErrorHandler; }
开启该模式后,错误处理器会遍历批次中的每个记录:
- 触发不可重试异常的记录会直接进入
ConsumerRecordRecoverer - 可重试的记录会单独执行重试逻辑
- 处理成功的记录会正常提交偏移量,不会因单个记录失败而重试整个批次
补充说明
spring-kafka 2.8+版本引入的CommonErrorHandler中,DefaultErrorHandler同时支持单条和批量消费场景,但批量场景下默认行为是批次级重试。通过setRetryOnEachRecord(true)切换到单记录级错误处理后,addNotRetryableExceptions()的配置才会对单个记录的异常生效。
内容的提问来源于stack exchange,提问作者jd1289

