RabbitMQ ACK重试和死信
RabbitMQ 消费可靠性的关键在 ACK。ACK 表示消费者已经成功处理消息,Broker 可以把消息从队列中移除。如果消费者在业务没成功时就 ACK,消息就可能丢失。
这页要解决的不是“知道 ACK 这个词”,而是让你能从零解释:一条 RabbitMQ 消息为什么会丢、为什么会重复、为什么会进死信、为什么重试不能无限快、为什么消费者明明在线但 Unacked 越来越高,以及生产系统应该怎么配置、怎么写代码、怎么排查。
学习目标
学完本页你要能做到:
- 区分自动 ACK、手动 ACK、
basicAck、basicNack、basicReject。 - 解释 ACK 为什么必须放在业务成功之后。
- 解释
requeue=true为什么可能造成失败消息立刻循环消费。 - 设计“业务队列 -> 重试队列 -> 原队列 -> 死信队列”的延迟重试链路。
- 说清死信产生条件、死信不是垃圾桶、死信为什么必须告警。
- 用 Spring AMQP 写出可靠消费、手动 ACK、幂等、重试和死信配置。
- 排查
Ready高、Unacked高、消费者掉线、重试风暴和死信暴涨。 - 在订单、支付回调、短信通知、ES 同步等商业场景中选择不同失败策略。
为什么 ACK 是消费可靠性的边界
Broker 不知道你的业务是否真的执行成功,它只能根据 ACK 判断“这条消息能不能删除”。所以 ACK 的时机必须和业务事务边界对齐:业务处理成功后 ACK,业务失败时拒绝或进入重试/死信。
flowchart TD
A["消息投递给 Consumer"] --> B["执行业务逻辑"]
B --> C{"业务是否真正成功"}
C -- "成功" --> D["basicAck<br/>Broker 删除消息"]
C -- "失败" --> E["basicNack / basicReject"]
E --> F["重试或死信"]如果提前 ACK,后续业务失败时 Broker 已经删除消息,只能靠人工补偿;如果一直不 ACK,消息会堆积或被重复投递。
更准确地说,RabbitMQ 只知道两件事:
| Broker 能知道 | Broker 不知道 |
|---|---|
| 消息是否投递给了消费者 | 你的订单表是否更新成功 |
| 消费者有没有 ACK | 你的 ES 是否写入成功 |
| 消费者连接是否断开 | 你的接口是否完成幂等处理 |
| 消息是否过期或被拒绝 | 这条消息对业务是否还能重试 |
所以可靠消费一定要把“业务成功”和“ACK”绑定起来。如果业务没成功就 ACK,消息丢;如果业务成功但 ACK 失败,消息可能重复投递。因此 RabbitMQ 消费端的真实语义更接近:至少处理一次,业务自己保证幂等。
ACK 流程
flowchart TD
A["Broker 投递消息"] --> B["Consumer 收到消息"]
B --> C["解析消息并做幂等检查"]
C --> D{"是否已经处理过"}
D -- "是" --> E["直接 basicAck"]
D -- "否" --> F["执行业务逻辑"]
F --> G{"业务是否成功"}
G -- "成功" --> H["记录成功状态"]
H --> I["basicAck"]
G -- "可重试失败" --> J["basicNack 或投递重试队列"]
G -- "不可重试失败" --> K["basicReject requeue=false"]
J --> L["延迟重试"]
K --> M["死信队列"]这里有两个关键点:
- 幂等检查要在业务执行前。重复消息到来时,如果之前已经处理成功,应直接 ACK,不要再次执行业务。
- ACK 要在业务成功后。如果先 ACK 再写库,写库失败时 RabbitMQ 已经删除消息。
自动 ACK 和手动 ACK
| 模式 | 说明 | 风险 |
|---|---|---|
| 自动 ACK | 消息投递给消费者后立即认为成功 | 消费者宕机或业务失败会丢消息 |
| 手动 ACK | 业务处理成功后手动确认 | 代码稍复杂,但可靠性更高 |
生产环境中,重要业务建议使用手动 ACK。
自动 ACK 为什么危险
自动 ACK 可以理解成:
Broker 把消息发给消费者 -> 立刻认为成功 -> 从队列删除问题在于消费者拿到消息后还没有真正执行业务。下面这些情况都会导致业务没完成但消息没了:
- 消费者 JVM 收到消息后宕机。
- 反序列化成功后写数据库失败。
- 调用三方接口超时。
- 写 ES 失败。
- 业务代码抛异常。
自动 ACK 只适合不重要、可丢弃、可重建的消息,例如某些低价值日志采样。订单、支付、库存、资产采集入库这类业务不要用自动 ACK。
手动 ACK 的三个动作
| 方法 | 含义 | 是否支持批量 | 常见使用 |
|---|---|---|---|
basicAck(tag, multiple) | 确认成功,Broker 删除消息 | 支持 | 业务成功后 |
basicNack(tag, multiple, requeue) | 拒绝,可批量,可选择重新入队 | 支持 | 失败重试或拒绝 |
basicReject(tag, requeue) | 拒绝单条,可选择重新入队 | 不支持批量 | 单条失败进入死信 |
multiple=true 表示确认当前 delivery tag 及之前未确认的消息。生产代码里如果对 delivery tag 管理不严,批量 ACK 可能误确认还没处理完的消息。初学和普通业务建议先用 multiple=false。
requeue=true 表示重新放回队列;requeue=false 表示不重新入队,如果队列配置了死信交换机,就会进入死信链路。
重试策略
RabbitMQ 本身可以通过 basicNack 让消息重新入队,但直接 requeue 容易造成失败消息被立即重复消费,形成死循环。
更常见的做法是:
- 消费失败后把消息投递到重试队列。
- 重试队列设置 TTL。
- TTL 到期后通过死信交换机回到原队列。
- 超过最大次数后进入真正的死信队列。
flowchart TD
A["业务队列"] --> B["消费失败"]
B --> C{"超过最大次数?"}
C -- "否" --> D["重试队列 TTL"]
D --> E["死信交换机"]
E --> A
C -- "是" --> F["死信队列"]为什么不建议无限 requeue=true
假设消费者处理订单消息时,数据库字段缺失导致每次都抛异常。如果代码里写:
channel.basicNack(tag, false, true);这条消息会马上回到原队列,又马上被同一个或另一个消费者拿到,再失败,再回队列。后果是:
- 消费者 CPU 被失败消息占满。
- 正常消息被排在后面处理不了。
- 错误日志暴涨。
- 下游数据库或接口被重复打爆。
Ready和Unacked来回抖动,排查困难。
正确思路是把失败分成两类:
| 失败类型 | 示例 | 策略 |
|---|---|---|
| 可重试 | DB 连接短暂失败、ES 写入超时、三方接口 503 | 延迟重试,限制次数 |
| 不可重试 | JSON 格式错误、必填字段缺失、业务状态非法 | 直接死信,告警人工处理 |
延迟重试为什么要递增间隔
如果下游数据库已经慢了,1000 条失败消息每秒立刻重试,只会让数据库更慢。延迟重试要给下游恢复时间。
常见重试间隔:
| 第几次失败 | 延迟 |
|---|---|
| 1 | 10 秒 |
| 2 | 1 分钟 |
| 3 | 5 分钟 |
| 4 | 30 分钟 |
| 超过上限 | 进入死信队列 |
RabbitMQ 可以通过多个重试队列实现不同 TTL,也可以用延迟消息插件。没有插件时,最常见是 TTL + DLX。
flowchart TD
A["order.paid.queue<br/>业务队列"] --> B["Consumer 消费"]
B --> C{"处理成功"}
C -- "成功" --> D["basicAck"]
C -- "失败且未超限" --> E["发送到 retry.10s.queue"]
E --> F["TTL 到期"]
F --> G["DLX 路由回业务交换机"]
G --> A
C -- "失败且超限" --> H["order.paid.dlq<br/>死信队列"]死信产生条件
消息进入死信队列常见原因:
- 消息被拒绝,并且
requeue=false。 - 消息过期。
- 队列达到最大长度。
- 消费重试次数超过业务限制。
死信队列不是垃圾桶,而是异常消息的观察窗口。重要业务必须对死信数量做监控。
死信队列应该怎么用
死信队列的价值不是“把异常消息放一边不管”,而是把异常从主消费链路隔离出来,让正常消息继续处理,同时给排查和补偿留下证据。
死信消息至少要能回答:
| 问题 | 为什么重要 |
|---|---|
| 哪个业务主键失败 | 能定位订单、支付单、资产编号 |
| 原始消息是什么 | 能判断消息格式和字段是否正确 |
| 失败原因是什么 | 区分代码 Bug、数据问题、下游故障 |
| 失败了几次 | 判断是否达到重试上限 |
| 最后失败时间 | 判断影响范围和恢复窗口 |
| 来自哪个 exchange / routing key | 判断路由配置是否错 |
生产上通常会有一个“死信处理后台”或“补偿任务”:
flowchart TD
A["死信队列"] --> B["告警"]
A --> C["死信消费服务"]
C --> D["落异常消息表"]
D --> E["人工查看失败原因"]
E --> F{"是否可修复重放"}
F -- "是" --> G["修复数据或代码后重新投递"]
F -- "否" --> H["标记废弃并记录原因"]注意:重放死信前必须确认消费者幂等,否则修复后重复投递可能产生重复扣款、重复通知、重复入库。
消费幂等
只要使用重试,就必须考虑重复消费。RabbitMQ 也无法保证业务只执行一次。
常见幂等方式:
- 业务唯一键加数据库唯一索引。
- 消息 ID 去重表。
- 状态机流转判断。
- Redis 短期去重。
例如订单关闭消息重复到达时,可以先判断订单是否仍是“待支付”。如果已经支付或已经关闭,则直接 ACK。
幂等为什么是必须的
重复消费可能来自:
- 消费者业务成功后,ACK 前宕机。
- ACK 网络包丢失,Broker 以为没成功。
- 消费超时后 Broker 重新投递。
- 消费者连接断开,未 ACK 消息重新入队。
- 失败重试和死信重放。
所以不要问“RabbitMQ 能不能保证只消费一次”。工程上应该问:重复投递发生时,我的业务会不会重复产生副作用。
幂等表示例
create table mq_consume_log (
id bigint primary key auto_increment,
message_id varchar(128) not null,
business_key varchar(128) not null,
consumer_group varchar(128) not null,
status varchar(32) not null,
error_msg varchar(1000),
created_at datetime not null,
updated_at datetime not null,
unique key uk_msg_group (message_id, consumer_group)
);消费时先插入消费记录,依赖唯一索引防重:
@Transactional
public boolean handleOrderPaid(OrderPaidEvent event) {
boolean firstConsume = consumeLogRepository.tryInsert(
event.getMessageId(),
event.getOrderNo(),
"order-paid-consumer"
);
if (!firstConsume) {
return true;
}
orderService.markPaid(event.getOrderNo(), event.getPayTime());
consumeLogRepository.markSuccess(event.getMessageId(), "order-paid-consumer");
return true;
}这里的关键不是代码长什么样,而是唯一索引把并发重复消息挡住。即使两个消费者同时拿到重复消息,也只有一个能插入成功。
排查建议
消费失败时日志至少包含:
- exchange
- routing key
- queue
- message id
- delivery tag
- 业务主键
- 异常堆栈
否则进入死信后很难判断是消息问题、代码问题还是下游系统问题。
生产排查:Ready 和 Unacked 怎么看
RabbitMQ 堆积排查最常看两个指标:
| 指标 | 含义 | 常见原因 |
|---|---|---|
Ready 高 | 消息还在队列中,等待投递 | 消费者数量不足、消费者掉线、prefetch 太小、路由到错队列 |
Unacked 高 | 消息已投递给消费者,但还没确认 | 业务处理慢、ACK 没执行、线程池卡住、prefetch 太大 |
Ready 高的排查链路
flowchart TD
A["Ready 高"] --> B{"消费者是否在线"}
B -- "否" --> C["检查消费者进程、连接、权限、队列名"]
B -- "是" --> D{"投递速率是否低"}
D -- "是" --> E["检查 prefetch、网络、Broker 负载"]
D -- "否" --> F["检查消费者处理速度和 ACK"]Unacked 高的排查链路
flowchart TD
A["Unacked 高"] --> B["说明消息已到消费者"]
B --> C{"业务耗时是否升高"}
C -- "是" --> D["查 DB、ES、HTTP、锁等待"]
C -- "否" --> E{"ACK 是否执行"}
E -- "否" --> F["查异常分支、线程池、代码逻辑"]
E -- "是" --> G["查网络、连接、channel 异常"]Unacked 高时不要第一反应重启消费者。重启会让未 ACK 消息重新入队,短时间可能造成重复消费和流量抖动。应该先看线程栈、业务耗时、错误日志、连接池、下游指标。
prefetch 为什么重要
prefetch 控制一个消费者最多能拿多少条未 ACK 消息。
| 配置 | 结果 | 风险 |
|---|---|---|
| 太大 | 一个消费者一次拿很多消息 | Unacked 高,其他消费者分不到,宕机后大量重投 |
| 太小 | 消费者每次拿很少消息 | 吞吐上不去,网络往返多 |
| 合理 | 按处理耗时和并发能力控制 | 既有吞吐,又不会本地堆积 |
例如单条消息平均处理 100ms,一个消费者容器并发 10,prefetch 可以从 10 到 50 之间压测。不要无脑设置几千,否则 Broker 中看似没多少 Ready,实际消息都堆在消费者本地未确认区。
Spring Boot 配置示例:
spring:
rabbitmq:
listener:
simple:
acknowledge-mode: manual
prefetch: 20
concurrency: 4
max-concurrency: 12完整 Demo:订单支付消息可靠消费
下面用订单支付成功消息演示生产常用配置。
1. 队列和交换机配置
@Configuration
public class RabbitOrderConfig {
public static final String ORDER_EXCHANGE = "order.exchange";
public static final String ORDER_PAID_QUEUE = "order.paid.queue";
public static final String ORDER_PAID_RETRY_QUEUE = "order.paid.retry.10s.queue";
public static final String ORDER_PAID_DLQ = "order.paid.dlq";
@Bean
public DirectExchange orderExchange() {
return ExchangeBuilder.directExchange(ORDER_EXCHANGE)
.durable(true)
.build();
}
@Bean
public Queue orderPaidQueue() {
return QueueBuilder.durable(ORDER_PAID_QUEUE)
.withArgument("x-dead-letter-exchange", ORDER_EXCHANGE)
.withArgument("x-dead-letter-routing-key", "order.paid.dlq")
.build();
}
@Bean
public Queue orderPaidRetryQueue() {
return QueueBuilder.durable(ORDER_PAID_RETRY_QUEUE)
.withArgument("x-message-ttl", 10_000)
.withArgument("x-dead-letter-exchange", ORDER_EXCHANGE)
.withArgument("x-dead-letter-routing-key", "order.paid")
.build();
}
@Bean
public Queue orderPaidDlq() {
return QueueBuilder.durable(ORDER_PAID_DLQ).build();
}
@Bean
public Binding orderPaidBinding() {
return BindingBuilder.bind(orderPaidQueue())
.to(orderExchange())
.with("order.paid");
}
@Bean
public Binding orderPaidRetryBinding() {
return BindingBuilder.bind(orderPaidRetryQueue())
.to(orderExchange())
.with("order.paid.retry.10s");
}
@Bean
public Binding orderPaidDlqBinding() {
return BindingBuilder.bind(orderPaidDlq())
.to(orderExchange())
.with("order.paid.dlq");
}
}这个配置里有三个队列:
| 队列 | 作用 |
|---|---|
order.paid.queue | 正常消费订单支付消息 |
order.paid.retry.10s.queue | 失败后延迟 10 秒再回原队列 |
order.paid.dlq | 超过重试次数或不可重试异常进入死信 |
2. 消息体设计
public class OrderPaidEvent {
private String messageId;
private String orderNo;
private String payNo;
private BigDecimal amount;
private LocalDateTime paidAt;
private Integer retryCount;
// getter/setter
}messageId 用于幂等,orderNo 是业务主键,retryCount 用于控制最大重试次数。实际项目也可以把重试次数放到 header。
3. 消费者代码
@Component
public class OrderPaidConsumer {
private final ObjectMapper objectMapper;
private final RabbitTemplate rabbitTemplate;
private final OrderPaidService orderPaidService;
public OrderPaidConsumer(
ObjectMapper objectMapper,
RabbitTemplate rabbitTemplate,
OrderPaidService orderPaidService
) {
this.objectMapper = objectMapper;
this.rabbitTemplate = rabbitTemplate;
this.orderPaidService = orderPaidService;
}
@RabbitListener(queues = RabbitOrderConfig.ORDER_PAID_QUEUE, ackMode = "MANUAL")
public void onMessage(Message message, Channel channel) throws IOException {
long tag = message.getMessageProperties().getDeliveryTag();
String body = new String(message.getBody(), StandardCharsets.UTF_8);
try {
OrderPaidEvent event = objectMapper.readValue(body, OrderPaidEvent.class);
orderPaidService.handlePaidEvent(event);
channel.basicAck(tag, false);
} catch (InvalidMessageException e) {
sendToDlq(body, e.getMessage());
channel.basicAck(tag, false);
} catch (TemporaryBusinessException e) {
retryOrDlq(body, e.getMessage());
channel.basicAck(tag, false);
} catch (Exception e) {
channel.basicNack(tag, false, false);
}
}
private void retryOrDlq(String body, String reason) {
OrderPaidEvent event = parse(body);
int retryCount = event.getRetryCount() == null ? 0 : event.getRetryCount();
if (retryCount >= 3) {
sendToDlq(body, reason);
return;
}
event.setRetryCount(retryCount + 1);
rabbitTemplate.convertAndSend(
RabbitOrderConfig.ORDER_EXCHANGE,
"order.paid.retry.10s",
event
);
}
private void sendToDlq(String body, String reason) {
rabbitTemplate.convertAndSend(
RabbitOrderConfig.ORDER_EXCHANGE,
"order.paid.dlq",
Map.of("body", body, "reason", reason, "time", LocalDateTime.now().toString())
);
}
private OrderPaidEvent parse(String body) {
try {
return objectMapper.readValue(body, OrderPaidEvent.class);
} catch (Exception e) {
throw new InvalidMessageException("message parse failed", e);
}
}
}这段代码的关键思路:
- 业务成功才 ACK。
- 不可重试错误进入死信。
- 可重试错误投递到重试队列后 ACK 原消息,避免原消息反复立即投递。
- 未预期异常
basicNack(requeue=false),让原消息按队列 DLX 进入死信,避免无限循环。 - 业务服务内部必须做幂等。
4. 业务幂等代码
@Service
public class OrderPaidService {
private final ConsumeLogMapper consumeLogMapper;
private final OrderMapper orderMapper;
@Transactional
public void handlePaidEvent(OrderPaidEvent event) {
int inserted = consumeLogMapper.insertIgnore(
event.getMessageId(),
"order-paid-consumer",
event.getOrderNo()
);
if (inserted == 0) {
return;
}
int updated = orderMapper.markPaid(
event.getOrderNo(),
event.getPayNo(),
event.getAmount(),
event.getPaidAt()
);
if (updated == 0) {
throw new InvalidMessageException("order status cannot mark paid");
}
consumeLogMapper.markSuccess(event.getMessageId(), "order-paid-consumer");
}
}订单更新也要带状态条件:
update t_order
set status = 'PAID',
pay_no = ?,
paid_at = ?,
updated_at = now()
where order_no = ?
and status = 'WAIT_PAY';这样重复的支付成功消息不会把已取消、已关闭、已退款订单乱改。
商业场景怎么选策略
| 场景 | 是否允许丢 | 失败策略 |
|---|---|---|
| 支付成功回调 | 不允许 | 手动 ACK、幂等、延迟重试、死信告警、人工补偿 |
| 订单同步 ES | 不建议丢,但可重建 | 重试、死信、后续按主库补偿重建索引 |
| 短信通知 | 部分可降级 | 重试几次后死信,人工或批量补发 |
| 操作日志 | 可按等级区分 | 重要审计不丢,普通行为日志可采样或降级 |
| 缓存刷新 | 可重建 | 失败后短重试,也可以依赖 TTL 自动恢复 |
商业系统里不要只说“失败进死信”。你必须回答:死信后谁看、怎么看、怎么修、能不能重放、重放会不会重复执行。
代码 Demo:手动 ACK 消费
下面示例中,业务成功才 basicAck;失败时拒绝消息并不重新入队,让它进入死信队列。
@RabbitListener(queues = "order.paid.queue", ackMode = "MANUAL")
public void onMessage(Message message, Channel channel) throws IOException {
long tag = message.getMessageProperties().getDeliveryTag();
try {
String body = new String(message.getBody(), StandardCharsets.UTF_8);
handleOrderPaid(body);
channel.basicAck(tag, false);
} catch (Exception e) {
channel.basicReject(tag, false);
}
}如果失败原因是数据库短暂不可用,可以进入重试队列;如果是消息格式错误,应该进入死信并报警。
常见坑
| 坑 | 后果 | 正确做法 |
|---|---|---|
| 自动 ACK 处理核心消息 | 业务失败但消息已删除 | 手动 ACK |
失败后无限 requeue=true | 重试风暴、正常消息被饿死 | 延迟重试 + 最大次数 + 死信 |
| 没有幂等 | 重试或重放导致重复扣款、重复通知 | 消息 ID、业务唯一键、状态机 |
| prefetch 太大 | 大量消息堆在消费者 Unacked | 按并发和耗时压测 |
| 死信没人处理 | 异常消息长期堆积,业务缺口没人知道 | 死信告警、异常表、补偿后台 |
| 只记录异常不记录业务主键 | 无法定位影响数据 | 日志带 messageId、orderNo、queue |
| ACK 放在事务提交前 | 事务回滚但消息删除 | 事务成功后 ACK |
面试标准回答
RabbitMQ 如何保证消费可靠性
RabbitMQ 消费可靠性的边界是 ACK。Broker 不知道业务是否成功,只能根据消费者 ACK 判断消息能否删除。生产上核心业务会使用手动 ACK,业务处理成功并完成幂等落库后再 basicAck;可重试异常进入延迟重试队列,不可重试异常进入死信队列并告警。因为 ACK 失败、消费者宕机、重试和死信重放都可能导致重复消费,所以消费端必须用消息 ID、业务唯一键、消费日志表、状态机或唯一索引保证幂等。RabbitMQ 死信队列有什么用
死信队列用于隔离无法正常消费的异常消息,避免毒消息反复重试拖垮主消费链路。消息被 reject 且 requeue=false、过期、队列满或超过重试次数时可以进入死信队列。死信不是垃圾桶,生产上要对死信数量告警,记录原始消息、业务主键、失败原因和重试次数,修复后按幂等规则重放或人工补偿。Ready 高和 Unacked 高怎么排查
Ready 高说明消息还在队列里等待投递,优先看消费者是否在线、消费实例是否足够、prefetch 是否太小、Broker 投递是否正常。Unacked 高说明消息已经投递给消费者但还没 ACK,优先看业务耗时、线程池、DB/ES/HTTP 下游、ACK 代码分支和 prefetch 是否太大。排查时不能只重启消费者,因为重启会让未 ACK 消息重新入队并造成重复消费。关联知识点
| 知识点 | 跳转 |
|---|---|
| RabbitMQ 模型 | RabbitMQ 总览 |
| Exchange 路由 | 交换机和路由 |
| 延迟任务 | 延迟队列 |
| MQ 堆积和背压 | 消息堆积与背压 |
| MQ 面试 | 消息队列面试题 |
本章小结
RabbitMQ 可靠消费不是一个 ACK 开关,而是一条链:手动 ACK 决定消息什么时候删除,延迟重试避免短暂故障放大,死信队列隔离毒消息,幂等保证重复投递不会造成重复副作用,监控排查负责发现 Ready、Unacked、死信和消费耗时异常。真正的生产级方案一定同时包含代码、配置、监控、告警和补偿。
