一、ack 是什么
ack(acknowledge)是消费者告诉 broker “这条消息我处理完了,你可以删除了” 的信号。没有 ack,broker 不知道消息是否被成功处理,也就无法决定是否删除消息。
producer → broker(queue) → consumer
│
├─ 处理成功 → ack → broker 删除消息
│
├─ 处理失败 → nack/reject → broker 重新入队或进死信
│
└─ 消费者宕机 → 超时无 ack → broker 重新分发给其他消费者二、为什么需要 ack
没有 ack 机制会怎样:
| 场景 | 无 ack 的后果 |
|---|---|
| 消费者收到消息后处理到一半崩溃 | 消息丢失,业务不完整 |
| 消费者处理失败(业务异常) | 无法重试,数据不一致 |
| 网络闪断,消费者没收到消息 | broker 以为已推送,消息丢失 |
ack 机制的核心保证:消息至少被成功处理一次(at least once)。
三、三种确认模式
1. 自动确认(auto)
spring:
rabbitmq:
listener:
simple:
acknowledge-mode: auto行为:
- 消息推送到消费者方法后,spring 根据方法执行结果自动决定:
- 方法正常返回 → 自动 ack
- 方法抛异常 → 自动 nack + requeue
@rabbitlistener(queues = "order.queue")
public void consume(integer orderid) {
orderservice.process(orderid);
// 正常返回 → 框架自动 ack
// 抛异常 → 框架自动 nack,消息重新入队
}优点:代码简单,不需要手动管理 缺点:异常时无限 requeue 可能导致消息循环(需配合重试策略)
2. 手动确认(manual)
spring:
rabbitmq:
listener:
simple:
acknowledge-mode: manual行为:
- 消费者必须在代码中显式调用
basicack或basicnack - 不调用则消息一直处于 unacked 状态,消费者断连后消息重新入队
@rabbitlistener(queues = "order.queue")
public void consume(integer orderid, channel channel,
@header(amqpheaders.delivery_tag) long deliverytag) {
try {
orderservice.process(orderid);
// 处理成功,确认消息
channel.basicack(deliverytag, false);
} catch (businessexception e) {
// 业务异常,不重试,直接丢弃(或进死信)
channel.basicreject(deliverytag, false);
} catch (exception e) {
// 临时异常,重新入队等待重试
channel.basicnack(deliverytag, false, true);
}
}优点:精细控制,可以区分不同异常做不同处理 缺点:代码复杂度增加,忘记 ack 会导致消息堆积
3. 无确认(none)
spring:
rabbitmq:
listener:
simple:
acknowledge-mode: none行为:
- broker 推送后立即删除消息,不等待消费者确认
- 等同于"发后即忘"
优点:最高吞吐 缺点:消息可能丢失,仅适用于可丢失的场景(如日志采集、监控指标)
四、ack 相关的三个操作
basicack — 确认成功
channel.basicack(deliverytag, multiple);
| 参数 | 说明 |
|---|---|
| deliverytag | 消息的唯一标识(broker 分配,单 channel 内递增) |
| multiple | false=只确认当前消息;true=确认 deliverytag 及之前所有未确认的消息 |
// 逐条确认 channel.basicack(deliverytag, false); // 批量确认(确认 tag<=5 的所有消息) channel.basicack(5, true);
basicnack — 否定确认(批量)
channel.basicnack(deliverytag, multiple, requeue);
| 参数 | 说明 |
|---|---|
| deliverytag | 消息标识 |
| multiple | 是否批量否定 |
| requeue | true=重新入队;false=丢弃或进死信 |
// 单条否定,重新入队 channel.basicnack(deliverytag, false, true); // 单条否定,不重回队列(进入死信队列或丢弃) channel.basicnack(deliverytag, false, false);
basicreject — 否定确认(单条)
channel.basicreject(deliverytag, requeue);
功能与 basicnack 相同,但只能操作单条消息(没有 multiple 参数)。
五、deliverytag 详解
channel 1: msg_a(tag=1) msg_b(tag=2) msg_c(tag=3) channel 2: msg_x(tag=1) msg_y(tag=2)
- deliverytag 是每个 channel 独立递增的序号
- 消费者收到消息时携带 tag,ack 时回传这个 tag
- broker 通过 (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 模式下忘记 ack
@rabbitlistener(queues = "order.queue")
public void consume(integer orderid, channel channel,
@header(amqpheaders.delivery_tag) long tag) {
orderservice.process(orderid);
// 忘记调用 channel.basicack(tag, false);
// 后果:消息一直 unacked,达到 prefetch 后不再推送新消息
}排查方式:rabbitmq 管理后台看 queue 的 unacked 数量持续增长。
陷阱2:异常时 requeue 导致无限循环
// 错误示例:所有异常都 requeue
@rabbitlistener(queues = "order.queue")
public void consume(integer orderid, channel channel,
@header(amqpheaders.delivery_tag) long tag) {
try {
orderservice.process(orderid);
channel.basicack(tag, false);
} catch (exception e) {
// 如果是参数错误(永远不会成功),每次 requeue 都会再失败
channel.basicnack(tag, false, true); // 无限循环!
}
}正确做法:区分可重试异常和不可重试异常:
try {
orderservice.process(orderid);
channel.basicack(tag, false);
} catch (retryableexception e) {
// 临时故障(网络超时、锁冲突),重新入队
channel.basicnack(tag, false, true);
} catch (exception e) {
// 不可恢复(参数错误、数据不存在),拒绝不重回
log.error("消费失败且不可重试, orderid={}", orderid, e);
channel.basicreject(tag, false); // 进死信或丢弃
}
陷阱3:批量确认丢消息
// 风险示例:multiple=true 批量确认 channel.basicack(tag, true); // 如果 tag=5,则 tag 1~5 全部确认 // 如果 tag=3 的消息其实还没处理完,也会被确认掉
建议:除非有明确的批量处理逻辑,否则始终用 multiple=false。
八、重试策略配置
spring 内置重试(auto 模式下)
spring:
rabbitmq:
listener:
simple:
acknowledge-mode: auto
retry:
enabled: true
initial-interval: 1000 # 第一次重试间隔 1s
max-interval: 10000 # 最大间隔 10s
multiplier: 2.0 # 间隔倍增
max-attempts: 3 # 最大重试次数流程:
第1次消费失败 → 等1s → 第2次重试 → 等2s → 第3次重试 → 仍失败 → 进入 messagerecoverer
自定义失败处理器
@bean
public messagerecoverer messagerecoverer(rabbittemplate rabbittemplate) {
// 重试耗尽后发送到死信队列
return new republishmessagerecoverer(rabbittemplate, "dlx.exchange", "dlx.order");
}
死信队列(dlx)配置
@bean
public queue orderqueue() {
map<string, object> args = new hashmap<>();
args.put("x-dead-letter-exchange", "dlx.exchange"); // 死信交换机
args.put("x-dead-letter-routing-key", "dlx.order"); // 死信路由键
return new queue("order.queue", true, false, false, args);
}
@bean
public queue deadletterqueue() {
return new queue("order.dlq", true);
}
@bean
public binding dlqbinding() {
return bindingbuilder.bind(deadletterqueue())
.to(new directexchange("dlx.exchange"))
.with("dlx.order");
}消息进入死信队列的条件:
- 被 reject/nack 且 requeue=false
- 消息 ttl 过期
- 队列达到最大长度
九、生产者确认(publisher confirm)
ack 不仅消费端有,生产端也有——确认消息成功到达 broker:
spring:
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仅用于明确标注"允许丢失"的非关键数据
十一、完整示例:可靠消费模板
@component
@slf4j
public class ordermqconsumer {
@resource private orderservice orderservice;
@resource private orderfaillogrepository faillogrepository;
@rabbitlistener(queues = "${mq.queue.order-process}")
public void consume(integer orderid, channel channel,
@header(amqpheaders.delivery_tag) long tag,
@header(value = "x-death", required = false) list<map<string, object>> xdeath) {
try {
orderservice.processorder(orderid);
channel.basicack(tag, false);
} catch (retryableexception e) {
log.warn("订单处理临时失败将重试, orderid={}", orderid, e);
channel.basicnack(tag, false, true);
} catch (exception e) {
log.error("订单处理不可恢复失败, orderid={}", orderid, e);
// 记录失败日志,支持运维排查和手动重试
faillogrepository.save(new orderfaillog(orderid, e.getmessage()));
// 拒绝消息,不重回队列(进死信)
channel.basicreject(tag, false);
}
}
}十二、总结
| 概念 | 一句话说明 |
|---|---|
| ack | 消费者告诉 broker “消息已处理完,可以删了” |
| nack | 消费者告诉 broker “消息处理失败” |
| requeue | nack 时是否让消息重新回到队列 |
| deliverytag | 消息在 channel 内的递增序号 |
| prefetch | broker 最多同时推给消费者多少条未确认消息 |
| 死信队列 | 处理失败的消息的"垃圾桶",可后续人工介入 |
| publisher confirm | 生产者确认消息到达 broker |
| at least once | ack 机制保证的语义:消息至少被成功处理一次 |
到此这篇关于rabbitmq 消息确认机制(ack):技术解析与实践指南的文章就介绍到这了,更多相关rabbitmq 消息确认机制内容请搜索代码网以前的文章或继续浏览下面的相关文章希望大家以后多多支持代码网!
发表评论