Chapter 01
模型与路由:消息不是发给队列的
RabbitMQ 的一切都立在一个反直觉的事实上:生产者从不把消息发给队列。这一章拆开消息从 producer 到 consumer 的完整路径——以及它在 ack 之后消失的那一刻。
本章你将建立的 schema
- Connection 与 Channel 的多路复用:一条 TCP,多个 channel,靠帧里的编号区分
- exchange / queue / binding 三件套:消息发给 exchange,由 binding 路由到 queue
- 四种 exchange 的路由规则,以及 routing key 与 binding key 的区别
- 消息生命周期:ack 之后即删除——RabbitMQ 没有 offset,这是它与 Kafka 的根本分叉
1.1Connection 与 Channel
Channel 是复用在一条 TCP Connection 上的轻量虚拟连接;所有协议交互跑在 channel 上,而不是直接跑在 TCP 上。
一个高并发服务可能有几百个线程同时收发消息。每个线程开一条 TCP 连接,意味着几百次握手、几百个文件描述符、几百份内核缓冲。channel 把这些复用到一条(或少数几条)TCP 连接上:每个 channel 一个编号,broker 靠编号把会话区分开。
底层机制(比文档深一层):AMQP 把每一帧(frame)都打上 2 字节的 channel 编号,broker 按编号把帧分发给对应的会话状态机。代价直接来自这个设计——同一个 channel 上的帧严格有序、串行处理,所以一个 channel 不能跨线程共享:两个线程往同一 channel 写,帧会交错,broker 解析到一半发现结构不对,报 unexpected frame 甚至断开连接。帧结构本身留到 03 章拆开。
channel 像 HTTP/2 在一条 TCP 上跑多个 stream。但 AMQP channel 没有 HTTP/2 那种独立流控:一旦 broker 因内存告警阻塞了某条连接,这条连接上所有 channel 一起被阻塞,不是只卡住一个。这个"一损俱损"的边界在 03 章流控里会再出现。
场景走查:一个订单服务跑 8 个工作线程。启动时建立 1 条 Connection,每个线程开自己的 channel(共 8 个)。线程 A 在 channel 1 上 publish,线程 B 在 channel 2 上 consume,互不干扰。若图省事让两个线程共享 channel 1——表现为偶发的 unexpected frame、消息体错位,而且低负载下不复现,上线放量才炸。
把 100 个线程压到 1 个 channel 上收发,会出什么问题?
展开答案(先停 10 秒再点)
两件事同时发生:① 帧交错损坏——多个线程的帧在同一 channel 序列里穿插,broker 解析失败;② 串行瓶颈——就算不崩,channel 内帧是串行的,100 个线程也只能排队走一条道。正解是每线程一个 channel,连接可以共享。
与下一节的关系:channel 是协议交互的通道;真正决定消息去哪的,是 channel 上声明的 exchange 和 binding。
1.2Exchange / Queue / Binding
生产者把消息发给 exchange,exchange 按 binding 规则把消息路由到 0 个或多个 queue;生产者不知道、也不关心 queue 的存在。
若生产者直接发给队列,生产者就得知道有哪些消费者、各要什么——拓扑硬编码在生产端。exchange 把"发布"和"路由"解耦:生产者只管发给一个命名的 exchange 加一个 routing key,由 binding(运维或消费端声明)决定消息落到哪些队列。新增一个消费者,只加一条 binding,生产者代码一行不动。
底层机制(比文档深一层):exchange 本身不存储消息——它是一张路由表加一个匹配函数。消息到达 exchange,broker 用该 exchange 类型对应的匹配算法,把消息的 routing key 与所有 binding 的 binding key 比对,命中的每个 queue 各收到一份消息引用。一个 binding 都没命中的消息被直接丢弃(或在 mandatory 标志下退回生产者)。所以"消息发出去了"不等于"消息进队列了"——这是新手最常栽的认知缺口。
exchange 像邮局的分拣台,binding 是"这个邮编送这条街"的规则,queue 是街道信箱。但分拣台不留底:没有匹配地址的信直接销毁,不会退回——除非你寄信时贴了 mandatory 回执要求。
场景走查:声明一个 topic 类型的 exchange orders。队列 order-email 用 binding key order.created 绑上去;队列 order-audit 用 order.# 绑上去。生产者 publish 到 orders,routing key 写 order.created:两个队列都命中(email 精确匹配,audit 的 # 匹配任意层级),同一条消息进了两个队列。换成 routing key order.shipped:只有 audit 命中,email 收不到。
"我 publish 返回成功了,broker 没报错,但队列里空的。" 根因:没有 binding 命中,消息被静默丢弃。publish 成功只表示broker 收到了,不表示进了某个队列。要确认确实入队,靠 publisher confirms(02 章)加 mandatory 标志,缺一个都看不见这次丢弃。
与下一节的关系:命不命中,取决于 exchange 的类型——四种类型,四套匹配规则。
1.3四种 exchange 类型
direct、fanout、topic、headers——区别只在一件事:拿什么和 binding key 比,怎么比。
| 类型 | 匹配规则 | 典型用途 |
|---|---|---|
| direct | routing key 与 binding key 完全相等才路由 | 按确定的 key 点对点分发,如按任务类型分队列 |
| fanout | 忽略 routing key,广播给所有绑定的队列 | 发布/订阅、广播配置变更 |
| topic | routing key 与 binding key 做模式匹配:* 匹配一个单词,# 匹配零或多个单词(以 . 分隔) | 按层级主题订阅,如 service.level 日志 |
| headers | 忽略 routing key,改用消息 headers 属性匹配,x-match=all/any | 多维度条件匹配,routing key 表达不了时 |
底层机制(比文档深一层):direct 和 fanout 本质是 topic 的两个特例——direct 等于"无通配的精确匹配",fanout 等于"全部匹配"。headers 走的是另一条匹配路径,根本不碰 routing key,而是遍历消息头字典。性能排序也由此而来:direct/fanout 最快,topic 的通配匹配略贵,headers 最灵活但最慢。选型时,能用 direct/topic 表达的,就别用 headers。
payment.error 的一条消息,被 topic exchange 按 binding key 同时投进 alerts 和 payment-log。注意:order.# 不匹配,order-flow 一份都收不到——一条消息可命中多个队列,也可能一个都不命中。场景走查(topic 通配):日志系统的 routing key 形如 <service>.<level>,例如 auth.error、payment.info。队列 all-errors 绑 *.error,payment-all 绑 payment.*,everything 绑 #。一条 payment.error 同时命中三者。
topic binding key *.error 能匹配 routing key auth.login.error 吗?
展开答案(先停 10 秒再点)
不能。* 只匹配一个单词,而 auth.login.error 是三个单词。要匹配任意前缀加 .error 结尾,得用 #.error(# 匹配零或多个单词)。这一个字符的差别,是 topic 路由最常见的失误来源。
1.4消息生命周期:ack 之后即删除
一条消息的命运:被路由进队列(Ready)→ 推送给消费者(Unacked)→ 消费者 ack → broker 永久删除。没有 offset,没有"再读一遍"。
为什么这是全教程的轴心:RabbitMQ 的队列是"消费即销毁"的。消息被 ack 后,broker 删除它,不留任何记录。这意味着 RabbitMQ 没有 offset、没有重放、没有从头再读一遍。这正是它与 Kafka 的根本分叉:Kafka 是不可变日志,消息留存、消费者用 offset 自己记进度,可以倒回去重读;RabbitMQ 是"智能 broker"——它替消费者记账(谁 ack 了、谁还欠着),记完账就把消息扔了。抓住这一点,后面所有特性都是推论。
底层机制(比文档深一层):消息在队列里有两个关键状态——Ready(已入队、待投递)和 Unacked(已推送给某消费者、等它 ack)。broker 为每一次投递分配一个 delivery tag(投递编号),挂在消费它的那个 channel 上。消费者 basic.ack(tag) → broker 删除消息;basic.nack/basic.reject → 按参数 requeue 或转入 dead-letter;消费者的 channel/连接断开 → 它名下所有 Unacked 消息自动 requeue 回队列(通常回到头部,这会破坏顺序,04 章详谈)。
消费者拿到消息、处理到一半进程崩了(没来得及 ack),这条消息会怎样?
展开答案(先停 10 秒再点)
broker 检测到该消费者的 channel/连接断开,把它名下所有 Unacked 消息 requeue,另一个消费者会重新拿到这条。后果是消费可能重复——所以消费逻辑必须幂等。代价的另一面:如果用了 auto-ack(消息一推送就当已确认),崩溃时这条消息已被删除,直接丢失、毫无痕迹。auto-ack 的取舍在 02 章。
一句话记住两者的分叉:RabbitMQ 的 broker 替你记账并在 ack 后删除消息;Kafka 的 broker 只追加日志,消费者自己记 offset。前者给你灵活路由和每消息确认,代价是没有重放;后者给你重放和高吞吐,代价是路由和每消息工作流得自己在消费端实现。这句话会在 04 章被反复用到。
1.5vhost 与三种队列类型
vhost 是命名空间隔离;queue 有三种实现——classic(单节点)、quorum(Raft 多副本,4.x 默认)、stream(可重放的追加日志)。
一个 broker 集群常被多个应用或团队共用。vhost 把 exchange、queue、权限隔成互不可见的命名空间——/prod 和 /staging 各有自己的 orders exchange,同名也互不干扰。客户端连接时指定要进哪个 vhost。它是隔离边界,不是性能边界。
队列类型这里只给地图,机制留到 03 章展开:
| 类型 | 一句话 | 什么时候用 |
|---|---|---|
| classic | 消息存在单个节点,无副本;节点挂了队列不可用 | 非关键、可容忍丢失的临时队列 |
| quorum | 基于 Raft 的多副本队列,写入要多数派确认;4.x 默认 | 需要可靠性、不丢消息的绝大多数场景 |
| stream | 追加型不可变日志,消费非破坏性(按 offset 读,不删除) | 要重放、要多个独立消费者读同一份数据 |
网上大量老教程教你用 classic mirrored queues(镜像队列)做高可用,配置里有 ha-mode、ha-policy、x-ha-policy。RabbitMQ 4.0(2024)已经彻底移除镜像队列。现在做副本/高可用用 quorum 队列。看到 ha-mode 这类配置,基本能判定那是 3.x 时代、已经不能照搬的内容。这条分界在 03 章讲清楚为什么。
§本章 self-check
先合上教程,把答案写在纸上或编辑器里。写完再点开对照——直接点开等于把这一节当再读一遍。
- 生产者 publish 一条消息时,需要指定队列名吗?为什么?它指定的是什么?
- 一条 routing key 没有任何 binding 匹配的消息,默认下场是什么?怎么让这次丢弃变得"看得见"?
- 消费者 ack 一条消息后,broker 对这条消息做了什么?这和 Kafka 消费者提交 offset 有什么本质不同?
- 为什么同一个 channel 不能被多个线程共享?根因在协议的哪一层?
答案(先做完再展开)
- 不需要队列名。它指定的是 exchange 名 + routing key。消息进哪些队列,由 binding 决定,生产者完全不感知队列的存在——这正是"发布/路由解耦"。
- 默认被静默丢弃。要让它看得见:publish 时带
mandatory标志(无法路由就退回生产者),并开启 publisher confirms 接收 broker 的basic.return。 - broker 永久删除这条消息,不留记录。Kafka 提交 offset 只是移动消费者的读指针,消息仍留在日志里,可被任何消费者按 offset 重读。RabbitMQ 没有这个能力(除非用 stream 队列)。
- 因为 AMQP 给每帧打 channel 编号、同一 channel 的帧严格有序串行处理。两个线程并发写同一 channel 会让帧交错,broker 解析出错。根因在帧层(03 章)。
用两条 binding 实现一组路由
一个 topic exchange 接收 <service>.<level> 形式的日志。要求:(a) 所有 error 级日志(任意 service)进 alerts 队列;(b) 来自 payment 服务的所有级别日志进 payment-log 队列;(c) 一条 payment.error 必须同时进这两个队列。写出 alerts 和 payment-log 各自的 binding key。
提示(卡住再展开)
alerts 用 *.error(任意 service 的 error 级);payment-log 用 payment.*(payment 的任意级别)。payment.error 同时满足两者,所以两个队列都收到——验证了"一条消息命中多个 binding"。想想:如果日志 key 可能是三段 payment.api.error,*.error 还够用吗?