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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 12:27:05