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

如何修改KafkaListener失败前的DELIVERY_ATTEMPT重试次数?

问题分析
  1. NullPointerException原因:默认情况下Spring Kafka不会自动添加KafkaHeaders.DELIVERY_ATTEMPT头,需显式配置重试机制并启用该头,才能在ErrorHandler中获取到重试次数。
  2. 重试行为未改变:仅自定义KafkaListenerErrorHandler不足以覆盖默认的10次重试,需结合容器的重试配置(如SeekToCurrentErrorHandler或RetryTemplate)来控制重试次数和异常类型。
解决方案:区分可重试/不可重试异常的配置示例

1. 自定义异常类型(可选,也可直接用内置异常)

// 可重试异常
public class RetriableException extends RuntimeException {
    public RetriableException(String message) {
        super(message);
    }
}

// 不可重试异常
public class NonRetriableException extends RuntimeException {
    public NonRetriableException(String message) {
        super(message);
    }
}

2. 配置RetryTemplate(定义重试规则)

配置重试模板,指定仅对RetriableException进行重试,最多重试3次,同时设置重试间隔:

@Bean
public RetryTemplate retryTemplate() {
    RetryTemplate retryTemplate = new RetryTemplate();

    // 重试策略:仅对RetriableException重试,最多3次
    SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy();
    retryPolicy.setMaxAttempts(3);
    Map<Class<? extends Throwable>, Boolean> retryableExceptions = new HashMap<>();
    retryableExceptions.put(RetriableException.class, true); // 可重试
    retryableExceptions.put(NonRetriableException.class, false); // 不可重试
    retryPolicy.setRetryableExceptions(retryableExceptions);
    retryTemplate.setRetryPolicy(retryPolicy);

    // 退避策略:每次重试间隔1秒
    FixedBackOffPolicy backOffPolicy = new FixedBackOffPolicy();
    backOffPolicy.setBackOffPeriod(1000);
    retryTemplate.setBackOffPolicy(backOffPolicy);

    return retryTemplate;
}

3. 配置KafkaListenerContainerFactory(关联重试模板与错误处理器)

配置容器工厂,启用DELIVERY_ATTEMPT头,并设置自定义错误处理器:

@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(
        ConsumerFactory<String, String> consumerFactory,
        RetryTemplate retryTemplate) {

    ConcurrentKafkaListenerContainerFactory<String, String> factory =
            new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory);

    // 配置SeekToCurrentErrorHandler,结合重试模板,启用deliveryAttempt头
    SeekToCurrentErrorHandler errorHandler = new SeekToCurrentErrorHandler(
            (record, ex) -> {
                // 重试耗尽后的处理逻辑:比如记录日志、死信队列等
                log.error("消息重试耗尽,丢弃记录: {}", record.value(), ex);
            },
            new FixedBackOff(1000L, 0L)) // 重试耗尽后不再重试
            .setRetryTemplate(retryTemplate)
            .setDeliveryAttemptHeader(true); // 启用DELIVERY_ATTEMPT头

    factory.setErrorHandler(errorHandler);
    return factory;
}

4. 修改Kafka监听器(使用自定义异常)

@KafkaListener(topics = "${some_topic}", autoStartup = "true", containerFactory = "kafkaListenerContainerFactory")
public void listenForMessage(String message) {
    log.warn("Accepted message: {}", message);

    // 模拟可重试异常场景
    if (message.contains("retry")) {
        throw new RetriableException("需要重试的错误");
    }

    // 模拟不可重试异常场景
    if (message.contains("no-retry")) {
        throw new NonRetriableException("无需重试的错误");
    }

    // 其他逻辑...
}

5. 自定义KafkaListenerErrorHandler(可选,细粒度控制)

如果需要在ErrorHandler中直接处理重试次数,可这样配置(前提是已启用deliveryAttemptHeader):

@Bean
public KafkaListenerErrorHandler customErrorHandler() {
    return (message, exception) -> {
        Integer deliveryAttempt = message.getHeaders().get(KafkaHeaders.DELIVERY_ATTEMPT, Integer.class);
        // 已启用deliveryAttemptHeader,此处不会抛出NPE
        if (deliveryAttempt != null && deliveryAttempt >= 3) {
            log.error("第{}次重试失败,终止处理", deliveryAttempt);
            return "FAILED";
        }
        // 抛出异常继续重试(仅当异常为可重试类型时)
        throw exception;
    };
}

然后在监听器中引用:

@KafkaListener(
        topics = "${some_topic}",
        autoStartup = "true",
        containerFactory = "kafkaListenerContainerFactory",
        errorHandler = "customErrorHandler"
)
关键说明
  • DELIVERY_ATTEMPT头需通过SeekToCurrentErrorHandler.setDeliveryAttemptHeader(true)启用,否则无法获取。
  • 重试次数优先通过RetryTemplate的setMaxAttempts配置,配合异常类型过滤,实现精准的可重试/不可重试区分。
  • 重试耗尽时,可在SeekToCurrentErrorHandler的失败回调中实现死信队列转发、日志记录等后续处理。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 01:53:32