跳转到主要内容

数据库/协同工具类型

RabbitMQ 深度指南

RabbitMQ 的作用定位、何时该用、核心机制(Exchange/Queue/持久化/ACK)、常见问题排查(消息丢失、重复消费、堆积、断连)、难点解决方案(可靠性三连、顺序消费、死信队列、幂等)与生产要点。

  • RabbitMQ
  • 消息队列

这是 RabbitMQ 的深度参考文档。上手实操见 Nest 课程第二十章

一、基础入门(5 分钟上手)

这一节带你把 RabbitMQ 跑起来、打开管理界面、亲手发一条收一条消息;Nest 课程那篇讲的是在应用里封装生产者/消费者——两篇不重复。

1. 用 Docker 起一个 RabbitMQ

docker run -d \
  --name rabbitmq \
  -p 5672:5672 -p 15672:15672 \
  rabbitmq:3-management
  • 5672:AMQP 协议端口(代码连这个)。
  • 15672:管理界面端口。
  • 3-management 镜像才自带 Web 控制台。
  • 浏览器打开 http://localhost:15672 ,默认账号 guest / guest

2. 在管理界面认识它

控制台顶部几个 Tab:

  • Overview:整体状态。
  • Connections / Channels:连进来的客户端。
  • Queues:队列,新手重点看这个。
  • Exchanges:交换机,消息先进这里再分发到队列。

3. 发一条 / 收一条消息

左侧 Queues → Add a new queue,Name 填 hello,点 Add queue

  1. hello 队列页,此时 Ready = 0(队列空)。
  2. 往下翻到 Publish message,Payload 填 {"msg":"hi"},点 Publish message。 → 页面上方提示 Message published.,队列的 Ready 变成 1——有一条消息在等消费。
  3. 再翻到 Get Messages,Ack mode 选 Basic Ack,点 Get Message(s)。 → 下方展开那条 {"msg":"hi"}Ready 归 0——消息被取走并确认了。

这就是 MQ 的本质:生产者丢进去(Ready +1)、消费者取出来(Ready -1)

4. 核心组件与 Exchange 四类型

回顾组件:Producer(生产者)→ Exchange(交换机)→ Queue(队列)→ Consumer(消费者),中间靠 binding + routing key 把交换机和队列连起来。

交换机怎么把消息分发到队列,由它的类型决定(掘金《四种类型交换机》):

类型routing key 匹配规则典型场景
direct精确相等点对点(按日志级别 error/warn)
topic通配符:* 匹配一个词、# 匹配零或多个词灵活路由(order.*.created
fanout忽略 routing key,广播到所有绑定队列群发(配置更新通知所有节点)
headers按消息头匹配(不看 routing key)少用

四种 Exchange 类型的路由规则对比:

flowchart LR P(["Producer 生产者"]) E(("Exchange 交换机")) Q1["队列1 direct"] Q2["队列2 fanout"] Q3["队列3 topic"] Q4["队列4 headers"] P -->|"publish"| E E -->|"direct 精确匹配 key"| Q1 E -->|"fanout 广播所有绑定"| Q2 E -->|"topic 通配符 order.*"| Q3 E -->|"headers 按消息头"| Q4

举例(topic):队列绑 order.*.paid,那么 order.shanghai.paid 能进,order.paid 进不了。

5. 生产 / 消费流程

概念上四步(代码层在 Nest 课程):

  1. 声明 Exchange(指定类型 + durable: true 持久化)。
  2. 声明 Queuedurable: true)。
  3. Binding:把 Queue 绑到 Exchange,带上 routing key。
  4. 生产者 publish(exchange, routingKey, msg);消费者 consume(queue) 拿到消息,处理完手动 ack

消息流向:生产者只管发到 Exchange + routing key,由 Exchange 按类型决定进哪些队列,生产者不关心具体队列。

6. 基本使用规则

  • 队列和消息都要持久化durable: true + 消息 persistent: true),MQ 重启才不丢。
  • 消费端手动 ack:处理完再确认;自动确认会在拿到消息就删,处理失败会丢。
  • 消费端做幂等:MQ 是「至少一次」投递,可能重复,按 messageId 去重或业务幂等。
  • 命名规范:Exchange order.created,Queue order.created.queue,routing key 跟着业务事件。
  • prefetch 控制流量:限制消费者未确认的消息数,避免一次拉太多撑爆。
  • 别拿 MQ 当数据库:消息有生命周期,堆积的转 DB 或清理。

7. 核心词汇速记

术语一句话
Producer生产者,发消息的
Consumer消费者,收消息的
Queue队列,消息存这儿等被消费
Exchange交换机,按规则把消息分发到队列
binding / routing key交换机和队列的绑定规则
ACK消费者确认「处理完了」,MQ 才删消息

二、作用与定位

RabbitMQ 是消息队列(MQ),核心价值三个字:异步、解耦、削峰

  • 异步:发消息方不用等处理完(注册后异步发邮件、下单后异步扣库存)。
  • 解耦:生产者消费者互不依赖(一个挂了恢复继续消费)。
  • 削峰:高并发写库前先入队,消费端按 DB 承受速率慢慢消费,保护 DB。

一句话:有「先接住、慢慢处理」或「A 触发但不想等 B」的需求,用 MQ

三、何时该用 / 何时用别的

场景选择
任务异步、削峰、复杂路由(按规则分发)RabbitMQ
日志流、超高吞吐、事件溯源Kafka
简单队列、量不大、已有 RedisRedis 的 list/stream
实时双向通信WebSocket(不是 MQ)

RabbitMQ vs Kafka:RabbitMQ 路由灵活(Exchange 4 种)、消息确认完善、单条延迟低,适合业务消息;Kafka 吞吐极高、按日志持久化、适合大数据流。Nest 业务系统大多用 RabbitMQ。

