Spring Boot中Kafka Offset10次重试后递增问题及无限重试配置咨询
解决Kafka消费者无限重试的配置修改
当前你的Kafka消费者重试10次后跳过消息,是因为Spring Kafka默认的DefaultErrorHandler限制了最多10次重试。要实现仅通过acknowledge和nack无限重试,需修改错误处理器的重试次数限制,具体如下:
- 自定义
DefaultErrorHandler,将最大重试次数设置为Integer.MAX_VALUE(模拟无限次),同时关闭默认的退避策略(符合你不使用Fixed/ExponentialBackOff的要求)。 - 在消费者容器工厂中注册这个自定义错误处理器,替换默认实现。
示例配置代码:
@Configuration public class KafkaConsumerConfig { @Bean public DefaultErrorHandler kafkaErrorHandler() { DefaultErrorHandler errorHandler = new DefaultErrorHandler(); // 设置最大重试次数为极大值,实现无限重试 errorHandler.setMaxAttempts(Integer.MAX_VALUE); // 关闭默认退避策略,仅依赖代码中nack指定的间隔 errorHandler.setBackOffFunction((record, ex) -> BackOffExecution.STOP); return errorHandler; } @Bean public ConcurrentKafkaListenerContainerFactory<?, ?> kafkaListenerContainerFactory( ConsumerFactory<Object, Object> consumerFactory, DefaultErrorHandler errorHandler) { ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); // 保持手动提交配置 factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); // 绑定自定义错误处理器 factory.setCommonErrorHandler(errorHandler); return factory; } }
修改后,当消费逻辑抛出异常时,nack(Duration.ofMillis(1000))会触发每隔1秒重试一次,且不会因为达到重试次数上限而跳过消息,直到消息处理成功并调用acknowledge()为止。
内容的提问来源于stack exchange,提问作者Ravi
相关产品推荐
相关产品推荐

