Files
autumn-recruitment/04.MQ/rabbitmq/ACK 确认与死信队列.md

7.6 KiB
Raw Permalink Blame History

tags, create time, update time
tags create time update time
mq/rabbitmq
manual-ack
dead-letter-exchange
dlq
retry-pattern
prefetch
2026-08-08 18:52 2026-08-08 18:52

ACK 确认与死信队列

概述

Consumer Acknowledgment(ACK)确认机制是 RabbitMQ 消费者可靠性保障的最后防线。Auto ACK 虽然高效但一旦消费者崩溃就会丢消息;Manual ACK 提供了精细的控制能力但增加了复杂度。配合死信队列(DLQ)和重试模式,可以构建出完整的异常消息处理流水线。

核心原理

Auto ACK vs Manual ACK

RabbitMQ 的消费确认模式分为两类:

特性 Auto ACK Manual ACK
确认时机 消息交付即确认 消费者显式调用 basic.ack
安全性 低(消费者 crash 会丢未处理消息) 高(处理完才确认)
吞吐量 略高 略低(额外的网络往返)
适用场景 无状态处理、容忍少量丢失 持久化处理、金融类业务

Auto ACK 的隐患在于:AMQP 协议层面,Basic.GetOk / Basic.ConsumeOk 本身就是"投递动作",而 Auto ACK 把"投递"等同于"处理成功"。如果消费者在收到消息后尚未执行业务逻辑就宕机了,这条消息已经被认为被消费完了。

Manual ACK 的正确使用方式是在业务逻辑执行完毕后发送 basic.ack,如果出现异常则发送 basic.nack 并指定 requeue 策略。

// Go 代码示例:Manual ACK 基本模式
err := ch.Qos(1, 0, false) // 预取 1 条,公平分发

msgs, err := ch.Consume("task.queue", "", false, // autoAck=false
    false, false, false, nil)

for d := range msgs {
    result, err := processJob(d.Body)
    if err != nil {
        // 业务异常:拒绝并路由到死信
        ch.Nack(d.DeliveryTag, false, false) // requeue=false
        continue
    }
    ch.Ack(d.DeliveryTag, false) // 处理完成
}

unacked 消息的重放机制

当消费者未发送 ACK 时,消息处于 unacked 状态,存储在 Broker 内存中。此时有三种处置方式:

  1. basic.ack:永久标记为已消费,从队列中移除
  2. basic.nack(requeue=true):放回队列头部,可能被同一消费者再次拉取
  3. basic.nack(requeue=false):根据队列的 DLX 配置路由到死信队列

Tip

面试高频追问:为什么手动 ACK 会比 Auto ACK 更安全但又可能更慢?答案在于 trade-off——手动 ACK 把控制权交给了消费者,你可以在数据库事务提交后才 ACK,这样即使消费者崩溃,未提交事务对应的消息也不会丢。代价是多了一次 RPC 调用(ack 请求到 Broker)。

Dead Letter Exchange(DLX)

DLX 不是独立的消息队列,而是一种通过队列属性触发的路由行为。当队列配置了 x-dead-letter-exchange,满足以下条件的消息会被自动重新发布到指定的 DLX:

  • 被 basic.nack 且 requeue=false
  • 消息 TTL 到期(x-message-ttl 或单消息 TTL)
  • 队列达到最大长度(x-max-length)

DLX 本身是一个普通的 Exchange,可以配置自己的 binding key(x-dead-letter-routing-key)来控制死信消息的去向。

sequenceDiagram
    participant P as 生产者
    participant REX as retry-exchange(Direct)
    participant RQ as retry-queue<br/>(TTL+DLX)
    participant DX as dead-letter-exchange
    participant DQ as dead-letter-queue
    participant C as 消费者

    P->>REX: publish(order.id=123)
    REX->>RQ: route to retry-queue<br/>(TTL=30s)
    RQ->>RQ: 消息等待 TTL 过期
    Note right of RQ: unacked / 被 nack(requeue=false) / TTL到期
    RQ->>DX: DLX 捕获消息
    DX->>DQ: 投递到死信队列
    DQ->>C: consumer 消费死信

