如何修改KafkaListener失败前的DELIVERY_ATTEMPT重试次数?
问题分析
- NullPointerException原因:默认情况下Spring Kafka不会自动添加
KafkaHeaders.DELIVERY_ATTEMPT头,需显式配置重试机制并启用该头,才能在ErrorHandler中获取到重试次数。 - 重试行为未改变:仅自定义
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
相关产品推荐
相关产品推荐

