190 lines
8.2 KiB
Markdown
190 lines
8.2 KiB
Markdown
|
|
---
|
|||
|
|
tags: [mq/rabbitmq, publisher-confirm, message-persistence, rabbitmq-transaction, return-callback]
|
|||
|
|
create time: 2026-08-08 18:51
|
|||
|
|
update time: 2026-08-08 18:51
|
|||
|
|
---
|
|||
|
|
|
|||
|
|
# 消息持久化与可靠性投递
|
|||
|
|
|
|||
|
|
## 概述
|
|||
|
|
|
|||
|
|
在分布式系统中,消息丢失是比重复消费更严重的问题。RabbitMQ 的可靠性投递涉及三个层面的协同:Exchange 和 Queue 的持久化配置、消息本身的 persistent 标记、以及 Publisher Confirm 机制提供的发送端反馈。理解这些机制的设计动机和取舍,才能在实际业务中做出正确的选型。
|
|||
|
|
|
|||
|
|
## 核心原理
|
|||
|
|
|
|||
|
|
### 三阶段持久化
|
|||
|
|
|
|||
|
|
消息从生产者发出到被消费者消费的整个链路中,任何一个环节未持久化都会导致消息丢失风险。
|
|||
|
|
|
|||
|
|
```mermaid
|
|||
|
|
graph TD
|
|||
|
|
A[生产者] -->|"1. Exchange durable"| B["Direct Exchange<br/>(durable=true)"]
|
|||
|
|
B -->|"2. Queue durable"| C["Order Queue<br/>(durable=true)"]
|
|||
|
|
C -->|"3. Message persistent"| D["磁盘存储"]
|
|||
|
|
|
|||
|
|
style A fill:#e1f5fe
|
|||
|
|
style B fill:#fff3e0
|
|||
|
|
style C fill:#e8f5e9
|
|||
|
|
style D fill:#fce4ec
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
第一阶段:Exchange 声明时设置 `durable=true`。如果 Exchange 不存在,RabbitMQ 会在启动时自动重建;若未设置 durable,重启后 Exchange 会被删除。
|
|||
|
|
|
|||
|
|
第二阶段:Queue 声明时设置 `durable=true`。注意:**队列持久化只保证元数据不丢**。实际存储在 queue 中的消息默认为内存页缓存,除非消息带有 persistent 标记才会写入磁盘。
|
|||
|
|
|
|||
|
|
第三阶段:消息发送时设置 `DeliveryMode=2`(persistent)。持久化消息在抵达 Queue 后会刷盘,而非仅仅留在内存 buffer 中。
|
|||
|
|
|
|||
|
|
> [!WARNING]
|
|||
|
|
> 常见误解:仅设置 Queue durable 并不能保证消息不丢。如果消息 DeliveryMode=1(transient),即使队列本身是持久的,这些瞬态消息在服务器重启时也会被丢弃。
|
|||
|
|
|
|||
|
|
### Publisher Confirm 机制
|
|||
|
|
|
|||
|
|
Publisher Confirm 是 RabbitMQ 提供的一种异步确认机制。生产者在开启 confirm 模式后,每条消息都会被分配一个唯一的 `delivery-tag`,Broker 处理完成后会通过回调通知生产者结果。
|
|||
|
|
|
|||
|
|
两种工作模式:
|
|||
|
|
|
|||
|
|
| 模式 | 说明 | 适用场景 |
|
|||
|
|
|------|------|---------|
|
|||
|
|
| 单条确认 | 每发一条消息等一次 ack/nack,串行等待 | 吞吐量要求极低、需严格逐条确认的场景 |
|
|||
|
|
| 批量确认 | 一次发送多条,一次性收到确认回调 | **绝大多数场景的首选**,吞吐量接近未开启 confirm 的水平 |
|
|||
|
|
|
|||
|
|
Confirm 模式的回调有两种触发条件:
|
|||
|
|
|
|||
|
|
```mermaid
|
|||
|
|
sequenceDiagram
|
|||
|
|
participant P as 生产者
|
|||
|
|
participant RMQ as RabbitMQ Broker
|
|||
|
|
participant CC as ConfirmCallback
|
|||
|
|
participant RC as ReturnCallback
|
|||
|
|
|
|||
|
|
P->>RMQ: Publish(msg, mandatory=true)
|
|||
|
|
activate RMQ
|
|||
|
|
RMQ-->>CC: ack(tag=1001)
|
|||
|
|
activate CC
|
|||
|
|
Note right of CC: 消息路由到 Queue 成功
|
|||
|
|
CC-->>P: 回调确认
|
|||
|
|
deactivate CC
|
|||
|
|
deactivate RMQ
|
|||
|
|
|
|||
|
|
Note over P,RMQ: 另一条消息
|
|||
|
|
P->>RMQ: Publish(msg2, mandatory=true)
|
|||
|
|
activate RMQ
|
|||
|
|
alt Exchange存在但无绑定队列
|
|||
|
|
RMQ-->>RC: return(msg2)
|
|||
|
|
activate RC
|
|||
|
|
Note right of RC: 路由失败但非不可恢复
|
|||
|
|
RC-->>P: 回调退回
|
|||
|
|
deactivate RC
|
|||
|
|
else Exchange不存在
|
|||
|
|
RMQ-->>RC: return(msg2)
|
|||
|
|
activate RC
|
|||
|
|
Note right of RC: Exchange 不可达
|
|||
|
|
RC-->>P: 回调退回
|
|||
|
|
deactivate RC
|
|||
|
|
end
|
|||
|
|
deactivate RMQ
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
`ConfirmCallback` 告诉生产者消息是否被 Broker 安全接收(路由完成),`ReturnCallback` 处理强制性消息因无法路由而被退回的情况。只有同时开启 `mandatory=true`,ReturnCallback 才会在路由失败时被触发。
|
|||
|
|
|
|||
|
|
### RabbitMQ Transaction 模式
|
|||
|
|
|
|||
|
|
RabbitMQ 的 AMQP 协议原生支持事务。开启事务后,发送的消息只有在事务提交后才会真正进入队列。如果事务回滚,期间发送的消息不会进入任何队列。
|
|||
|
|
|
|||
|
|
事务消息的四项基本原则:
|
|||
|
|
|
|||
|
|
1. **Channel 必须声明为事务模式**:调用 `channel.txSelect()`
|
|||
|
|
2. **消息发送后必须提交事务**:`channel.txCommit()`
|
|||
|
|
3. **事务内发送的消息不会立即投递**:直到 txCommit 后才进入队列
|
|||
|
|
4. **txCommit 失败则事务回滚**:可通过 `txRollback()` 手动回滚
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
// Go 代码示例:Transaction 模式(慎用)
|
|||
|
|
func sendWithTx(ch *amqp.Channel, body string) error {
|
|||
|
|
if err := ch.Tx(); err != nil {
|
|||
|
|
return fmt.Errorf("tx select failed: %w", err)
|
|||
|
|
}
|
|||
|
|
defer func() {
|
|||
|
|
_ = ch.TxRollback() // 出错时回滚
|
|||
|
|
}()
|
|||
|
|
if err := ch.Publish("", "my.queue", false, false,
|
|||
|
|
amqp.Publishing{Body: []byte(body)}); err != nil {
|
|||
|
|
return err
|
|||
|
|
}
|
|||
|
|
return ch.TxCommit()
|
|||
|
|
}
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
### Confirm vs Transaction:对比分析
|
|||
|
|
|
|||
|
|
| 维度 | Publisher Confirm | Transaction |
|
|||
|
|
|------|------------------|-------------|
|
|||
|
|
| 吞吐量 | 高(异步批量) | 低(同步阻塞) |
|
|||
|
|
| 实现复杂度 | 简单 | 简单 |
|
|||
|
|
| 故障恢复 | 需要重试逻辑 | Broker 自动回滚 |
|
|||
|
|
| 资源占用 | 低 | 高(需维护事务日志) |
|
|||
|
|
| 语义强度 | 最多一次(至少一次需重试) | 恰好一次(配合幂等设计) |
|
|||
|
|
| 适用性 | **生产环境首选** | 几乎不推荐使用 |
|
|||
|
|
|
|||
|
|
为什么 Confirm 优于 Transaction?
|
|||
|
|
|
|||
|
|
1. **性能差距巨大**:Transaction 采用同步提交方式,每一条消息的事务提交都要等待 Broker 的 fsync 操作完成,而 Confirm 可以批量异步返回。实测差距可达数倍甚至十倍以上。
|
|||
|
|
2. **资源开销更小**:Transaction 需要维护完整的事务日志和回滚状态,在大量并发下容易成为瓶颈。Confirm 模式无需这些额外开销。
|
|||
|
|
3. **生态兼容性更好**:Spring AMQP、Go amqp 等主流客户端对 Confirm 的支持远成熟于 Transaction。
|
|||
|
|
|
|||
|
|
> [!TIP]
|
|||
|
|
> 面试回答要点:如果需要"恰好一次"语义,不要指望 RabbitMQ 本身保证——无论是 Confirm 还是 Transaction 都不够。正确做法是:Confirm 作为投递保障 + 业务层幂等键(如 orderId 唯一索引)兜底。这才是工业级方案的思路。
|
|||
|
|
|
|||
|
|
## 代码示例
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
// Go 代码示例:Confirm + Return 完整回调链路
|
|||
|
|
ch.NotifyPublish(make(chan amqp.Confirmation)) // 注册 Confirm
|
|||
|
|
ch.NotifyReturn(make(chan amqp.Return)) // 注册 Return
|
|||
|
|
|
|||
|
|
confirms := ch.NotifyPublish(nil)
|
|||
|
|
for conf := range confirms {
|
|||
|
|
if conf.Ack {
|
|||
|
|
log.Printf("消息 delivery-tag=%d 已确认", conf.DeliveryTag)
|
|||
|
|
} else {
|
|||
|
|
// nack 时需要业务侧重试或写入死信表
|
|||
|
|
log.Printf("消息 delivery-tag=%d 未被确认", conf.DeliveryTag)
|
|||
|
|
}
|
|||
|
|
}
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
这段代码展示了 Confirm 的核心结构:通过 `NotifyPublish` 注册一个 Confirmation 通道,每个发出的消息对应一个 Confirmation 对象,其 `Ack` 字段为 true 表示 Broker 已成功处理,false 则需要应用层处理重发或降级逻辑。
|
|||
|
|
|
|||
|
|
```java
|
|||
|
|
// Java 代码示例:Spring AMQP 中的 Confirm 配置
|
|||
|
|
@Bean
|
|||
|
|
public RabbitTemplate rabbitTemplate(RabbitConnectionFactory factory) {
|
|||
|
|
RabbitTemplate template = new RabbitTemplate(factory);
|
|||
|
|
template.setMandatory(true); // 开启 ReturnCallback
|
|||
|
|
template.setConfirmCallback((correlation, ack, reason) -> {
|
|||
|
|
if (!ack) log.warn("消息未确认: {}", reason);
|
|||
|
|
});
|
|||
|
|
template.setReturnsCallback(returned -> {
|
|||
|
|
log.warn("消息被退回: {} - {}", returned.getMessage(), returned.getReason());
|
|||
|
|
});
|
|||
|
|
return template;
|
|||
|
|
}
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
## 实践场景
|
|||
|
|
|
|||
|
|
**订单支付后的状态推送**:用户支付完成后需要将状态变更推送给物流、库存等多个下游系统。在这个场景下:
|
|||
|
|
|
|||
|
|
- Exchange 和 Queue 均设置为 `durable=true` 防止服务重启丢数据
|
|||
|
|
- 消息设置为 `DeliveryMode=2`(persistent)确保磁盘持久化
|
|||
|
|
- 开启 Publisher Confirm,nack 的消息放入本地重试表,定时任务补偿
|
|||
|
|
- 配合业务幂等键(支付流水号 + 目标系统标识),确保重复收到的消息不会产生副作用
|
|||
|
|
|
|||
|
|
> [!WARNING]
|
|||
|
|
> 持久化消息会带来明显的性能代价(大约 10 倍吞吐下降)。在对丢消息敏感度低的场景(如埋点日志采集),可以使用 transient 消息获得更高吞吐。关键是评估业务的容错阈值后再决定。
|
|||
|
|
|
|||
|
|
## 关联笔记
|
|||
|
|
- [[Exchange 路由机制]]
|
|||
|
|
- [[ACK 确认与死信队列]]
|
|||
|
|
- [[推拉结合消费模式]]
|