Skip to content

RabbitMQ ACK重试和死信

RabbitMQ 消费可靠性的关键在 ACK。ACK 表示消费者已经成功处理消息,Broker 可以把消息从队列中移除。如果消费者在业务没成功时就 ACK,消息就可能丢失。

这页要解决的不是“知道 ACK 这个词”,而是让你能从零解释:一条 RabbitMQ 消息为什么会丢、为什么会重复、为什么会进死信、为什么重试不能无限快、为什么消费者明明在线但 Unacked 越来越高,以及生产系统应该怎么配置、怎么写代码、怎么排查。

学习目标

学完本页你要能做到:

  1. 区分自动 ACK、手动 ACK、basicAckbasicNackbasicReject
  2. 解释 ACK 为什么必须放在业务成功之后。
  3. 解释 requeue=true 为什么可能造成失败消息立刻循环消费。
  4. 设计“业务队列 -> 重试队列 -> 原队列 -> 死信队列”的延迟重试链路。
  5. 说清死信产生条件、死信不是垃圾桶、死信为什么必须告警。
  6. 用 Spring AMQP 写出可靠消费、手动 ACK、幂等、重试和死信配置。
  7. 排查 Ready 高、Unacked 高、消费者掉线、重试风暴和死信暴涨。
  8. 在订单、支付回调、短信通知、ES 同步等商业场景中选择不同失败策略。

为什么 ACK 是消费可靠性的边界

Broker 不知道你的业务是否真的执行成功,它只能根据 ACK 判断“这条消息能不能删除”。所以 ACK 的时机必须和业务事务边界对齐:业务处理成功后 ACK,业务失败时拒绝或进入重试/死信。

mermaid
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 流程

mermaid
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["死信队列"]

这里有两个关键点:

  1. 幂等检查要在业务执行前。重复消息到来时,如果之前已经处理成功,应直接 ACK,不要再次执行业务。
  2. ACK 要在业务成功后。如果先 ACK 再写库,写库失败时 RabbitMQ 已经删除消息。

自动 ACK 和手动 ACK

模式说明风险
自动 ACK消息投递给消费者后立即认为成功消费者宕机或业务失败会丢消息
手动 ACK业务处理成功后手动确认代码稍复杂,但可靠性更高

生产环境中,重要业务建议使用手动 ACK。

自动 ACK 为什么危险

自动 ACK 可以理解成:

text
Broker 把消息发给消费者 -> 立刻认为成功 -> 从队列删除

问题在于消费者拿到消息后还没有真正执行业务。下面这些情况都会导致业务没完成但消息没了:

  1. 消费者 JVM 收到消息后宕机。
  2. 反序列化成功后写数据库失败。
  3. 调用三方接口超时。
  4. 写 ES 失败。
  5. 业务代码抛异常。

自动 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 容易造成失败消息被立即重复消费,形成死循环。

更常见的做法是:

  1. 消费失败后把消息投递到重试队列。
  2. 重试队列设置 TTL。
  3. TTL 到期后通过死信交换机回到原队列。
  4. 超过最大次数后进入真正的死信队列。
mermaid
flowchart TD
    A["业务队列"] --> B["消费失败"]
    B --> C{"超过最大次数?"}
    C -- "否" --> D["重试队列 TTL"]
    D --> E["死信交换机"]
    E --> A
    C -- "是" --> F["死信队列"]

为什么不建议无限 requeue=true

假设消费者处理订单消息时,数据库字段缺失导致每次都抛异常。如果代码里写:

java
channel.basicNack(tag, false, true);

这条消息会马上回到原队列,又马上被同一个或另一个消费者拿到,再失败,再回队列。后果是:

  1. 消费者 CPU 被失败消息占满。
  2. 正常消息被排在后面处理不了。
  3. 错误日志暴涨。
  4. 下游数据库或接口被重复打爆。
  5. ReadyUnacked 来回抖动,排查困难。

正确思路是把失败分成两类:

失败类型示例策略
可重试DB 连接短暂失败、ES 写入超时、三方接口 503延迟重试,限制次数
不可重试JSON 格式错误、必填字段缺失、业务状态非法直接死信,告警人工处理

延迟重试为什么要递增间隔

如果下游数据库已经慢了,1000 条失败消息每秒立刻重试,只会让数据库更慢。延迟重试要给下游恢复时间。

常见重试间隔:

第几次失败延迟
110 秒
21 分钟
35 分钟
430 分钟
超过上限进入死信队列

RabbitMQ 可以通过多个重试队列实现不同 TTL,也可以用延迟消息插件。没有插件时,最常见是 TTL + DLX

mermaid
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判断路由配置是否错

生产上通常会有一个“死信处理后台”或“补偿任务”:

