180 lines
7.6 KiB
Markdown
180 lines
7.6 KiB
Markdown
|
|
---
|
|||
|
|
tags: [mq/rabbitmq, manual-ack, dead-letter-exchange, dlq, retry-pattern, prefetch]
|
|||
|
|
create time: 2026-08-08 18:52
|
|||
|
|
update time: 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
|
|||
|
|
// 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`)来控制死信消息的去向。
|
|||
|
|
|
|||
|
|
```mermaid
|
|||
|
|
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 次 → 最终死信"的自动降级流程。
|
|||
|
|
|
|||
|
|
```yaml
|
|||
|
|
# 重试队列层级设计
|
|||
|
|
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
|
|||
|
|
// 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。
|
|||
|
|
|
|||
|
|
## 关联笔记
|
|||
|
|
- [[Exchange 路由机制]]
|
|||
|
|
- [[消息持久化与可靠性投递]]
|
|||
|
|
- [[推拉结合消费模式]]
|