Spring Boot 2.1.8中Acknowledgment.nack()等效功能及实现方案问询
Spring Boot 2.1.8 Kafka批量消费失败重试解决方案
Spring Boot 2.1.8 内置的 spring-kafka 版本为 2.2.x,Acknowledgment.nack() 是 spring-kafka 2.3 版本才新增的API,因此该版本确实无法直接使用nack方法。
替代方案:使用SeekToCurrentBatchErrorHandler
完全可以使用SeekToCurrentBatchErrorHandler实现你的需求:批量消费过程中如果单条消息处理失败抛出异常时,该错误处理器会自动将当前批次所有未处理的消息(包含消费失败的消息)seek回对应分区的当前消费位置,下次poll时会重新拉取整个批次的消息处理,符合你要求的剩余未消费记录下次重新拉取的逻辑。
具体配置步骤
1. 基础消费参数配置
在application.yml中开启批量消费、关闭自动提交:
spring: kafka: consumer: enable-auto-commit: false max-poll-records: 10 # 按需配置单次拉取最大消息数 key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer group-id: 你的消费组ID listener: type: batch ack-mode: manual_immediate
2. 注入SeekToCurrentBatchErrorHandler配置
编写Kafka配置类,将错误处理器绑定到监听容器工厂:
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.listener.SeekToCurrentBatchErrorHandler; import org.springframework.kafka.listener.ContainerProperties.AckMode; import java.time.Duration; @Configuration public class KafkaConsumerConfig { @Bean public ConcurrentKafkaListenerContainerFactory<Object, Object> kafkaListenerContainerFactory( ConsumerFactory<Object, Object> consumerFactory) { ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); // 开启批量监听 factory.setBatchListener(true); // 配置批量消费错误处理器 SeekToCurrentBatchErrorHandler errorHandler = new SeekToCurrentBatchErrorHandler(); // 可选配置:失败后重试间隔,不配置默认立即重试 errorHandler.setBackOffFunction((record, exception) -> Duration.ofSeconds(3)); factory.setBatchErrorHandler(errorHandler); // 设置手动提交模式 factory.getContainerProperties().setAckMode(AckMode.MANUAL_IMMEDIATE); return factory; } }
3. 批量消费逻辑编写
单条消息处理失败时直接抛出异常即可触发重试,整个批次处理成功后再提交偏移量:
import org.apache.kafka.clients.consumer.ConsumerRecord; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.support.Acknowledgment; import org.springframework.stereotype.Component; import java.util.List; @Component public class BatchMessageConsumer { @KafkaListener(topics = "你的业务Topic名称") public void consume(List<ConsumerRecord<String, String>> records, Acknowledgment ack) { for (ConsumerRecord<String, String> record : records) { try { // 单条消息业务处理逻辑 // 模拟处理失败抛出异常,触发错误处理器 if (isProcessFail(record)) { throw new RuntimeException("单条消息处理失败"); } } catch (Exception e) { // 抛出异常触发seek逻辑 throw e; } } // 整个批次全部处理成功后提交偏移量 ack.acknowledge(); } private boolean isProcessFail(ConsumerRecord<String, String> record) { // 自定义业务错误判断逻辑,按需修改 return false; } }
注意事项
- 重试的是整个批次的消息,业务逻辑必须保证幂等性,避免重复消费导致数据异常
- 如果需要限制最大重试次数,可在
SeekToCurrentBatchErrorHandler中扩展逻辑,超过重试次数的消息可以转发至死信队列存储处理 - 偏移量提交必须在整个批次所有消息处理成功后调用,不要中途提交部分偏移量,否则会导致失败消息丢失
内容的提问来源于stack exchange,提问作者ppb
相关产品推荐
相关产品推荐