TTL + DLX 实现重试队列模式

重试是 DLX 最经典的应用场景。通过多层队列接力,实现"尝试 N 次 → 最终死信"的自动降级流程。

# 重试队列层级设计
retry-layer-1:
  exchange: retry-exchange
  queue: retry-queue-1
  arguments:
    x-message-ttl: 5000          # 5 秒后重试
    x-dead-letter-exchange: retry-exchange
    x-dead-letter-routing-key: retry-layer-2  # 第一次重试仍进 retry-exchange

retry-layer-2:
  exchange: retry-exchange
  queue: retry-queue-2
  arguments:
    x-message-ttl: 30000         # 30 秒后重试

retry-layer-3:
  exchange: retry-exchange
  queue: retry-queue-3
  arguments:
    x-message-ttl: 120000        # 2 分钟后重试

dead-letter-queue:
  exchange: dead-letter-exchange
  queue: dead-letter-queue
  arguments:
    x-dead-letter-exchange: monitoring-exchange  # 最终告警

这种 N 层队列模式的优势在于每条消息只需要携带一个简单的 TTL header,而不需要维护计数器。当队列配置 x-dead-letter-routing-key 指向下一层重试队列所在的 Exchange,消息自然形成递增的时间链。

Note

如果不想用多层队列,也可以在消息 body 中嵌入一个 retryCount 字段,消费者处理失败时增加计数并重新发送到原 Exchange。这种方式更灵活,但需要应用层维护计数逻辑。

Nack + Requeue 策略

Nack 是 Negative Acknowledgement 的缩写,用于告知 Broker 某条消息不应被继续交给当前消费者。关键的决策树如下:

  • nack(tag, multiple=false, requeue=true):单条重放,回到队列头。适用于临时故障(如下游短暂不可用)。
  • nack(tag, multiple=false, requeue=false):单条不重放,走 DLX。适用于不可恢复的错误(如数据校验失败)。
  • nack(tag, multiple=true, ...):一次性拒绝当前通道中所有 unacked 消息。谨慎使用,可能造成大面积消息堆积。
// Go 代码示例:带指数退避的重试 Nack
const maxRetries = 3

func handleMessage(d amqp.Delivery) {
    retries := d.Headers["x-retry-count"].(int32)
    if retries >= maxRetries {
        ch.Nack(d.DeliveryTag, false, false) // 进入死信
        return
    }
    d.Headers["x-retry-count"] = retries + 1

    if err := process(d.Body); err != nil {
        backoff := math.Pow(2, float64(retries)) * 1000
        d.Headers["x-delay"] = uint64(backoff)
        ch.Publish("retry-exchange", "retry", false, false,
            amqp.Publishing{
                Body:  d.Body,
                Headers: d.Headers,
            })
        ch.Ack(d.DeliveryTag, false)
    } else {
        ch.Ack(d.DeliveryTag, false)
    }
}

这段代码的关键在于:当处理失败时,先在消息头中递增重试次数,计算指数退避延迟,然后重新发布到重试 Exchange,最后对原消息发送 ACK。这样既避免了 nack 导致的重复拉取,又能让下次投递有时间间隔。

实践场景

电商订单支付回调处理:第三方支付平台发起异步回调时,可能出现网络抖动或服务短暂的负载过高。采用 Manual ACK + 三层重试队列 + DLX 的方案:

  1. 第一层重试间隔 5 秒,适合处理瞬时网络问题
  2. 第二层重试间隔 30 秒,给下游留出恢复窗口
  3. 第三层重试间隔 2 分钟,覆盖短时服务重启
  4. 三次重试均失败后进入死信队列,人工介入或对账修复

这个模式下,prefetch_count 通常设为与消费者并发度相当的值,配合 no_ack=false 实现精准控制。

Tip

实践中常见的坑:如果消费者处理耗时过长(比如几分钟),unacked 的消息会持续占用 prefetch 配额,导致其他消费者无法获得新消息。解决方案是将 prefetch 适当放大、或者在处理过程中定期发送 heartbeat ack。

关联笔记