1. 项目背景
第 5 章用 HTTP API 留下了ex.promo.direct和q.promo.sms。产品把需求摊开后,这一对远远不够:
- 订单域:支付成功只应进支付队列,发货只应进履约队列,key 写错必须能发现,不能 silently 丢进「看起来像成功」的 HTTP 200。
- 营销域:大促开场同一条「满减开始」要同时打到短信、邮件、App Push 三个窗口,且以后加「站内信」不能改生产者。
- 审计域:所有服务打日志事件,测试只要
*.error,数据组只要order.#,运维只要全量。
若继续「一个 Direct 打天下」,会出现三种事故:
生产者写死队列名 ├─ 加一个下游 = 改所有发布代码 ├─ key 写错 → 消息蒸发,订单组说「MQ 丢了」 └─ 广播做成三次 publish,其中一次超时导致部分触达Exchange 的职责只有一句:按 Binding 决定这份消息复制到哪些队列。队列才是容器。把交换机当队列用、把队列名当 Routing Key 乱用,是推广中台联调第一周的头号混乱。
约束:单机promoVHost;先用经典队列;Confirm / Ack 留给第 8、9 章。本章必须让测试能用 API 断言routed与队列深度,而不是靠「我感觉收到了」。
联调周还会出现「绑定在 UI 里点出来、发布在 Java 里写死交换机名」两套真相。营销临时加了一个q.mkt.wecom企业微信队列,却忘了绑到 Fanout,活动开始只有短信和邮件。若验收只看发布速率,谁也发现不了第三通道是空的。所以本章把绑定表当成和代码同等的交付物:没有表,就没有发布列车。
2. 项目设计
小胖把快递分拣中心的视频一甩。
小胖:这不就是快递分拣吗?单号对上就进那条传送带。为啥还要四种交换机?食堂窗口不也是看号叫人,没听说窗口还分 Direct 窗口、广播窗口。
大师:食堂窗口只有一种叫号规则。快递分拣其实有三种:按运单号精确入格(Direct)、一车货整车卸到所有出口(Fanout)、按「华东.*.易碎」这种模式入格(Topic)。还有一种按贴纸(Headers):贴了「冷冻」和「医药」才进冷库。四种不是炫技,是四种业务问法。
技术映射:Exchange Type = 匹配算法;Binding = 格口订阅条件;Routing Key = 运单上的地址栏。
小白:Binding Key 和 Routing Key 是不是同一个东西?Default 交换机又是什么,文档说「没名字」是不是空字符串?Topic 的*和#谁更贪心?一个队列能不能绑两个交换机?Fanout 还看 Routing Key 吗?匹配失败默认丢弃还是回给生产者?
大师:Binding Key 是「格口声明的规则」,Routing Key 是「这件货上写的地址」,Direct 要求二者相等,Topic 用点分单词做通配,Fanout完全忽略Routing Key。Default 交换机名字就是"",类型 Direct,隐式把每个队列以队列名绑到自己身上——所以basic.publish到默认交换机、routing_key=q.promo.sms能进该队列,这是图省事不是架构。一个队列可以绑多个交换机,一份消息匹配到 N 个队列就复制 N 份(内存/磁盘成本按份算)。匹配失败默认丢弃;要发现,必须mandatory=true等 Basic.Return,第 8 章会把 Return 做成发布器的一等公民,今天先用 HTTP API 的routed:false和 AMQP mandatory 各看一眼。
小胖:那订单用 Direct、营销用 Fanout、日志用 Topic,Headers 是不是可以扔了?听着像过度设计。
大师:Headers 适合「条件在属性里、不想把条件塞进 Routing Key 字符串」的场景,例如region=cn且channel=sms。推广中台第一期可以不做,但测试要有一条 Headers 用例,免得半年后有人用 Headers 却不知道x-match默认是all(所有键都要匹配,x-match自身不参与)。源码里还有any、all-with-x、any-with-x。
技术映射:
x-match=all是与;any是或;带-with-x才把x-*头也纳入匹配。
小白:Topic 绑定写两个#会怎样?源码注释好像有上限。另外 Fanout 绑 50 个队列,一条 1MB 消息是不是瞬间 50MB?Direct 用哈希、Topic 用树,性能差一个数量级吗?生产者能否不声明交换机,只往名字上发?alternate-exchange 和 mandatory 谁优先?
大师:rabbit_exchange_type_topic.erl里MAX_HASH_WILDCARDS为 2,过多#会让匹配爆炸,绑定校验会拦。Fanout 就是复制,大促 1MB 乘 N 是真实流量,营销广播只放小 JSON。性能:Direct 按键精确匹配(rabbit_db_binding:match_routing_key),Fanout 用'_'取该交换机全部绑定,Topic 走rabbit_db_topic_exchange:match。绑定上万条 Topic 时才需要担心,第一期日志规则二三十条够用。第 34 章再压路由热点。未声明就发布会NOT_FOUND通道异常,生产必须「先拓扑后流量」。备用交换机(AE)在无法路由时把消息转走,和 mandatory 回客户端是两条路:中台订单域优先 mandatory 让发布器失败可观测;AE 适合「进垃圾桶队列」而不是让生产者感知。两者同时存在时要写进规范选一个主策略,避免有的组以为 Return 了、其实进了 AE。
小胖:白板上就三套:订单精确、营销复印、日志模式。错误 key 必须有人喊「没格口」。今天实验就这三枪。
大师:再补第四枪:同一条支付成功,既进 Direct 业务队列,也进 Topic 审计队列——证明「一消息多绑定」是拷贝不是移动。
3. 项目实战
3.1 环境准备
沿用rabbit-promo-1(第 3 章基线),用户promo/promo_dev_2026,VHostpromo。Pythonpika>=1.3.2。
exportMQAPI=http://127.0.0.1:15672/apiexportAUTH=promo:promo_dev_2026exportVH=promo3.2 步骤一:Direct —— 订单精确投递
步骤目标:支付与发货两个 key 进不同队列;错 key 时routed=false。
# 交换机(可与第 5 章已有的 ex.promo.direct 并存或复用)curl-s-u"$AUTH"-H"content-type: application/json"-XPUT\$MQAPI/exchanges/$VH/ex.order.direct\-d'{"type":"direct","durable":true,"arguments":{}}'curl-s-u"$AUTH"-H"content-type: application/json"-XPUT\$MQAPI/queues/$VH/q.order.pay-d'{"durable":true,"arguments":{}}'curl-s-u"$AUTH"-H"content-type: application/json"-XPUT\$MQAPI/queues/$VH/q.order.ship-d'{"durable":true,"arguments":{}}'curl-s-u"$AUTH"-H"content-type: application/json"-XPOST\$MQAPI/bindings/$VH/e/ex.order.direct/q/q.order.pay\-d'{"routing_key":"pay.ok","arguments":{}}'curl-s-u"$AUTH"-H"content-type: application/json"-XPOST\$MQAPI/bindings/$VH/e/ex.order.direct/q/q.order.ship\-d'{"routing_key":"ship.ok","arguments":{}}'发布:
curl-s-u"$AUTH"-H"content-type: application/json"-XPOST\$MQAPI/exchanges/$VH/ex.order.direct/publish\-d'{"properties":{"delivery_mode":2},"routing_key":"pay.ok","payload":"PAY1","payload_encoding":"string"}'# 期望 "routed":truecurl-s-u"$AUTH"-H"content-type: application/json"-XPOST\$MQAPI/exchanges/$VH/ex.order.direct/publish\-d'{"properties":{},"routing_key":"pay.typo","payload":"LOST","payload_encoding":"string"}'# 期望 "routed":false运行结果:q.order.pay深度 1;q.order.ship为 0;错 key 不增加任何队列深度。
坑:Direct 是全等匹配,pay.ok带空格都不行。
坑:绑定在ex.promo.direct上却往ex.order.direct发,表现为 routed false,不是 Broker 坏了。
源码对照:Direct 把消息里的 routing keys 拿去精确匹配绑定。
route(#exchange{name = Name, type = Type}, Msg) -> route(#exchange{name = Name, type = Type}, Msg, #{}). route(#exchange{name = Name}, Msg, _Opts) -> Routes = mc:routing_keys(Msg), rabbit_db_binding:match_routing_key(Name, Routes).3.3 步骤二:Fanout —— 营销一发三
步骤目标:一条消息进入短信、邮件、Push 三个队列,Routing Key 随意。
curl-s-u"$AUTH"-H"content-type: application/json"-XPUT\$MQAPI/exchanges/$VH/ex.mkt.fanout\-d'{"type":"fanout","durable":true,"arguments":{}}'forqinq.mkt.sms q.mkt.mail q.mkt.push;docurl-s-u"$AUTH"-H"content-type: application/json"-XPUT\$MQAPI/queues/$VH/$q-d'{"durable":true,"arguments":{}}'curl-s-u"$AUTH"-H"content-type: application/json"-XPOST\$MQAPI/bindings/$VH/e/ex.mkt.fanout/q/$q\-d'{"routing_key":"ignored","arguments":{}}'donecurl-s-u"$AUTH"-H"content-type: application/json"-XPOST\$MQAPI/exchanges/$VH/ex.mkt.fanout/publish\-d'{"properties":{},"routing_key":"whatever","payload":"SALE_ON","payload_encoding":"string"}'运行结果:三个队列深度均为 1。Fanout 源码用'_'匹配该交换机全部绑定,不读 Routing Key:
route(#exchange{name = Name}, _Message) -> route(#exchange{name = Name}, _Message, #{}). route(#exchange{name = Name}, _Message, _Opts) -> rabbit_router:match_routing_key(Name, ['_']).坑:Fanout 上「用不同 Binding Key 分流」是无效的,要分流请用 Direct/Topic。
坑:三个队列都 durable,消息也要delivery_mode=2才谈得上重启还在(第 7 章)。
3.4 步骤三:Topic —— 日志订阅
步骤目标:order.error同时命中*.error与order.#;user.info只命中谁都不订则 routed false。
curl-s-u"$AUTH"-H"content-type: application/json"-XPUT\$MQAPI/exchanges/$VH/ex.log.topic\-d'{"type":"topic","durable":true,"arguments":{}}'curl-s-u"$AUTH"-H"content-type: application/json"-XPUT\$MQAPI/queues/$VH/q.log.errors-d'{"durable":true,"arguments":{}}'curl-s-u"$AUTH"-H"content-type: application/json"-XPUT\$MQAPI/queues/$VH/q.log.order-d'{"durable":true,"arguments":{}}'curl-s-u"$AUTH"-H"content-type: application/json"-XPOST\$MQAPI/bindings/$VH/e/ex.log.topic/q/q.log.errors\-d'{"routing_key":"*.error","arguments":{}}'curl-s-u"$AUTH"-H"content-type: application/json"-XPOST\$MQAPI/bindings/$VH/e/ex.log.topic/q/q.log.order\-d'{"routing_key":"order.#","arguments":{}}'# promo-mq/ch06/topic_pub.pyimportpika conn=pika.BlockingConnection(pika.ConnectionParameters("127.0.0.1",5672,"promo",pika.PlainCredentials("promo","promo_dev_2026"),client_properties={"connection_name":"ch06-topic"}))ch=conn.channel()forkey,bodyin[("order.error",b"E1"),("order.info",b"I1"),("user.error",b"E2")]:ch.basic_publish("ex.log.topic",key,body)print("published",key)conn.close()运行结果(文字):
| 消息 key | q.log.errors | q.log.order |
|---|---|---|
| order.error | 有 | 有 |
| order.info | 无 | 有 |
| user.error | 有 | 无 |
*匹配恰好一段;#匹配零段或多段。order.error两段,*.error命中;order.#也命中。
坑:#.error与*.error不是一回事;order*没有点,Topic 不会当通配符。
坑:源码限制过多#:
%% More than two '#' segments should not be necessary -define(MAX_HASH_WILDCARDS, 2).绑定a.#.b.#.c.#这类会在校验期失败或匹配极慢,禁止进生产。
3.5 步骤四:mandatory 看 Return(为第 8 章打样)
步骤目标:错 key +mandatory=True时客户端收到 Return,而不是以为 publish 返回了就进了队列。
# promo-mq/ch06/mandatory_miss.pyimportpika returned=[]defon_return(ch,method,props,body):returned.append((method.reply_text,body))conn=pika.BlockingConnection(pika.ConnectionParameters("127.0.0.1",5672,"promo",pika.PlainCredentials("promo","promo_dev_2026")))ch=conn.channel()ch.add_on_return_callback(on_return)ch.confirm_delivery()ok=ch.basic_publish("ex.order.direct","no.slot",b"MISS",mandatory=True)print("publish returned-from-lib",ok,"broker-returns",returned)# BlockingConnection 在 confirm+mandatory 下,未路由常以 UnroutableError 形式出现conn.close()运行结果:pika 在 confirm 模式下对不可路由常抛UnroutableError;若关闭 confirm 只开 mandatory,则走 Return 回调。两种都证明消息没进队列。第 8 章把 Confirm 与 Return 拆开讲。
坑:HTTP APIrouted:false与 AMQP mandatory 不是同一条代码路径,但验收语义一致:没格口。
坑:不设 mandatory 时,AMQP 发布「成功」只表示 Broker 收下了帧,可以零队列。
3.6 步骤五:一消息两套交换机(拷贝不是移动)
把q.order.pay额外绑到ex.log.topic,keyorder.pay:
curl-s-u"$AUTH"-H"content-type: application/json"-XPOST\$MQAPI/bindings/$VH/e/ex.log.topic/q/q.order.pay\-d'{"routing_key":"order.pay","arguments":{}}'只往ex.log.topic发order.pay:业务队列和日志队列策略不同——这里演示队列可多绑。往ex.order.direct发pay.ok不会自动进 Topic,因为那是另一扇分拣口。
坑:「绑到两个交换机」≠「发一次进两个交换机」。生产者仍要选一个交换机发布;要两套都进,要么 Fanout 前置再由各队列转、要么应用发两次,要么用交换机到交换机绑定(本期不做)。
Headers 最小用例务必跑通一次,避免半年后踩x-match默认 all:
# promo-mq/ch06/headers_demo.pyimportpika conn=pika.BlockingConnection(pika.ConnectionParameters("127.0.0.1",5672,"promo",pika.PlainCredentials("promo","promo_dev_2026")))ch=conn.channel()ch.exchange_declare("ex.hdr","headers",durable=True)ch.queue_declare("q.hdr.sms",durable=True)ch.queue_bind("q.hdr.sms","ex.hdr",routing_key="",arguments={"x-match":"all","region":"cn","channel":"sms"})props_ok=pika.BasicProperties(headers={"region":"cn","channel":"sms"})props_bad=pika.BasicProperties(headers={"region":"cn"})ch.basic_publish("ex.hdr","",b"H-OK",properties=props_ok)ch.basic_publish("ex.hdr","",b"H-MISS",properties=props_bad)q=ch.queue_declare("q.hdr.sms",durable=True,passive=True)print("headers queue depth",q.method.message_count)# 期望 1conn.close()运行结果:缺channel头的那条不进队列。源码默认走match_all,只有声明x-match=any才是或。
route(#exchange{name = Name}, Msg, _Opts) -> Headers = mc:routing_headers(Msg, [x_headers]), rabbit_router:match_bindings( Name, fun(#binding{args = Args}) -> case rabbit_misc:table_lookup(Args, <<"x-match">>) of {longstr, <<"any">>} -> match_any(Args, Headers, fun match/2); ... _ -> match_all(Args, Headers, fun match/2)3.7 完整代码清单
rabbitmq-server/column/samples/ch06/ topic_pub.py mandatory_miss.py headers_demo.py README.md # 拓扑图拓扑验收图:
[订单服务] --pay.ok--> ex.order.direct --pay.ok--> q.order.pay --ship.ok--> q.order.ship [营销] ----*----> ex.mkt.fanout --> q.mkt.sms --> q.mkt.mail --> q.mkt.push [各服务] --order.error--> ex.log.topic --*.error--> q.log.errors --order.#--------------> q.log.order3.8 测试验证
| 编号 | 操作 | 期望 |
|---|---|---|
| TC-CH06-01 | Directpay.ok | 仅 pay 队列 +1 |
| TC-CH06-02 | Direct 错 key | routed false,深度不变 |
| TC-CH06-03 | Fanout 一条 | 三队列各 +1 |
| TC-CH06-04 | Topicorder.error | errors 与 order 各 +1 |
| TC-CH06-05 | mandatory 错 key | Unroutable 或 Return,队列不增 |
| TC-CH06-06 | Headers 缺字段 | 深度不加 |
值班检查单(中文):发版前打开绑定表,核对生产者使用的交换机名与表中一致;抽一条错误 Routing Key,确认业务队列深度不变;Fanout 活动前数绑定个数是否等于触达通道数。这三步比看 Overview 的发布曲线更接近真实事故。测试还要保存一次GET /api/bindings/{vhost}的快照进 CI 产物,和 Git 里的期望绑定表做差集:多了的是手滑,少了的是漏绑。差集为空才允许发布列车继续跑第 7 章以后的用例。绑定即合同。少一条绑定就少一条触达。
curl-s-u"$AUTH"$MQAPI/queues/$VH/q.mkt.sms|rg messages4. 项目总结
优点与缺点
| 类型 | 优点 | 缺点 |
|---|---|---|
| Direct | 精确、好测、订单域默认 | 下游变多要加绑定或改 key 规范 |
| Fanout | 加下游只加队列+绑定 | 无视 key;复制放大流量 |
| Topic | 一条 key 多种订阅 | 规则一乱就重叠或漏订;#过多伤性能 |
| Headers | 条件在属性里 | 调试不直观;默认 all 易踩 |
Default"" | 演示快 | 把队列名泄漏给生产者,无法做广播演进 |
优点:1)生产者与消费者解耦。2)一消息多队列是显式拷贝。3)匹配失败可观测(mandatory/routed)。
缺点:1)四种规则混用时文档必须跟上。2)静默丢弃是默认。3)绑定错误要到运行时才发现。
适用场景
- 订单状态精确分发(Direct)。
- 营销触达多通道(Fanout)。
- 日志/事件多订阅(Topic)。
- 少量属性条件路由(Headers)。
不适用:用 Topic 模拟数据库查询;用 Fanout 传大附件;用 Default 交换机当中台总线。
注意事项
- 交换机
durable与队列durable是两件事,只持久化其中一个,重启后绑定可能残缺。 - 4.x 声明参数冲突仍是 406,先
list_bindings再改。 - 安全:configure 权限才能绑;业务账号不要给
#配置权。 - 版本:Headers 的
any-with-x是较新扩展,老客户端文档可能没有。 - Default 交换机无法在 UI 里「删掉重建」,不要把生产流量建立在它上面。
- 绑定是元数据,消息是拷贝:删绑定不会删已经在队列里的货,但会让新消息不再进来。
常见踩坑(生产)
- 生产者往错误交换机发,监控只看发布速率。速率很健康,业务队列永远 0。根因:没有
routed/mandatory 验收。处理:发布器 mandatory + 测试断言深度。 - Topic 写成
order*漏掉所有order.pay。根因:通配符只作用于点分单词。处理:key 规范评审。 - Fanout 接了会写库的重消费者,一条活动打出 20 次下单。根因:广播当 Direct 用。处理:触达与订单状态分交换机。
思考题
- 同一条消息命中同一队列两次(两个绑定都匹配),消费者会收到几条?谁负责去重?
- 若必须「发一次,Direct 业务 + Topic 审计都进」,又不想改生产者发两次,有哪些 Broker 侧选项与代价?
(答案见第 7 章附录 C。)
推广计划提示
| 部门 | 本章怎么用 | 协作 |
|---|---|---|
| 开发 | 输出《routing key 规范》;禁止 Default 交换机上生产 | 与测试共享绑定表 |
| 测试 | 主责 TC-CH06-*,错 key 必测 | 用 API 断言深度,不要只看 publish 200 |
| 运维 | 变更绑定走 definitions diff | 大 Fanout 评估复制流量 |
| 架构 | 冻结「订单 Direct / 营销 Fanout / 日志 Topic」 | 拒绝新业务再发明第四套命名 |
第 7 章进入队列本身:durable、exclusive、4.3 对临时队列的拒绝,以及堆积上限。
附录 A:完整清单与仓库位置
rabbitmq-server/column/samples/ch06/放置topic_pub.py、mandatory_miss.py。Git 提交建议带上绑定表 Markdown,便于测试做 diff。
附录 C:第 5 章思考题参考答案
题 1:routed: true为何不等于消费者成功。
缺口至少三条:① 消息可能非持久,重启丢失(第 7 章);② 尚未被 Ack,处理失败或崩溃会重投或丢失取决于 autoAck(第 8–9 章);③ 过期/拒绝可进死信,业务队列为空不代表成功(第 10 章)。此外 Confirm、磁盘、无消费者堆积都不在routed里。
题 2:VHost 含/与空格的 URL。
对 vhost、队列名、交换机名分别做 UTF-8 百分号编码(Pythonurllib.parse.quote(name, safe=''),注意safe为空才能把/编成%2F)。单测:/、%、空格、中文、+、已编码输入不要二次编码。用quote而非手工 replace。
延伸阅读与资源
Dify 从入门到进阶:LLM 应用平台实战修炼
Java 工程师进阶:从 JVM 生产排障到OpenJDK原理
NumPy 从入门到生产落地:全链路实战指南(科学计算/向量化)
Redis 8 实战精讲:从 CRUD 到源码,构建高可用缓存系统
Redis 实战修炼与原理进阶
Python 3实战精进:从脚本到高并发订单引擎
python入门:Rquests从菜鸟脚本到企业级SDK的网络实战圣经
Milvus向量数据库实战修炼:从 0 到 1精通向量检索与生产落地
MongoDB 实战进阶与内核修炼
后端工程师的 AI 转型第一课:Ollama 与私有化大模型实战
10倍开发者的 Dify 魔法书:从零构建全栈 AI 应用
后端工程师转型AI第一课-Ollama 与私有化大模型实战
大型语言模型(LLM) vLLM 高性能推理落地实战
Agent开发之LlamaIndex 实战修炼与源码进阶
大语言模型Transformers 实战修炼与源码剖析