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

四种模式的前三种本质上是前两章知识的重新命名:工作队列对应"一个队列多消费者加 prefetch 分摊",发布订阅对应"fanout 交换机",路由模式对应"direct 或 topic 交换机按绑定分发"。模式的价值不在新知识,而在给讨论提供公共词汇——评审时说"这是发布订阅模式,注意下游容量",比描述三个箭头高效得多。
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 模式里是配对的唯一凭证,省不得。
下一节盘点插件生态:原生不够用的能力,哪些值得引入,哪些宁可绕路。