1. 项目背景
推广中台营销 Topic 曾经「绑着绑着就慢」:审计日志用 stock.#,有人再加 *.*.*.# 这种多段 #,发布 P99 被路由吃掉,值班却去调 prefetch。另一拨人打开 mandatory 后发现 Return 了,以为队列丢消息,其实 绑定表就没有那条 key。还有人给交换机加了 alternate-exchange,死信和 AE 搅在一起,测试分不清是未路由还是业务失败。
痛点:
publish rk=stock.change.sku
│
├─ Direct:精确命中或空
├─ Fanout:忽略 rk,所有绑定
├─ Topic:Trie 走 * 与 #,绑定越多越贵
└─ Headers:扫绑定 args 做 all/any
│
▼
空列表 + mandatory → Return(不是「MQ 丢了」)
空列表 + AE → 再进另一个交换机(要防环)
一万条 Topic 绑定 → CPU 在 rabbit_db_topic_exchange:match
四种交换机声明人人会做;本章把匹配 落到表与树:绑定持久在元数据存储,运行时 Direct/Fanout 走 rabbit_db_binding:match_routing_key,Topic 走独立 Trie 投影。4.x Topic 绑定 最多两个 # 段(?MAX_HASH_WILDCARDS = 2),再多 validate_binding 直接 {error, binding_invalid}——这是源码硬限制,不是风格建议。
业务约束:压测用 marketing 实验室交换机 ex.ch34.topic,不要在生产审计 Topic 上造 1 万绑定。测试要对比 Direct 精确匹配与 Topic 通配的耗时数量级,不要求实验室打出生产级纳秒报表。
大促前营销还容易犯「一个 Topic 包打天下」:支付状态、库存变更、优惠券、风控灰名单全绑到同一交换机,绑定数随活动线性涨,路由 CPU 把支付 Direct 同节点的调度挤瘦。正确拆法是:一种业务问法一种交换机类型。本章源码证明 Topic 的成本在 Trie 与绑定基数,不在「交换机名字好不好听」。运维侧要给 每交换机绑定数 做告警,超过实验室标定(例如单交换机数千精确键)就该考虑拆 Direct 或分区流,而不是先加 CPU。实验室绑定造完必须删除,禁止留在营销生产交换机上过夜。
2. 项目设计
小胖把快递面单上的「省.市.区」和食堂「窗口号」对比。
小胖: Direct 不就是窗口号对上就送吗?Fanout 不就是大喇叭,绑了的窗口人手一份?Topic 不就是用 * 当通配符,跟 glob 一样。为啥还要 Trie,还要限制两个 #。Headers 听着像 HTTP,大促谁用啊。AE 不就是没人收就转前台吗,跟 DLX 有啥区别。
大师: Direct/Fanout 的直觉对:源码里 Fanout 甚至是 match_routing_key(Name, ['_'])——用通配键取出该交换机上全部目的地。Topic 不是每次把所有绑定当 glob 扫一遍(绑定上万时会爆),而是 按点分段走进 Trie(rabbit_db_topic_exchange,Khepri 投影表)。每个额外 # 让匹配器做的工作倍增,文件头写明 MQTTv5 都只允许至多一个 #,Broker 放宽到两个封顶。Headers 给「按消息头分流」用,大促主路径很少靠它吃吞吐,但客服工单、灰度标签会用;匹配是 rabbit_router:match_bindings 加谓词,绑定多时比 Direct 更贵。AE 是 路由结果为空 时把交换机当下一跳再 route;DLX 是消息 从队列死信出去。一个发生在进队前,一个发生在出队/拒绝/过期后,画在同一箭头上就会把未路由当成支付失败。
技术映射: Direct = 键等值;Fanout = 全绑定;Topic = Trie;Headers = 扫 args;AE = 空结果再路由;DLX = 队列死信。
小白: rabbit_router 现在几乎只是转发到 rabbit_db_binding?默认交换机为什么不用绑定表?装饰器何时插入?AE 如何防环?return_binding_keys 这个 option 谁用?mandatory Return 在 channel 哪一段、还是 route 里?1 万绑定会不会把 Khepri 打满?
大师: rabbit_router:match_routing_key/2 就是一行 rabbit_db_binding:match_routing_key;Fanout/Direct 类型模块直接调 DB 或经 router。默认交换机(空名)在 rabbit_exchange:route/3 第一子句:Routing Key 当队列名 拼 #resource{kind=queue},不查绑定。装饰器 rabbit_exchange_decorator:select(route, Decorators),process_decorators 把额外目的地并进工作列表,consistent-hash 等插件走这条。AE:process_alternate 仅当 ExchangeDests 为空才取 alternate-exchange 参数/Policy,把 另一个交换机 放进工作列表;process_route 用 gb_sets 记见过的交换机,环会停。return_binding_keys 让调用方拿到「命中了哪条绑定键」,channel 的 publish 路径传 true,方便后续头信息。Return 不在 route 里发,route 只给队列名列表;channel deliver_to_queues 发现 mandatory 且无目标再让 writer 发 basic.return。1 万绑定是实验室数量级,元数据在 Khepri,内存投影也占 RAM,所以实验用完 删交换机,不要留在营销生产名上。
小胖: 实验:四种类型各打一条对照;Topic 造 1 万绑定测耗时;错误 rk + mandatory 看 Return;AE 接到 ex.ch34.ae。别在支付 Direct 上绑 1 万条。
大师: 支付保持精确 Direct。营销 Topic 绑定规范写进评审:禁止超过两个 #,禁止 stock.#.#.# 这种过时写法(声明阶段就会被拒)。
技术映射: route/3 返回队列集合(可带 binding_keys);空集合不是异常,是「没人收」。
小白: Topic 匹配热点函数具体是哪几个?重复目的地谁去重?同一消息进两个队列是拷贝还是引用?Headers 的 x-match=all-with-x 和默认 all 差在哪?
大师: 热点:split_topic_key_binary(按 . 切,pattern 缓存在 persistent_term)、trie_match。类型模块注释写明 route 可以返回重复 destination,调用方负责去重——rabbit_exchange:route 用 map 键收队列名。进两个队列是 两份投递(经典队列两份存储,第 35 章),不是零拷贝共享。Headers:默认 all 忽略 x- 开头的绑定参数当匹配条件;all-with-x / any-with-x 连 x- 头也参与(match_x),x-match 自身仍 skip。
小白: 默认交换机走队列名,会不会绕过权限检查?绑定删除后路由表多久失效?装饰器返回的交换机再 AE,会不会绕过 gb_sets 环检测?压测 1 万绑定是否该关 Confirm 以免测到的是刷盘不是路由?
大师: 默认交换机仍然走 channel 的 check_write_permitted,队列资源名要有 write;不是「空名就免检」。绑定变更经 Khepri 投影,实验室几乎立刻,测试不要依赖「睡 1ms」。装饰器目的地同样进 process_route,交换机会进 Seen 集合;队列名进 map。压测对比路由必须 同一 Confirm 设置,否则一组开 Confirm 一组不开,数字没有可比性;本章 bench 两边都 confirm_delivery。
技术映射: 权限在 channel;匹配在 exchange type / DB;去重与破环在 route1。
小胖: 检查单:耗时对比表、Return、AE、绑定非法 # 被拒。收工。
3. 项目实战
3.1 环境准备
| VHost | marketing(或实验室专用 order 下的实验交换机,避免污染支付) |
| 用户 | app_mkt |
| 目录 | promo-mq/ch34/ |
| 对照 | rabbit_exchange.erl 的 route/3、rabbit_exchange_type_topic.erl、rabbit_db_topic_exchange.erl |
建议单独交换机名,实验结束删除。
3.2 步骤一:阅读 route/3 骨架
步骤目标: 能讲清默认交换机、类型回调、AE、装饰器四条支路。
%% rabbit_exchange.erl
route(#exchange{name = #resource{name = ?DEFAULT_EXCHANGE_NAME, virtual_host = VHost}},
Message, _Opts) –>
RKs = lists:usort(mc:routing_keys(Message)),
[rabbit_misc:r(VHost, queue, RK) || RK <– RKs];
route(X = #exchange{decorators = Decorators}, Message, Opts) –>
Decs = rabbit_exchange_decorator:select(route, Decorators),
QNamesToBKeys = route1(Message, Decs, Opts, {[X], XName, #{}}),
%% Opts 含 return_binding_keys 时返回 {QName, #{binding_keys := …}}
%% route1 内:
%% ExchangeDests = Type:route(X, Message [, Opts])
%% DecorateDests = process_decorators(…)
%% AlternateDests = process_alternate(X, ExchangeDests) %% 仅 Dest 为空
%% rabbit_exchange_type_direct.erl
route(#exchange{name = Name}, Msg, _Opts) –>
rabbit_db_binding:match_routing_key(Name, mc:routing_keys(Msg)).
%% rabbit_exchange_type_fanout.erl
route(#exchange{name = Name}, _Message, _Opts) –>
rabbit_router:match_routing_key(Name, ['_']).
运行结果: 读者能在 IDE 跳转上述函数。 坑: 把 rabbit_router 当成「还有一套独立算法」——Direct 已绕过它直调 DB,Fanout 经 router 到同一 DB 接口。
3.3 步骤二:Topic Trie 与 # 上限
步骤目标: 非法绑定被拒;合法绑定能匹配。
# promo-mq/ch34/topic_validate.py
import pika
from pika.exceptions import ChannelClosedByBroker
c = pika.BlockingConnection(pika.ConnectionParameters(
"127.0.0.1", 5672, "marketing",
pika.PlainCredentials("app_mkt", "mkt_dev_2026"),
client_properties={"connection_name": "ch34-val"}))
ch = c.channel()
ch.exchange_declare("ex.ch34.topic", "topic", durable=True)
ch.queue_declare("q.ch34.ok", durable=True)
ch.queue_bind("q.ch34.ok", "ex.ch34.topic", "stock.#")
try:
ch.queue_bind("q.ch34.ok", "ex.ch34.topic", "a.#.b.#.c.#")
print("UNEXPECTED success")
except ChannelClosedByBroker as e:
print("rejected", e.reply_code, e.reply_text)
c.close()
源码:
%% rabbit_exchange_type_topic.erl
–define(MAX_HASH_WILDCARDS, 2).
validate_binding(_X, #binding{key = BindingKey}) –>
Words = rabbit_db_topic_exchange:split_topic_key_binary(BindingKey),
case count_hash_wildcards(Words) of
N when N > ?MAX_HASH_WILDCARDS –>
{error, {binding_invalid, "Topic binding key '~ts' uses ~b '#' wildcards, "
"at most ~b are allowed", [BindingKey, N, ?MAX_HASH_WILDCARDS]}};
匹配:
match(...) –>
Words = split_topic_key_binary(RoutingKey),
%% Khepri 投影版本 >=4:trie_match(TrieTab, BindingTab, …)
运行结果: stock.# 成功;三个 # 的绑定 406/绑定无效。 坑: 有人用队列 arguments 当绑定键,声明成功但永远匹配不到。
3.4 步骤三:1 万绑定耗时对比
步骤目标: 同一发布循环下,Direct 精确键 vs Topic 多绑定的相对耗时。
# promo-mq/ch34/bench_route.py
import time, pika
def conn():
return pika.BlockingConnection(pika.ConnectionParameters(
"127.0.0.1", 5672, "marketing",
pika.PlainCredentials("app_mkt", "mkt_dev_2026"),
heartbeat=30, client_properties={"connection_name": "ch34-bench"}))
def setup(n=10000):
c = conn(); ch = c.channel()
ch.exchange_declare("ex.ch34.direct", "direct", durable=True)
ch.exchange_declare("ex.ch34.topic", "topic", durable=True)
ch.queue_declare("q.ch34.d", durable=True)
ch.queue_declare("q.ch34.t", durable=True)
ch.queue_bind("q.ch34.d", "ex.ch34.direct", "exact.key")
for i in range(n):
ch.queue_bind("q.ch34.t", "ex.ch34.topic", f"seg.{i}.#")
c.close()
def bench(ex, rk, n=2000):
c = conn(); ch = c.channel(); ch.confirm_delivery()
t0 = time.perf_counter()
for _ in range(n):
ch.basic_publish(ex, rk, b"x", mandatory=True)
ch.wait_for_confirms()
dt = time.perf_counter() – t0
c.close()
return dt
if __name__ == "__main__":
setup(10000)
print("direct", bench("ex.ch34.direct", "exact.key"))
print("topic", bench("ex.ch34.topic", "seg.9999.leaf"))
运行结果: Direct 明显更快;Topic 在 1 万绑定时发布循环更慢(实验室记录两列数字即可)。热点在 rabbit_db_topic_exchange:match/3 与 trie_match。 把两列耗时、绑定条数、Broker CPU 截图贴进 promo-mq/ch34/bench.txt。数字只用于 相对比较,不得写进第 30 章容量评估当生产 TPS。若 Topic 并不比 Direct 慢多少,检查是否 rk 未命中而在走 Return,或绑定其实没创建成功。
坑: 用 mandatory+未命中去测 Topic,测到的是 Return 路径。rk 必须命中其中一条。 坑: 绑定循环不要每条一个队列,否则队列数爆炸;本脚本多绑定到同一 q.ch34.t。 坑: 实验后删除 ex.ch34.* 与队列,避免管理面卡顿。
# 清理
python – <<'PY'
import pika
c=pika.BlockingConnection(pika.ConnectionParameters("127.0.0.1",5672,"marketing",
pika.PlainCredentials("app_mkt","mkt_dev_2026")))
ch=c.channel()
for q in ("q.ch34.d","q.ch34.t","q.ch34.ok","q.ch34.ae"):
ch.queue_delete(q)
for x in ("ex.ch34.direct","ex.ch34.topic","ex.ch34.ae"):
ch.exchange_delete(x)
c.close()
print("cleaned")
PY
3.5 步骤四:mandatory Return 与 AE
步骤目标: 空路由与 AE 再投递分开验收。
# promo-mq/ch34/ae_and_return.py
import pika
from pika.exceptions import UnroutableError
c = pika.BlockingConnection(pika.ConnectionParameters(
"127.0.0.1", 5672, "marketing",
pika.PlainCredentials("app_mkt", "mkt_dev_2026"),
client_properties={"connection_name": "ch34-ae"}))
ch = c.channel()
ch.confirm_delivery()
ch.exchange_declare("ex.ch34.ae", "fanout", durable=True)
ch.queue_declare("q.ch34.ae", durable=True)
ch.queue_bind("q.ch34.ae", "ex.ch34.ae", "")
ch.exchange_declare("ex.ch34.miss", "direct", durable=True,
arguments={"alternate-exchange": "ex.ch34.ae"})
ch.basic_publish("ex.ch34.miss", "no.such", b"via-ae", mandatory=False)
print("ae path published")
try:
ch.exchange_declare("ex.ch34.empty", "direct", durable=True)
ch.basic_publish("ex.ch34.empty", "no.such", b"return-me", mandatory=True)
except UnroutableError:
print("return as expected")
c.close()
运行结果: q.ch34.ae 深度 +1;无 AE 的 empty 交换机 mandatory 走 Return。 坑: AE 与 DLX 名字都叫「备用」,评审必须写清阶段。 坑: AE 指回自己或互相指,靠 gb_sets 破环,表现为 静默不投递,比死循环好,但业务仍丢,要监控 unroutable。
3.6 步骤五:Headers 最小对照
# promo-mq/ch34/headers_min.py
import pika
c = pika.BlockingConnection(pika.ConnectionParameters(
"127.0.0.1", 5672, "marketing",
pika.PlainCredentials("app_mkt", "mkt_dev_2026"),
client_properties={"connection_name": "ch34-hdr"}))
ch = c.channel()
ch.exchange_declare("ex.ch34.hdr", "headers", durable=True)
ch.queue_declare("q.ch34.hdr", durable=True)
ch.queue_bind("q.ch34.hdr", "ex.ch34.hdr", "", arguments={
"x-match": "all", "kind": "coupon", "lane": "a"})
ch.confirm_delivery()
ch.basic_publish("ex.ch34.hdr", "", b"h", mandatory=True,
properties=pika.BasicProperties(headers={"kind": "coupon", "lane": "a"}))
print("headers ok")
c.close()
源码 x-match:any / any-with-x / all-with-x / 默认 all。 坑: 消息头缺字段时 all 失败,mandatory 则 Return,容易误报「交换机坏了」。
3.7 测试验证
| TC-CH34-01 | 三个 # 绑定 | 失败 |
| TC-CH34-02 | stock.# 匹配 stock.change | 入队 |
| TC-CH34-03 | bench Direct vs Topic | Direct 更快,数字入库 |
| TC-CH34-04 | mandatory 空路由 | Unroutable |
| TC-CH34-05 | AE | 备用队列 +1 |
| TC-CH34-06 | 清理实验对象 | list 不再出现 ex.ch34.* |
4. 项目总结
优点与缺点
| Direct | 匹配 O(键),支付主路径 | 无模式 |
| Fanout | 实现极简 | 绑定全量拷贝,扇出放大存储 |
| Topic | 灵活审计 | Trie/绑定规模敏感,# 有硬上限 |
| Headers | 按属性分流 | 扫绑定,不适合超高 TPS 主路径 |
| AE | 未路由兜底 | 与 DLX 易混;环被静默吃掉 |
rabbit_exchange:route 统一工作列表,插件类型只要实现 rabbit_exchange_type 回调(第 39 章)。
适用场景
- 解释「发布成功但队列空」是未路由还是没消费者。
- 评审 Topic 绑定规范(# 个数、绑定基数)。
- 排障 AE/DLX 混淆。
不适用:用 1 万 Topic 绑定当支付路由;用 Headers 扛大促主路径。
绑定变更要进变更单:谁在哪个 VHost 给哪个交换机加了 # 通配,影响面等于「未来所有匹配该模式的消息多一份拷贝」。管理面手工点绑定和生产代码 declare 双通道,最容易出现「测试环境多一条幽灵绑定」。CI 应对交换机绑定做快照 diff(第 6 章做法),综合实战第 31 章的营销 Fanout 三队列也要纳入差集检查,防止少绑站内信。
注意事项
- Topic route 可重复 destination,外层 map 去重。
- 默认交换机不建绑定。
- Policy 里的 AE 与参数合并走 rabbit_policy:get_arg。
- 绑定存在 Khepri,节点内存有投影,删绑定要等投影更新,测试别瞬间断言。
- 安全:绑定权限是 configure,乱绑等于乱路由。
- 1 万绑定实验结束后必须删交换机,否则管理面与内存投影一直为实验室数据付费。
- 路由压测与容量压测分开:前者比 Direct/Topic,后者用第 30 章 Persist+Confirm 口径。营销 Topic 的绑定变更走评审,禁止大促当天在管理面手工加 #。Headers 只给灰度/工单类低 TPS 路径,主路径继续 Direct。AE 与 DLX 的评审词必须写成「进队前兜底」和「出队后死信」,禁止在架构图上共用一个「备用」框。实验室 bench 记录相对倍数即可,禁止把 1 万绑定的本地数字写进采购单。清理脚本失败则第二天管理面仍卡,算本章验收未通过。绑定快照 CI 失败不得合并支付相关拓扑变更。
常见踩坑(生产)
思考题
附录 A:完整清单与绑定规范
promo-mq/ch34/
topic_validate.py
bench_route.py
ae_and_return.py
headers_min.py
cleanup.py # 删除 ex.ch34.* / q.ch34.*
Topic 绑定规范(贴进开发 Wiki,代码评审用):
思考题 1 要点:精确键堆在 Topic 上仍走 Trie,通常应迁 Direct,代价是客户端 rk 与绑定表重做、以及迁移窗口双写。思考题 2 要点:Fanout 的 route 几乎是一次全绑定取出,贵在 每条队列一份投递与存储(第 35 章)。
LangChain从入门到进阶实战之旅 SQLAlchemy 2.0从入门到进阶的实战之旅 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 实战修炼与源码剖析
网硕互联帮助中心



评论前必须登录!
注册