一句话结论

RabbitMQ 的核心是 Exchange → Binding → Queue 的路由模型。消息可靠性靠三阶段保障(Confirm + 持久化 + ACK),死信队列处理失败消息,延迟消息靠 TTL+DLX 实现。

Exchange 类型

类型

路由方式

场景

Direct

精准匹配 routing_key

任务分发

Topic

通配符匹配(order.*, #.error)

灵活路由

Fanout

广播到所有绑定队列

配置更新通知

Headers

按消息头匹配

复杂条件

消息可靠性三阶段

生产 → 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 分钟未支付自动关闭

消息积压处理

  1. 紧急扩容消费者实例(水平扩展)

  2. 新开队列 + 新消费者(不影响旧队列)

  3. 启用批量消费(prefetch_count 调大)

  4. 丢弃低优先级消息(降级策略)

  5. 临时增加 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 复杂。

项目应用

项目

RabbitMQ 场景

AI-Agent工作流与评测平台

Agent 长任务异步执行、Worker 消费

分布式电商交易系统

秒杀异步下单、订单超时关闭(延迟队列)

速记

Exchange=Direct/Topic/Fanout。可靠性=Confirm+持久化+ACK。死信=拒绝/过期/满→DLX处理。延迟=TTL+DLX。积压=扩容+新队列+批量。至少一次交付→消费端必须幂等。