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

Spring Cloud Stream 3.2.4中如何在消费端手动确认(Ack)消息?

Spring Cloud Stream 3.2.4消费端手动确认(Ack)消息(函数式编程模型)

要在函数式编程模型下实现手动消息确认,需要完成配置禁用自动确认和消费逻辑中调用确认接口两步操作。

第一步:配置禁用自动确认

先在配置文件(如application.yml)中为对应的消费者绑定设置手动确认模式:

spring:
  cloud:
    stream:
      bindings:
        test-in-0: # 绑定名称规则:函数名-in-0,对应你定义的Consumer Bean名test
          destination: your-target-topic-or-queue
          group: your-consumer-group
          consumer:
            acknowledge-mode: manual

第二步:在消费逻辑中手动确认

在Consumer的实现代码里,可从Message的headers中获取Acknowledgment对象,调用ack()完成消息确认,或调用nack()拒绝消息(拒绝后的行为依中间件而定,比如重发或进入死信队列)。

修改后的完整代码如下:

import org.springframework.messaging.Message;
import org.springframework.cloud.stream.binder.Acknowledgment;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration
public class ConsumerConfig {

    @Bean
    public Consumer<Message<String>> test() {
        return msg -> {
            try {
                // 执行你的业务处理逻辑
                System.out.println("处理消息内容:" + msg.getPayload());
                
                // 手动确认消息处理成功
                Acknowledgment acknowledgment = msg.getHeaders().get(Acknowledgment.ACKNOWLEDGMENT_HEADER, Acknowledgment.class);
                if (acknowledgment != null) {
                    acknowledgment.ack();
                }
            } catch (Exception e) {
                // 处理异常,手动拒绝消息
                Acknowledgment acknowledgment = msg.getHeaders().get(Acknowledgment.ACKNOWLEDGMENT_HEADER, Acknowledgment.class);
                if (acknowledgment != null) {
                    // nack参数可指定是否重发,不同中间件支持度有差异
                    acknowledgment.nack(false);
                }
                // 按需抛出异常,让框架做后续处理
                throw new RuntimeException("消息处理失败", e);
            }
        };
    }
}

关键说明

  • Acknowledgment对象由框架注入到消息headers中,仅当acknowledge-mode设为manual时才会存在。
  • 调用ack()表示消息处理成功,中间件会移除该消息;调用nack()表示处理失败,具体行为取决于使用的中间件(比如RabbitMQ会按配置决定是否重发,Kafka会标记消息为未提交,后续可能重新消费)。
  • 若既不调用ack()也不调用nack(),消息会一直处于未确认状态,直到超时或中间件触发重发逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 09:18:29