ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

RabbitMQ 消息确认机制(ACK):技术解析与实践

RabbitMQ 消息确认机制(ACK):技术解析与实践 RabbitMQ 消息确认机制ACK技术解析与实践一、ACK 是什么ACKAcknowledge是消费者告诉 Broker “这条消息我处理完了你可以删除了” 的信号。没有 ACKBroker 不知道消息是否被成功处理也就无法决定是否删除消息。Producer → Broker(Queue) → Consumer │ ├─ 处理成功 → ACK → Broker 删除消息 │ ├─ 处理失败 → NACK/Reject → Broker 重新入队或进死信 │ └─ 消费者宕机 → 超时无 ACK → Broker 重新分发给其他消费者注博客https://blog.csdn.net/badao_liumang_qizhi二、为什么需要 ACK没有 ACK 机制会怎样场景无 ACK 的后果消费者收到消息后处理到一半崩溃消息丢失业务不完整消费者处理失败业务异常无法重试数据不一致网络闪断消费者没收到消息Broker 以为已推送消息丢失ACK 机制的核心保证消息至少被成功处理一次At Least Once。三、三种确认模式1. 自动确认autospring:rabbitmq:listener:simple:acknowledge-mode:auto行为消息推送到消费者方法后Spring 根据方法执行结果自动决定方法正常返回 → 自动 ACK方法抛异常 → 自动 NACK requeueRabbitListener(queuesorder.queue)publicvoidconsume(IntegerorderId){orderService.process(orderId);// 正常返回 → 框架自动 ACK// 抛异常 → 框架自动 NACK消息重新入队}优点代码简单不需要手动管理缺点异常时无限 requeue 可能导致消息循环需配合重试策略2. 手动确认manualspring:rabbitmq:listener:simple:acknowledge-mode:manual行为消费者必须在代码中显式调用basicAck或basicNack不调用则消息一直处于 Unacked 状态消费者断连后消息重新入队RabbitListener(queuesorder.queue)publicvoidconsume(IntegerorderId,Channelchannel,Header(AmqpHeaders.DELIVERY_TAG)longdeliveryTag){try{orderService.process(orderId);// 处理成功确认消息channel.basicAck(deliveryTag,false);}catch(BusinessExceptione){// 业务异常不重试直接丢弃或进死信channel.basicReject(deliveryTag,false);}catch(Exceptione){// 临时异常重新入队等待重试channel.basicNack(deliveryTag,false,true);}}优点精细控制可以区分不同异常做不同处理缺点代码复杂度增加忘记 ACK 会导致消息堆积3. 无确认nonespring:rabbitmq:listener:simple:acknowledge-mode:none行为Broker 推送后立即删除消息不等待消费者确认等同于发后即忘优点最高吞吐缺点消息可能丢失仅适用于可丢失的场景如日志采集、监控指标四、ACK 相关的三个操作basicAck — 确认成功channel.basicAck(deliveryTag,multiple);参数说明deliveryTag消息的唯一标识Broker 分配单 Channel 内递增multiplefalse只确认当前消息true确认 deliveryTag 及之前所有未确认的消息// 逐条确认channel.basicAck(deliveryTag,false);// 批量确认确认 tag5 的所有消息channel.basicAck(5,true);basicNack — 否定确认批量channel.basicNack(deliveryTag,multiple,requeue);参数说明deliveryTag消息标识multiple是否批量否定requeuetrue重新入队false丢弃或进死信// 单条否定重新入队channel.basicNack(deliveryTag,false,true);// 单条否定不重回队列进入死信队列或丢弃channel.basicNack(deliveryTag,false,false);basicReject — 否定确认单条channel.basicReject(deliveryTag,requeue);功能与basicNack相同但只能操作单条消息没有 multiple 参数。五、deliveryTag 详解Channel 1: msg_a(tag1) msg_b(tag2) msg_c(tag3) Channel 2: msg_x(tag1) msg_y(tag2)deliveryTag 是每个 Channel 独立递增的序号消费者收到消息时携带 tagACK 时回传这个 tagBroker 通过 (channel tag) 唯一定位一条消息六、Unacked 消息与 prefetch 的关系prefetch 3 时 Broker Queue: [msg4] [msg5] [msg6] ... ↑ 等待中需要消费者 ACK 释放额度 Consumer 手中: [msg1:处理中] [msg2:处理中] [msg3:处理中] ↑ Unacked 数量 3 prefetch 上限Broker 最多推prefetch条未确认消息给消费者消费者 ACK 一条后Broker 才会推下一条如果消费者始终不 ACK达到 prefetch 上限后不再推送新消息消费者宕机时Broker 检测到 Channel/Connection 断开所有 Unacked 消息自动回到 Queue 头部其他消费者可以重新消费这些消息七、常见问题与陷阱陷阱1忘记 ACK 导致消息堆积// 错误示例manual 模式下忘记 ACKRabbitListener(queuesorder.queue)publicvoidconsume(IntegerorderId,Channelchannel,Header(AmqpHeaders.DELIVERY_TAG)longtag){orderService.process(orderId);// 忘记调用 channel.basicAck(tag, false);// 后果消息一直 Unacked达到 prefetch 后不再推送新消息}排查方式RabbitMQ 管理后台看 Queue 的 Unacked 数量持续增长。陷阱2异常时 requeue 导致无限循环// 错误示例所有异常都 requeueRabbitListener(queuesorder.queue)publicvoidconsume(IntegerorderId,Channelchannel,Header(AmqpHeaders.DELIVERY_TAG)longtag){try{orderService.process(orderId);channel.basicAck(tag,false);}catch(Exceptione){// 如果是参数错误永远不会成功每次 requeue 都会再失败channel.basicNack(tag,false,true);// 无限循环}}正确做法区分可重试异常和不可重试异常try{orderService.process(orderId);channel.basicAck(tag,false);}catch(RetryableExceptione){// 临时故障网络超时、锁冲突重新入队channel.basicNack(tag,false,true);}catch(Exceptione){// 不可恢复参数错误、数据不存在拒绝不重回log.error(消费失败且不可重试, orderId{},orderId,e);channel.basicReject(tag,false);// 进死信或丢弃}陷阱3批量确认丢消息// 风险示例multipletrue 批量确认channel.basicAck(tag,true);// 如果 tag5则 tag 1~5 全部确认// 如果 tag3 的消息其实还没处理完也会被确认掉建议除非有明确的批量处理逻辑否则始终用multiplefalse。八、重试策略配置Spring 内置重试auto 模式下spring:rabbitmq:listener:simple:acknowledge-mode:autoretry:enabled:trueinitial-interval:1000# 第一次重试间隔 1smax-interval:10000# 最大间隔 10smultiplier:2.0# 间隔倍增max-attempts:3# 最大重试次数流程第1次消费失败 → 等1s → 第2次重试 → 等2s → 第3次重试 → 仍失败 → 进入 MessageRecoverer自定义失败处理器BeanpublicMessageRecoverermessageRecoverer(RabbitTemplaterabbitTemplate){// 重试耗尽后发送到死信队列returnnewRepublishMessageRecoverer(rabbitTemplate,dlx.exchange,dlx.order);}死信队列DLX配置BeanpublicQueueorderQueue(){MapString,ObjectargsnewHashMap();args.put(x-dead-letter-exchange,dlx.exchange);// 死信交换机args.put(x-dead-letter-routing-key,dlx.order);// 死信路由键returnnewQueue(order.queue,true,false,false,args);}BeanpublicQueuedeadLetterQueue(){returnnewQueue(order.dlq,true);}BeanpublicBindingdlqBinding(){returnBindingBuilder.bind(deadLetterQueue()).to(newDirectExchange(dlx.exchange)).with(dlx.order);}消息进入死信队列的条件被 reject/nack 且 requeuefalse消息 TTL 过期队列达到最大长度九、生产者确认Publisher ConfirmACK 不仅消费端有生产端也有——确认消息成功到达 Brokerspring:rabbitmq:publisher-confirm-type:correlated# 异步确认publisher-returns:true# 路由失败回调Component public class OrderMqProducer{Resource private RabbitTemplate rabbitTemplate; PostConstruct public void init(){// 消息到达 Exchange 的确认 rabbitTemplate.setConfirmCallback((data,ack,cause)-{if (!ack){log.error(消息未到达Exchange,cause{},cause); // 重发或记录}}); // 消息无法路由到 Queue 的回调 rabbitTemplate.setReturnsCallback(returned-{log.error(消息无法路由,exchange{},routingKey{},replyText{},returned.getExchange(),returned.getRoutingKey(),returned.getReplyText());});}}完整的消息可靠性链路Producer → Confirm → Exchange → Routing → Queue → Consumer → ACK │ │ └── 生产端确认消息到达 Broker 消费端确认消息处理完成──┘十、三种确认模式选型指南模式适用场景消息安全性吞吐量代码复杂度auto大多数业务场景高中低manual需要精细控制的核心业务最高中低高none日志采集、监控指标等可丢失场景无最高最低实际项目建议默认用auto 重试配置 死信队列覆盖 90% 场景核心资金类业务用manual精确控制每条消息的命运none仅用于明确标注允许丢失的非关键数据十一、完整示例可靠消费模板ComponentSlf4jpublicclassOrderMqConsumer{ResourceprivateOrderServiceorderService;ResourceprivateOrderFailLogRepositoryfailLogRepository;RabbitListener(queues${mq.queue.order-process})publicvoidconsume(IntegerorderId,Channelchannel,Header(AmqpHeaders.DELIVERY_TAG)longtag,Header(valuex-death,requiredfalse)ListMapString,ObjectxDeath){try{orderService.processOrder(orderId);channel.basicAck(tag,false);}catch(RetryableExceptione){log.warn(订单处理临时失败将重试, orderId{},orderId,e);channel.basicNack(tag,false,true);}catch(Exceptione){log.error(订单处理不可恢复失败, orderId{},orderId,e);// 记录失败日志支持运维排查和手动重试failLogRepository.save(newOrderFailLog(orderId,e.getMessage()));// 拒绝消息不重回队列进死信channel.basicReject(tag,false);}}}十二、总结概念一句话说明ACK消费者告诉 Broker “消息已处理完可以删了”NACK消费者告诉 Broker “消息处理失败”requeueNACK 时是否让消息重新回到队列deliveryTag消息在 Channel 内的递增序号prefetchBroker 最多同时推给消费者多少条未确认消息死信队列处理失败的消息的垃圾桶可后续人工介入Publisher Confirm生产者确认消息到达 BrokerAt Least OnceACK 机制保证的语义消息至少被成功处理一次
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进