5.1 消息模式:工作队列、发布订阅与 RPC


文档摘要

5.1 消息模式:工作队列、发布订阅与 RPC 本节摘要:消息模式是反复出现的拓扑套路的行业速记法:工作队列分摊任务、发布订阅扇出事件、路由模式按类分发、RPC 模式在异步管道上重建同步语义。本节给出四模式的结构对照与选型口径,重点拆解实现最繁也最易错的 RPC 模式。 工程化整理的第一刀,砍向"每次都重新设计拓扑"。仔细观察会发现,业务需要的拓扑其实就几类——别人早把它们画成了模式图。掌握模式语言后,设计讨论从"画图比划一小时"变成"用 RPC 模式加优先级"一句话。

5.1 消息模式:工作队列、发布订阅与 RPC

本节摘要:消息模式是反复出现的拓扑套路的行业速记法:工作队列分摊任务、发布订阅扇出事件、路由模式按类分发、RPC 模式在异步管道上重建同步语义。本节给出四模式的结构对照与选型口径,重点拆解实现最繁也最易错的 RPC 模式。

工程化整理的第一刀,砍向"每次都重新设计拓扑"。仔细观察会发现,业务需要的拓扑其实就几类——别人早把它们画成了模式图。掌握模式语言后,设计讨论从"画图比划一小时"变成"用 RPC 模式加优先级"一句话。

四种模式的形态对照

图 13 四种消息模式结构对照

图 8 四种消息模式结构对照

四种模式的前三种本质上是前两章知识的重新命名:工作队列对应"一个队列多消费者加 prefetch 分摊",发布订阅对应"fanout 交换机",路由模式对应"direct 或 topic 交换机按绑定分发"。模式的价值不在新知识,而在给讨论提供公共词汇——评审时说"这是发布订阅模式,注意下游容量",比描述三个箭头高效得多。

RPC 模式:最值得细拆的一个

RPC 模式解决一个看似矛盾的需求:团队想保留消息队列的解耦与缓冲,又需要"调用并拿回结果"的同步体验。它的实现利用了两个此前的知识点:回调队列(2.4 节的 exclusive 队列)与关联编号(2.1 节的 correlation_id 属性)。

背景:图片处理服务耗时数秒,前端需要同步等待处理结果。

操作:调用方代码:

import pika, uuid, json connection = pika.BlockingConnection(pika.ConnectionParameters(host="localhost")) channel = connection.channel() # 排他自动删除的回调队列:只属于本连接,断开自动清理 result = channel.queue_declare(queue="", exclusive=True, auto_delete=True) callback_queue = result.method.queue corr_id = str(uuid.uuid4()) response_holder = {} def on_response(ch, method, props, body): # 只有编号匹配的响应才认领,防止串包 if props.correlation_id == corr_id: response_holder["data"] = body ch.basic_ack(delivery_tag=method.delivery_tag) channel.basic_consume(queue=callback_queue, on_message_callback=on_response, auto_ack=False) channel.basic_publish( exchange="", routing_key="image.process.q", # 服务方监听的请求队列 body=json.dumps({"image_key": "cdn-bucket-8848/example-001.png"}).encode(), properties=pika.BasicProperties( reply_to=callback_queue, # 告诉服务方把结果发回哪 correlation_id=corr_id)) # 请求编号 # 等待响应帧(生产代码应设超时并做降级) connection.process_data_events(time_limit=30) print("处理结果:", response_holder.get("data", b"timeout").decode())

服务方是普通消费者,处理完向 props.reply_to 指定的队列回发结果,带上同一个 correlation_id——一对属性完成了"请求响应在异步管道上配对"的全部魔法。

结果与解读:结果正常返回。但请注意解读里的三个代价。代价一,排他回调队列与连接绑定:客户端多进程或重启后,回调队列消失,在途响应作废——客户端的重试逻辑必须考虑这一点。代价二,超时与降级必须显式设计:示例里的 time_limit 是省略号,生产代码要处理"服务方永远没回"的场景,超时策略与同步调用一样存在。代价三,吞吐受限:每个请求占一条请求队列消息与一次等待,RPC 模式不适合高并发低延迟场景——那本来就是同步调用的主场。

变式一:并发 RPC——发起多个请求后循环 process_data_events 收集响应,按编号配对,吞吐立涨,但代码复杂度同步上涨,先确认真的需要。变式二:用 AsyncBlockingConnection 或框架层(如 Spring 的 AsyncRabbitTemplate)重写,回调配对交给框架——这是工程上的正解,5.5 节会看到 Java 版本。

选型口径

模式选择的判定树很简单:一个任务多人干,选工作队列;一件事多方知,选发布订阅;一类消息一类处理,选路由;要结果但不想要耦合,选 RPC;要结果又要低延迟,回头用同步调用。 最后一条是护栏:RPC 模式不改变"跨网络拿结果"的物理延迟,用它包装高频同步调用是把消息队列用回它们想逃离的样子。

⚠️ 常见坑:RPC 请求消息不带 correlation_id 就发。服务方回不了包还不知道自己丢了什么,调用方超时后重试又造成重复处理。编号属性在 RPC 模式里是配对的唯一凭证,省不得。

本节要点回顾

  • 模式即词汇:四种套路覆盖绝大多数拓扑需求,评审沟通成本大降;
  • 前三模式:是队列、fanout、topic 的工程化命名,重点是识别与套用;
  • RPC 三件套:请求队列、回调队列、关联编号,缺一配不上对;
  • RPC 三代价:回调队列随连接生灭、超时须显式设计、吞吐天然受限;
  • 选型护栏:低延迟强同步的场景,正确答案是不用消息队列;
  • 模式混用:真实拓扑是模式的组合件,先用模式名搭骨架再补定制点,约定的守护靠团队基础设施库。

下一节盘点插件生态:原生不够用的能力,哪些值得引入,哪些宁可绕路。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U