mermaid
flowchart TD
    A["死信队列"] --> B["告警"]
    A --> C["死信消费服务"]
    C --> D["落异常消息表"]
    D --> E["人工查看失败原因"]
    E --> F{"是否可修复重放"}
    F -- "是" --> G["修复数据或代码后重新投递"]
    F -- "否" --> H["标记废弃并记录原因"]

注意:重放死信前必须确认消费者幂等,否则修复后重复投递可能产生重复扣款、重复通知、重复入库。

消费幂等

只要使用重试,就必须考虑重复消费。RabbitMQ 也无法保证业务只执行一次。

常见幂等方式:

  • 业务唯一键加数据库唯一索引。
  • 消息 ID 去重表。
  • 状态机流转判断。
  • Redis 短期去重。

例如订单关闭消息重复到达时,可以先判断订单是否仍是“待支付”。如果已经支付或已经关闭,则直接 ACK。

幂等为什么是必须的

重复消费可能来自:

  1. 消费者业务成功后,ACK 前宕机。
  2. ACK 网络包丢失,Broker 以为没成功。
  3. 消费超时后 Broker 重新投递。
  4. 消费者连接断开,未 ACK 消息重新入队。
  5. 失败重试和死信重放。

所以不要问“RabbitMQ 能不能保证只消费一次”。工程上应该问:重复投递发生时,我的业务会不会重复产生副作用。

幂等表示例

sql
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)
);

消费时先插入消费记录,依赖唯一索引防重:

java
@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 高的排查链路

mermaid
flowchart TD
    A["Ready 高"] --> B{"消费者是否在线"}
    B -- "否" --> C["检查消费者进程、连接、权限、队列名"]
    B -- "是" --> D{"投递速率是否低"}
    D -- "是" --> E["检查 prefetch、网络、Broker 负载"]
    D -- "否" --> F["检查消费者处理速度和 ACK"]

Unacked 高的排查链路

mermaid
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 配置示例:

yaml
spring:
  rabbitmq:
    listener:
      simple:
        acknowledge-mode: manual
        prefetch: 20
        concurrency: 4
        max-concurrency: 12

完整 Demo:订单支付消息可靠消费

下面用订单支付成功消息演示生产常用配置。

1. 队列和交换机配置

java
@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. 消息体设计

java
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. 消费者代码

java
@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);
        }
    }
}

这段代码的关键思路:

  1. 业务成功才 ACK。
  2. 不可重试错误进入死信。
  3. 可重试错误投递到重试队列后 ACK 原消息,避免原消息反复立即投递。
  4. 未预期异常 basicNack(requeue=false),让原消息按队列 DLX 进入死信,避免无限循环。
  5. 业务服务内部必须做幂等。

4. 业务幂等代码

java
@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");
    }
}

订单更新也要带状态条件:

sql
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;失败时拒绝消息并不重新入队,让它进入死信队列。

java
@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 如何保证消费可靠性

text
RabbitMQ 消费可靠性的边界是 ACK。Broker 不知道业务是否成功,只能根据消费者 ACK 判断消息能否删除。生产上核心业务会使用手动 ACK,业务处理成功并完成幂等落库后再 basicAck;可重试异常进入延迟重试队列,不可重试异常进入死信队列并告警。因为 ACK 失败、消费者宕机、重试和死信重放都可能导致重复消费,所以消费端必须用消息 ID、业务唯一键、消费日志表、状态机或唯一索引保证幂等。

RabbitMQ 死信队列有什么用

text
死信队列用于隔离无法正常消费的异常消息,避免毒消息反复重试拖垮主消费链路。消息被 reject 且 requeue=false、过期、队列满或超过重试次数时可以进入死信队列。死信不是垃圾桶,生产上要对死信数量告警,记录原始消息、业务主键、失败原因和重试次数,修复后按幂等规则重放或人工补偿。

Ready 高和 Unacked 高怎么排查

text
Ready 高说明消息还在队列里等待投递,优先看消费者是否在线、消费实例是否足够、prefetch 是否太小、Broker 投递是否正常。Unacked 高说明消息已经投递给消费者但还没 ACK,优先看业务耗时、线程池、DB/ES/HTTP 下游、ACK 代码分支和 prefetch 是否太大。排查时不能只重启消费者,因为重启会让未 ACK 消息重新入队并造成重复消费。

关联知识点

知识点跳转
RabbitMQ 模型RabbitMQ 总览
Exchange 路由交换机和路由
延迟任务延迟队列
MQ 堆积和背压消息堆积与背压
MQ 面试消息队列面试题

本章小结

RabbitMQ 可靠消费不是一个 ACK 开关,而是一条链:手动 ACK 决定消息什么时候删除,延迟重试避免短暂故障放大,死信队列隔离毒消息,幂等保证重复投递不会造成重复副作用,监控排查负责发现 Ready、Unacked、死信和消费耗时异常。真正的生产级方案一定同时包含代码、配置、监控、告警和补偿。