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
相关产品推荐
相关产品推荐

