第12章:RabbitMQ 消息属性、RPC 与 Correlation
1. 项目背景支付要先问库存服务「能否预占 2 件 SKU」。有人用同步 HTTP库存抖动时收银台一起抖。有人用 MQ 但只发 body库存回了三条响应订单服务把别人的预占结果安在当前用户头上超卖事故。还有人把订单号放在 JSON 里却不设message_id重试时库存按「新请求」再扣一次。属性不是可选彩头而是双方必须遵守的信封delivery_mode 持久化第 7 章 content_type json / 二进制避免乱解码 message_id 这条消息是谁重试仍是它 correlation_id 这次 RPC 对话的配对号 reply_to 响应打回哪里 timestamp 超时与审计 expiration 单条 TTL 字符串毫秒 priority 需队列支持 x-max-priority headers 业务扩展如 x-retry、tenantRPC 在 MQ 上的经典做法请求发到q.stock.rpcreply_to指向调用方的exclusive 回调队列第 7 章人走格子拆库存服务把同一correlation_id写回。超时则本地取消等待迟到的响应当孤儿丢掉。乱序时只认匹配的 id绝不按「先到先用」。本章仍用手动拓扑与 pika。不要把 RPC 当普通通知通知可以丢重试预占必须配对、超时、幂等。2. 项目设计小胖把餐厅取餐小票举起来。小胖这不就是叫号吗我报手机尾号厨房喊号。为啥还要 correlation_id、reply_to 这么多栏直接把库存结果发回订单队列不就行了反正都是 MQ。大师叫号必须核对小票号否则你端走邻桌的菜。订单队列是「通知广播」人人都能听RPC 响应必须回到这一次等待的人。reply_to 是你的桌号correlation_id 是这一餐的小票。厨房不看小票只按到达顺序上菜就会超卖或把失败结果给成功的人。技术映射reply_to 回调队列correlation_id 配对键message_id 请求本身的身份幂等。小白回调队列用 exclusive 还是每个调用方固定q.order.reply固定队列怎么避免别人收到我的响应超时后响应才到会不会又扣一次库存直接 reply-to 不声明队列行不行属性里 expiration 和队列 TTL 谁优先幂等键放 message_id 还是 headers 还是 bodyAMQP 0-9-1 的 RPC 和 gRPC 比对吗大师演示用 exclusive进程内并发用 correlation_id 区分多次等待。固定共享回复队列要靠 id 过滤还要防止别的服务误消费权限上只给库存 write 该队列。超时后库存仍可能执行成功——必须库存侧也用 message_id 幂等订单侧丢弃孤儿响应必要时查库存状态对账不能再发一条「看起来一样」的新 message_id。reply_to 写成队列名对方basic.publish到默认交换机 该 routing key未声明的 exclusive 名由 Broker 生成后填进去。expiration 与队列 TTL 取更短第 10 章。幂等键message_id 表示「这一封信」业务单号可同时放 headersx-order-id便于检索只放 body 时重试网关可能改 JSON 空白导致哈希变。gRPC 适合同步强依赖MQ RPC 适合解耦与削峰但多了超时与孤儿不要拿 MQ 模拟本地函数调用的所有语义。小胖那我一次预占就开一条 exclusive 队列3000 QPS 会不会把 Broker 队列数打爆大师会。高并发 RPC 应每进程一条回调队列连接级用 correlation_id 多路复用而不是每次请求 declare 新队列。这和第 4 章「连接贵、通道便宜」是同一哲学队列声明也贵。实验可以每次请求一条好懂生产规范写死进程级复用。技术映射回调队列数 ≈ 客户端进程数不是请求数。小白乱序实验怎么做库存能否在未声明 reply_to 时丢弃priority 不配 x-max-priority 有用吗content_encoding 要不要强制 utf-8大师服务端故意 sleep 短请求更久、先回后到的请求客户端必须仍匹配。无 reply_to 则库存 Ack 并打错误指标不要往猜测的队列发。priority 没有队列参数就是空操作。content_type 用application/json编码 utf-8 写进规范即可。小胖实验三枪配对正确、乱序不串单、超时丢孤儿。属性表贴进 Wiki。3. 项目实战3.1 环境准备使用第 11 章orderVHost 与app_order若尚未 provision用promo也可但正文按隔离后的账号写。库存服务可与订单同用户演示生产应拆app_stock只授 RPC 队列。# 若沿用 ch11exportVHorderexportUSERapp_orderexportPASSord_dev_20263.2 步骤一声明 RPC 工作队列步骤目标持久工作队列回调队列由客户端 exclusive 声明。# promo-mq/ch12/declare_rpc.pyimportpika connpika.BlockingConnection(pika.ConnectionParameters(127.0.0.1,5672,order,pika.PlainCredentials(app_order,ord_dev_2026)))chconn.channel()ch.queue_declare(q.stock.rpc,durableTrue)print(rpc work queue ready)conn.close()3.3 步骤二库存 worker回 echo 故意乱序步骤目标按 correlation_id 回对slowSKU sleep制造乱序。# promo-mq/ch12/stock_worker.pyimportjson,time,pika connpika.BlockingConnection(pika.ConnectionParameters(127.0.0.1,5672,order,pika.PlainCredentials(app_order,ord_dev_2026),heartbeat30))chconn.channel()ch.queue_declare(q.stock.rpc,durableTrue)ch.basic_qos(prefetch_count4)defon_req(ch,method,props,body):reqjson.loads(body)print(req,req,cid,props.correlation_id,reply,props.reply_to)ifreq.get(sku)SLOW:time.sleep(1.5)okreq.get(qty,0)10respjson.dumps({ok:ok,sku:req[sku],message_id:props.message_id})ifnotprops.reply_toornotprops.correlation_id:print(drop malformed rpc)ch.basic_ack(method.delivery_tag)returnch.basic_publish(,props.reply_to,resp.encode(),propertiespika.BasicProperties(correlation_idprops.correlation_id,content_typeapplication/json,delivery_mode1,# 回调队列短命可不 persistent),)ch.basic_ack(method.delivery_tag)ch.basic_consume(q.stock.rpc,on_req,auto_ackFalse)print(stock worker up)ch.start_consuming()坑worker 里sleep会挡心跳实验 1.5s 可接受生产用线程池。坑向已消失的 exclusive 队列 publish 会 Return要开 mandatory 或忽略。3.4 步骤三客户端配对、乱序、超时步骤目标并发两条FAST 与 SLOWSLOW 先发后到再测超时。# promo-mq/ch12/stock_client.pyimportjson,time,uuid,pikafrompika.exceptionsimportUnroutableErrorclassStockRpc:def__init__(self):self.connpika.BlockingConnection(pika.ConnectionParameters(127.0.0.1,5672,order,pika.PlainCredentials(app_order,ord_dev_2026),heartbeat30))self.chself.conn.channel()self.ch.queue_declare(q.stock.rpc,durableTrue)qself.ch.queue_declare(,exclusiveTrue)self.callbackq.method.queue self.pending{}self.ch.basic_consume(self.callback,self._on_reply,auto_ackTrue)def_on_reply(self,ch,method,props,body):futself.pending.pop(props.correlation_id,None)iffutisNone:print(orphan,props.correlation_id,body)returnfut[data]json.loads(body)fut[done]Truedefcall(self,sku,qty,timeout2.0):cidstr(uuid.uuid4())midstr(uuid.uuid4())self.pending[cid]{done:False,data:None}propspika.BasicProperties(reply_toself.callback,correlation_idcid,message_idmid,content_typeapplication/json,delivery_mode2,timestampint(time.time()),headers{x-order-id:P-RPC,x-sku:sku},)self.ch.confirm_delivery()self.ch.basic_publish(,q.stock.rpc,json.dumps({sku:sku,qty:qty}).encode(),propertiesprops,mandatoryTrue)t0time.time()whiletime.time()-t0timeout:self.conn.process_data_events(time_limit0.1)ifself.pending[cid][done]:returnself.pending.pop(cid)[data]self.pending.pop(cid,None)# 超时后迟到的视为孤儿raiseTimeoutError(frpc timeout sku{sku}cid{cid})if__name____main__:rpcStockRpc()# 先发慢的再发快的完成顺序应是 FAST 先返回但不能把 FAST 结果当成 SLOWimportthreading box{}defgo(name,sku):try:box[name]rpc.call(sku,1,timeout3.0)exceptExceptionase:box[name]str(e)t1threading.Thread(targetgo,args(slow,SLOW))t2threading.Thread(targetgo,args(fast,FAST))t1.start();time.sleep(0.05);t2.start()t1.join();t2.join()print(box,box)# 注意BlockingConnection 非线程安全生产用独立连接或串行 wait说明pikaBlockingConnection不要多线程共享第 4 章。上面 threading 仅示意乱序实验请串行先 publish SLOW 与 FAST 两个 cid在单线程process_data_events里看谁先done断言slow的 data.sku 仍是 SLOW。更稳的串行乱序脚本# promo-mq/ch12/rpc_reorder.pyimportjson,time,uuid,pika connpika.BlockingConnection(pika.ConnectionParameters(127.0.0.1,5672,order,pika.PlainCredentials(app_order,ord_dev_2026)))chconn.channel()cbch.queue_declare(,exclusiveTrue).method.queue got{}defon_reply(ch,method,props,body):got[props.correlation_id]json.loads(body)ch.basic_consume(cb,on_reply,auto_ackTrue)ids{}forskuin(SLOW,FAST):cidstr(uuid.uuid4())ids[sku]cid ch.basic_publish(,q.stock.rpc,json.dumps({sku:sku,qty:1}).encode(),propertiespika.BasicProperties(reply_tocb,correlation_idcid,message_idstr(uuid.uuid4()),content_typeapplication/json,delivery_mode2))t0time.time()whiletime.time()-t04andlen(got)2:conn.process_data_events(time_limit0.2)print(FAST maps,got[ids[FAST]][sku],SLOW maps,got[ids[SLOW]][sku])assertgot[ids[FAST]][sku]FASTassertgot[ids[SLOW]][sku]SLOWprint(reorder OK)conn.close()运行结果即使 FAST 响应先到字典仍按 cid 入座不会串单。超时实验把 worker 关掉call(..., timeout0.5)应TimeoutError随后启动 worker 打来的包打印orphan。坑超时后务必从pending删除否则迟到包会填进已放弃的请求。坑message_id每次重试若换新库存会当成新预占——重试应沿用同一 message_id。坑exclusive 回调队列不要 durable也不可多连接共享。属性在rabbit_basic里组装进内容帧不是 JSON 的一部分-export([message/3, ... properties/1, extract_headers/1, extract_timestamp/1,3.5 步骤四属性验收表测试用发一条通知非 RPC检查信封propspika.BasicProperties(delivery_mode2,content_typeapplication/json,message_idmid-1,timestampint(time.time()),headers{x-order-id:P1},expiration60000)HTTP Get 应能看到这些字段。缺content_type的发布在代码评审打回。3.6 完整代码清单column/samples/ch12/declare_rpc.py column/samples/ch12/stock_worker.py column/samples/ch12/rpc_reorder.py3.7 测试验证编号操作期望TC-CH12-01正常预占 qty≤10ok truecid 匹配TC-CH12-02SLOWFAST 乱序sku 不串TC-CH12-03worker 宕机短超时TimeoutError无脏 pendingTC-CH12-04迟到响应日志 orphanTC-CH12-05无 reply_toworker Ack 且不乱发值班检查单RPC 工作队列必须有消费者告警回调队列数异常上涨说明有人按请求 declare。超时率进仪表盘。库存幂等表按 message_id 唯一。禁止把 RPC 响应发到 Fanout。文档写明「MQ RPC 不是本地调用」产品不要假设一定 50ms 内返回。单线程 BlockingConnection 做 RPC 时消费回调与等待循环必须在同一线程跑process_data_events否则看起来像「MQ 丢响应」其实是客户端没泵事件。这和第 9 章 sleep 挡心跳是亲戚问题。信封规范建议打印成一页纸必填 message_id、content_type、delivery_modeRPC 另加 correlation_id 与 reply_to重试不得更换 message_idheaders 里的业务键只增不改名。代码评审对照这一页比看 JSON schema 更快。库存预占的成功响应也要带原 message_id便于订单侧日志对上。超时后的补偿不是再发一条全新 id而是查询库存状态或带着旧 id 重试。产品若要求「用户连点两次必须扣两次」那不是幂等场景应拆成两笔预占而不是共享 id。RPC 工作队列的消费者数告警与通知队列不同RPC 没人接会导致收银台线程/协程堆超时比短信晚到更伤转化。值班把q.stock.rpc的 consumers0 定为 P1短信队列可为 P2。不要用同一套「堆积深度」阈值套两种队列RPC 深度一小就该响因为它表示正在等待的用户。回调队列数突然升高去查是否有人按请求 declare exclusive。Direct Reply-to 能缓解队列数但排障更不直观团队不熟就先 exclusive 复用。乱序测试必须自动化断言 sku 字段禁止用肉眼看打印顺序。CI 里 worker 与 client 的启动竞争要用就绪探针队列存在且 consumers≥1再发请求否则用例红在「没 worker」而不是乱序逻辑。超时用例要关 worker 而不是把 timeout 设得比 sleep 还长还期望失败——那种用例测的是网络而不是客户端 pending 清理。孤儿响应打 metrics不要只 print。属性与 body 的职责切分能放信封的不要只放 JSON否则多语言客户端解码失败时连配对键都拿不到。correlation_id 尤其必须在属性上。body 可以缺字段信封不能缺。审计系统若要检索可冗余一份到 headers但主键仍是属性。RPC 超时时间应小于用户接口超时才能在 HTTP 层返回明确「预占未确认」而不是整串 504。超时后页面提示「请刷新订单状态」后台对账而不是让用户连点造成新的 message_id。把产品文案和 RPC 超时一起评审否则技术做对了客服仍教用户狂点。库存 worker 的 prefetch 不宜过大预占是重操作。prefetch4 的示例是起点按 DB 耗时调。与第 9 章相同窗口过大只是把等待从队列挪到 unacked。RPC 的 unacked 告警应比通知更敏感。worker 进程滚动时未 Ack 会回流新实例必须幂等这和第 9 章滚动双发是同一类事故只是伤的是库存不是短信。客户端 pending 字典要有上限防止超时路径漏删导致内存涨。定期扫描超时条目。这是应用层的 connection_max与 Broker 无关但故障看起来像 MQ 泄漏。排障时同时看两边。4. 项目总结优点与缺点做法优点缺点exclusive 回调 correlation_id隔离好、实现直观进程内要泵事件高并发需复用队列共享回复队列队列数少过滤与权限更难同步 HTTP语义简单与第 1 章相同的耦合只用 body 里的订单号少几个字段重试、乱序、网关改写 JSON 时脆弱优点1信封标准化可测。2乱序可证明。3超时与孤儿有明确行为。缺点1至少一次仍要幂等。2pika 线程模型易踩。3RPC 放大队列与通道成本。适用场景库存预占、风控询问等需要应答的解耦调用。通知类消息的信封规范即使不用 RPC。测试要断言配对而非到达顺序。不适用用 RPC 替代数据库事务每请求一条 exclusive 队列打满大促把 correlation_id 当 message_id 混用。注意事项correlation_id 管对话message_id 管信本身两者都要。expiration 类型是字符串。4.x 无 immediate。安全reply_to 不要指向别人的可写队列名伪造回调。回调队列权限库存需能向其 publish通常默认交换器 队列名VHost 内 write 正则要覆盖生成名exclusive 名随机这是正则.*对 reply 的压力。实践库存账号对amq.gen-*或改用 direct reply-to 扩展。若正则过死RPC 回调会 403。订单 VHost 内可给库存amq.gen-.*write或使用 RabbitMQDirect Reply-toamq.rabbitmq.reply-to减少临时队列。测试至少验证一种。Direct Reply-to 可作加分实验客户端reply_toamq.rabbitmq.reply-to无需 declare exclusive。不展开源码生产选型时与 exclusive 对比连接亲和。常见踩坑生产按到达顺序取响应串单超卖。根因无 correlation_id。每次请求 declare 队列队列数打满。根因把 exclusive 当请求级。超时又重试却换了 message_id库存双扣。根因幂等键不稳定。思考题Direct Reply-to 在连接断开时尚未取走的响应去哪与 exclusive 队列对比。若库存 worker 多实例如何保证同一 message_id 只预占一次MQ 能代替数据库唯一约束吗附录 C第 11 章思考题参考答案题 1跨 VHost 看见同一支付成功。Shovel/Federation 按权限最小搬运第 21 章或应用双发要幂等不要给订单应用 marketing 的 write。复制的是消息拷贝不是共享队列。题 2topic permission。可以普通 write 允许交换机名topic 权限再限制 routing key 模式。这样能发order.error不能发order.info。测试要单独声明 topic permission只测普通正则不够。延伸阅读与资源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 实战修炼与源码剖析