一句话结论
RabbitMQ 的核心是 Exchange → Binding → Queue 的路由模型。消息可靠性靠三阶段保障(Confirm + 持久化 + ACK),死信队列处理失败消息,延迟消息靠 TTL+DLX 实现。
Exchange 类型
消息可靠性三阶段
生产 → Broker: Publisher Confirm(确认到达)
Broker 存储: 队列+消息持久化(durable=true)
Broker → 消费: 手动 ACK(处理完才确认)
只能保证至少一次投递,消费端必须幂等!
详细拆解
阶段 1:Producer → Broker(Publisher Confirm)
// 开启 Confirm 模式
ch.Confirm(false) // false=不等待
confirms := ch.NotifyPublish(make(chan amqp.Confirmation, 1))
err := ch.Publish("orders", "order.created", true, false, msg)
if err != nil {
// 发布失败,需要重试或补偿
}
select {
case conf := <-confirms:
if !conf.Ack {
// Broker 拒绝(可能队列不存在等)
}
case <-time.After(5 * time.Second):
// 超时未确认 → 需要重试(但要幂等!)
}
阶段 2:Broker 存储(持久化)
durable=true:队列在 Broker 重启后不丢失persistent=true(deliveryMode=2):消息写入磁盘⚠️ 持久化 ≠ 每条消息都 fsync——RabbitMQ 有刷盘间隔
阶段 3:Broker → Consumer(ACK)
// ❌ autoAck=true → 消息投递后立即删除,消费者崩溃则消息丢失
// ✅ autoAck=false → 手动确认
msgs, _ := ch.Consume("orders", "", false, false, false, false, nil)
for msg := range msgs {
err := processOrder(msg.Body)
if err != nil {
msg.Nack(false, true) // 重回队列重试
continue
}
msg.Ack(false) // false=只确认当前这条
}
可靠性三角: Confirm + 持久化 + 手动 ACK → 至少一次交付。如果任何环节失败+重试 → 可能重复交付 → 消费端必须幂等。
死信队列(DLX)
消息被 NACK/reject + requeue=false → 进入 DLX
消息 TTL 过期 → 进入 DLX
队列满了 → 进入 DLX
用途: 失败消息统一处理、延迟队列
延迟消息
消息 TTL + 死信队列 = 延迟消息
1. 消息发到 TTL 队列(无消费者)
2. TTL 过期 → 自动路由到 DLX
3. 实际消费者监听 DLX 队列
场景: 订单 30 分钟未支付自动关闭
消息积压处理
紧急扩容消费者实例(水平扩展)
新开队列 + 新消费者(不影响旧队列)
启用批量消费(prefetch_count 调大)
丢弃低优先级消息(降级策略)
临时增加 Topic 分区(仅 Kafka)
Prefetch 与流量控制
// Qos 控制消费者预取数量(未 ACK 的消息数上限)
ch.Qos(
100, // prefetchCount: 最多同时处理 100 条
0, // prefetchSize: 0=不限大小
false, // global: false=当前 channel 独占
)
// prefetch=1: 每次只取一条,处理完再取
// prefetch=100: 允许批量取,提高吞吐但要注意公平性
关键: prefetch 太小 → 消费者空闲等待(吞吐低)。prefetch 太大 → 消息堆积在消费者内存(崩溃则丢大量未确认消息)。一般设为消费者处理能力的 1-2 秒量。
连接与通道模型
AMQP 连接层次:
物理 TCP 连接
└── Channel 1(轻量虚拟连接,有自己的 acks/confirms)
└── Channel 2(可并发,不需建立新 TCP 连接)
└── Channel N(一个 TCP 连接可承载上千个 Channel)
为什么要复用 Channel?
建立 TCP 连接(TCP 握手 + TLS)代价高 → 一次连接,多 Channel 并行
Channel 是逻辑隔离,close Channel 不影响其他 Channel
高频面试问题
Q: RabbitMQ 如何保证消息不丢?
30 秒回答: 三阶段保障。Producer 端用 Publisher Confirm 确认 Broker 收到。Broker 端队列和消息都设为 durable+persistent 持久化到磁盘。Consumer 端关闭 autoAck,处理完手动 ACK。三者缺一不可。
深入回答: 即使三层全做,也只是"至少一次交付"。因为 Confirm 超时或 ACK 超时会导致重试——消息可能被投递两次。所以消费端要做幂等(消息 ID 去重 / 业务幂等)。真正的"精确一次"在 RabbitMQ 中做不到,需要用事务性发件箱或消费端幂等来近似。
继续追问:
"持久化后一定不丢吗?" → 不保证。RabbitMQ 不是每条消息都 fsync——刷盘有间隔,如果 Broker 在两次刷盘之间宕机,未刷盘的消息会丢失。这就是为什么还要用镜像队列(Mirrored Queue)或 Quorum Queue 做多节点冗余。
"Confirm + 事务可以同时开吗?" → 不建议。事务(txSelect/txCommit)会严重降低吞吐(同步等回包)。Publisher Confirm 是异步的,性能好得多。RabbitMQ 官方推荐用 Confirm 替代事务。
Q: 死信队列的典型应用场景?
30 秒回答: ① 延迟队列(消息 TTL 过期 → DLX → 实际消费者,如订单 30 分钟未支付自动取消)。② 失败重试耗尽后的死信(重试 3 次仍失败 → 进入 DLX 供人工处理)。③ 队列溢出时的溢出处理(队列满了 → 溢出到 DLX 防止消息丢失)。
继续追问:
"延迟队列有什么精度问题?" → RabbitMQ TTL 的检查不是实时的——消息过期后未必立即被投递到 DLX,可能有秒级延迟。对精度要求高的场景(如秒级定时任务),RabbitMQ 的延迟消息不够精确,考虑用 Redis 的 ZSet(Score=过期时间戳)或者专门的延迟队列(如 SchedulerX)。
Q: RabbitMQ vs Kafka 怎么选?
30 秒回答: RabbitMQ 适合业务任务(订单处理、邮件发送),消息量万级,需要灵活路由和 ACK。Kafka 适合大数据流(日志采集、事件溯源),消息量百万级,需要消息回溯和顺序保证。
继续追问:
"能用 RabbitMQ 做日志采集吗?" → 可以但不合适。消息量上去后 RabbitMQ 吞吐不如 Kafka,且 RabbitMQ 不支持消息回溯(消息被 ACK 后就删了——做日志的话历史数据全丢了)。Kafka 按时间/offset 随意回溯。
"能用 Kafka 做秒杀异步下单吗?" → 可以,但要注意 Kafka 的 Pull 模型——消费者按 offset 拉取,不支持单个消息的 ACK/NACK。如果某条消息处理失败需要跳过的逻辑会比 RabbitMQ 复杂。
项目应用
速记
Exchange=Direct/Topic/Fanout。可靠性=Confirm+持久化+ACK。死信=拒绝/过期/满→DLX处理。延迟=TTL+DLX。积压=扩容+新队列+批量。至少一次交付→消费端必须幂等。