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

Spring Boot Kafka批量消费中DefaultErrorHandler.addNotRetryableExceptions失效问题

Spring Kafka批量消费中DefaultErrorHandler不可重试异常不生效问题

问题描述

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 21:11:39