四、核心机制速览

Producer → [ Exchange → binding(routing key) → Queue ] → Consumer

一张图看清消息从生产到消费的完整路由路径:

flowchart LR P(["Producer 生产者"]) E(("Exchange 交换机")) Q["Queue 队列"] C(["Consumer 消费者"]) P -->|"publish 加 routing key"| E E -->|"按 binding 规则路由"| Q Q -->|"取出消息消费"| C C -->|"ack 确认"| Q
  • Exchange 4 类型direct(key 精确匹配)、topic(key 通配符)、fanout(广播)、headers(按头匹配)。
  • Queue + binding:队列绑到 Exchange + routing key,消息按规则进队列。
  • 持久化:队列声明 durable: true、消息 persistent: true,MQ 重启不丢。
  • ACK:消费者手动确认(ack)后 MQ 才删消息,处理失败可 nack(重入队或丢弃)。
  • prefetch:限制消费者未确认的消息数,做流量控制。

五、常见问题与排查

1. 消息丢失(最致命)

三个环节都可能丢:

环节丢的原因解决
生产者 → MQ网络断、消息没真到publisher confirm(MQ 收到后回确认)
MQ 自身重启、宕机队列 + 消息都持久化durable + persistent
MQ → 消费者消费者拿到没处理完就挂关闭自动确认(noAck: false),处理完再 ack

可靠性三连:confirm + 持久化 + 手动 ACK,三个都开才不丢。

2. 重复消费

MQ 保证「至少一次」投递,不保证「恰好一次」——消费者可能收到重复消息(比如 ACK 丢了,MQ 重发)。

解决:消费端做幂等

  • 消息带唯一 id(messageId),消费前查「处理过没」(Redis/DB 去重)。
  • 或业务天然幂等(如 upsertDECR 后判断)。

3. 消息堆积

原因:消费速度跟不上生产,或消费者挂了消息堆队列里。

解决

  • 加消费者(水平扩容)。
  • 提高 prefetch、优化消费逻辑(批量处理)。
  • 紧急情况临时起一个「清扫消费者」快速消费(但要保证幂等)。
  • TTL + 死信队列:超时的消息转死信,别无限堆。

4. 连接频繁断开重连

原因:连接没复用、心跳超时、网络抖动。

解决:用连接池 / 长连接、开心跳(heartbeat)、客户端配自动重连(如 amqplib 的重连逻辑)。

5. 消费者处理慢导致 prefetch 不生效

prefetch 设了但一次还是收到很多 → 检查是不是用了同一个 channel 多次 consume,或 prefetch 作用域设错(per-channel vs per-consumer)。

六、难点与解决方案

1. 顺序消费

场景:同一笔订单的「创建 → 支付 → 发货」消息必须按顺序处理。

难点:多消费者并行消费时,顺序会乱。

方案

  • 需要保序的消息路由到同一队列(用 routing key,如 order.${orderId} 落同一 queue),且该队列只一个消费者串行处理。
  • Kafka 天然按 partition 保序(同 key 同 partition),这是它比 RabbitMQ 适合顺序场景的原因。

2. 死信队列(DLX)

消息「死了」(被 reject/nack 且不重入队、或 TTL 过期、或队列满)时,转到一个死信交换机/队列,便于后续排查或重试。

正常队列 --(失败/过期)--> 死信交换机 --> 死信队列 --> 人工/定时处理

死信队列的完整流转路径:

flowchart LR P(["Producer 生产者"]) NQ["正常队列"] C(["Consumer 消费者"]) DLX(("死信交换机 DLX")) DQ["死信队列"] H(["人工或告警兜底"]) P -->|"发消息"| NQ NQ -->|"消费"| C C -->|"nack 或 reject 失败"| NQ NQ -->|"超 TTL 或重试达上限"| DLX DLX --> DQ DQ -->|"兜底处理"| H

典型用法:消费失败 → 进死信 → 定时重试 N 次 → 还失败就告警人工介入。

3. 延迟队列

RabbitMQ 本身没原生延迟队列,靠 TTL + 死信实现:消息进一个带 TTL 的队列,过期后转死信队列被消费 = 延迟。或用 rabbitmq_delayed_message_exchange 插件。

4. 幂等设计

所有消费者都要按「可能重复」来设计。模式:

  • 唯一 id 去重if (redis.get(msgId)) return; process(); redis.set(msgId, 1)
  • 业务幂等:状态机(「未支付 → 已支付」只允许走一次)、INSERT ... ON DUPLICATE KEY、乐观锁版本号。

七、生产环境要点

  • 可靠性三连:publisher confirm + 持久化 + 手动 ACK。
  • 消费端幂等,别假设「只收到一次」。
  • 设监控告警:队列堆积深度、消费者数量、未确认消息数、连接数。
  • 命名规范:Exchange / Queue / routing key 命名清晰(如 order.createdorder.created.queue)。
  • 别让 MQ 当数据库:消息是有生命周期的,堆积的该转 DB 或清理。
  • 集群 + 镜像队列:RabbitMQ 高可用用镜像队列(或 Quorum Queue),单节点别上生产。

速查:常见问题对照

问题原因解决
消息丢了三连没开齐confirm + 持久化 + 手动 ACK
重复消费MQ 至少一次投递消费端幂等(唯一 id / 业务幂等)
队列堆积消费慢 / 消费者挂加消费者、优化、TTL+死信
顺序错乱并行消费同 key 同队列 + 单消费者
处理失败业务异常nack + 死信队列重试
延迟需求没原生延迟TTL + 死信 / 延迟插件