做后端时间长了你会发现不同团队对“发布/订阅模式”的认知经常对不上。有人觉得它不过就是换了个方式写回调有人觉得必须上重型消息中间件才算数。其实 Pub/Sub 是一种消息传递模式解决的问题一直很明确让消息的生产者不再关心消息由谁处理、什么时候处理、在哪个节点处理。很多人拿它和观察者模式对比两者确实长得像但 Pub/Sub 更强调解耦尤其适合分布式系统里的消息分发场景。我自己在某个订单事件系统上反反复复改过好几次架构从最开始的同步调用到后来自己写内存事件总线再到引入独立的消息中间件整个过程踩了很多坑。这篇文章不讲大道理就把我对发布/订阅模式的理解、设计思路、最小可用实现和实战中遇到的问题摊开说清楚。无论你刚开始接触还是已经在用消息队列应该都能从中找到可以落地的参考。1. 先想清楚发布/订阅模式到底解决什么问题1.1 观察者模式与 Pub/Sub 的界限观察者模式是面试题常客主题对象维护一组观察者状态变化时逐个通知。听起来和 Pub/Sub 很像但两者在分布式系统里的表现完全不同。观察者模式通常工作在单个应用进程内部观察者列表由主体对象自己持有对象之间要么存在直接引用要么通过事件回调接口建立联系。对应到代码层面一个订单状态变化了订单对象直接调用已注册的那些监听器。Pub/Sub 之所以被单独拎出来是因为它把观察者模式里隐含的那些“大家都认识”的假设全部拆掉了。订阅者不认识发布者发布者也不维护订阅者名单中间还有一层独立的通道负责转发。这个通道可以是一个线程内的列表也可以是一台独立部署的消息中间件。用生活化的类比讲观察者模式像你在公司群里直接 某个人“版本发完了你记得去验收”。你认识这个人也知道他的职责范围消息能不能送达取决于对方当时在不在线。Pub/Sub 则像在公告栏贴了一张公告“版本已发布”谁要响应谁自己来看。你不需要知道谁会关注也不关心对方当时在不在公告栏自己会保管消息。1.2 解耦到底解了什么很多人把“解耦”理解为代码层面不写依赖关系但这在分布式场景里远远不够。Pub/Sub 带来的解耦至少有三层含义每一层都在解决一类实际问题。第一层是时间解耦。发布者发完消息就能走订阅者可以稍后上线再消费。放在观察者模式里做不到观察者不在对象的内存里就收不到通知。对应到消息中间件上消费者的进程挂了一段时间重启之后还能从断点继续消费正是因为消息被持久化保存了。这个能力对于异步任务特别重要比如凌晨生成的报表任务消费者当时不在线没关系服务恢复后再处理就行。第二层是空间解耦。双方不需要知道对方的网络地址、实例数量、端口号只认一个逻辑上的主题或队列名。服务扩容缩容的时候只改变订阅方实例数发布者配置完全不动。做过微服务的人都有体会如果每次上游扩容都要回去改下游的调用地址整个发布流程就会变成噩梦。有了消息中间件做中转上下游之间的物理拓扑被彻底屏蔽。第三层是流量解耦。生产者产生的消息峰值可能非常尖锐比如大促开始的那几秒流量瞬间冲上来。如果没有中间缓冲下游服务的处理带宽就会被直接打爆。消息中间件在这里充当了一个非常大的蓄水池把瞬时流量分散到后续时间内平滑消费。这是同步请求/响应模式永远给不了的能力也是很多系统必须上消息队列的根本原因。1.3 和点对点队列、请求响应模式的区别实际项目里经常把 Pub/Sub 和点对点队列混在一起用。同样是发消息语义差别其实非常大。对比维度点对点队列模式发布/订阅模式请求/响应模式核心语义一条消息只有一个消费者消费一条消息可被多个订阅者消费一个请求对应一个响应耦合关系双方通过同一个队列间接关联双方只通过同一主题间接关联双方需要显式寻址消息生命周期被消费后通常移除按照订阅关系分别派发请求响应即时完成典型场景任务分发、轮询处理作业事件广播、系统间协同通知接口调用、数据库查询这里面的核心差异是“广播”与“竞争”。点对点队列强调一条消息只能被一个消费者拿走适合任务分配发布订阅模式则允许同一个事件被多个订阅方同时拿到适合事件驱动。比如订单创建后库存服务要扣减库存积分服务要增加积分消息通知服务要发短信三者都需要这个订单事件但又互不干扰。如果只用点对点队列就得把同一个消息复制三份消息中间件通常不鼓励这种用法。使用 Pub/Sub每个订阅组都能独立消费同一份消息逻辑才顺畅。请求/响应模式更不用多说它是同步的、有返回值的适合需要立刻知道结果的操作。Pub/Sub 是异步的注重“事件已经发生”的事实广播谁处理、处理结果如何发布者一般不关心。这两者最好不要混用一旦把一次 RPC 调用改成 Pub/Sub调用方就失去了拿返回值的能力很多事务逻辑会变得很难做。2. 核心设计发布方、消息通道、订阅方如何配合2.1 三个角色和主题机制一个标准的 Pub/Sub 系统里有三个角色发布者、消息通道、订阅者。消息通道是中间层它可以是内存事件总线的内部调度器也可以是一台分布式的消息中间件集群。既然是发布/订阅双方关心的不是对方而是一个叫“主题”的逻辑名字。主题是消息分类的基本维度。发布者发送一条订单支付成功消息就往 order_paid 这个主题上发发送一条用户注册消息就往 user_registered 上发。订阅者想接收哪类消息就订阅对应的主题。这里最需要注意的是主题的设计对未来扩展影响巨大。我给团队的建议是主题命名要带业务语义最好约定成“业务域.事件名称”的形式例如trade.order.paid、marketing.coupon.issued。命名太随意的话后面建了很多主题之后根本没法认消息中间件的管理后台会变得一团乱。订阅关系也值得专门讲一下。同一个主题下面订阅者往往不是一个实例而是一组实例。大多数消息中间件会把这一组实例抽象成“消费组”。同一个消费组里的多个消费者实例分摊主题消息也就是说一条消息只会被组内某一个实例消费不同消费组则会各自独立收到这条消息。这就是第1节里提到的广播与竞争的组合关系组内竞争组间广播。2.2 推模型与拉模型的选择消息通道怎么把消息交给订阅者业界基本有两种思路推和拉。推模型是消息中间件主动把消息送到订阅者进程里订阅者不需要轮询消息一到就触发回调。这种模式实时性高实现起来也直白但有个明显问题当生产速度远大于消费速度时中间件疯狂向消费者推送消费者可能直接被流量打挂。拉模型反过来消费者自己主动去消息中间件取消息取多少、什么时候取都由自己控制。常见的消费流程是消费者发一个拉取请求指定主题、消费组和一次最多拿多少条中间件拉出一批消息返回消费者处理完以后提交确认再拉下一批。这种模式天然带流量控制消费者处理不过来时可以少拉或者暂停拉取不容易被压垮。选推还是选拉没有绝对的好坏。轻量内存总线通常用推因为在一个进程内部推送开销可以忽略实时性还更好。分布式的消息中间件大多偏向拉模型或者至少支持消费者客户端主动拉取这是为了在大规模场景下保护下游。我实际使用中更推荐拉模型因为消费者可以自己做批量处理、自主限流线上出问题的时候处理速度至少是我们自己在掌控。2.3 消息过滤和路由细节主题这个维度有时候还不够细。比如同一类业务事件主题下面订阅者想只收其中一部分消息。常见做法有几类一类是主题通配符订阅者可以订阅一个主题模式而不是固定主题例如订阅trade.*.paid就能收到所有交易域下以 paid 结尾的事件另一类是消息标签发布者在消息头里带上标签订阅者在订阅时指定只接收某种标签的消息。这两种方案各有取舍。主题通配符灵活但会让路由逻辑复杂化尤其是在消息中间件底层做匹配时模式匹配引擎要额外计算主题数量大了之后性能会受影响。消息标签的实现更简单直观但标签粒度一旦定死扩展就受限制。实际项目里我一般建议在发布端和服务端维护一份主题清单尽量不要过度依赖通配符。真正需要做动态路由的时候可以用消息里的业务字段在订阅端自行过滤这比让中间件去做复杂匹配要可控得多。还有一个容易被忽略的细节是路由的确定性。有些业务要求事件必须按某种顺序处理比如同一个用户的操作日志不能被后面的操作提前消费。这时只靠主题没法解决通常会在消息里带一个分区键中间件根据这个键做哈希把同一个键的消息路由到同一个分区。这样同一个用户的事件就能保证进入同一个分区消费者按顺序处理。很多消息中间件内部就是靠这个机制保证局部顺序的。3. 实操落地从零搭一个模拟 Pub/Sub 系统3.1 选型判断内存总线还是独立消息中间件动手之前一定要先想清楚究竟是只需要一个进程内的事件总线还是需要一套独立的分布式消息中间件。这个选择做错了后面会付出很大的维护成本。如果项目是单体应用或者一个应用内部多个模块之间需要解耦比如订单模块要通知库存模块、积分模块用内存事件总线就够了。它轻、快、没有额外运维负担出问题也好调试。如果应用已经拆成了多个服务部署在不同机器上那内存总线天然过不去必须有一个网络层的消息通道这时才需要引入独立的消息中间件。独立消息中间件带来的不只是网络通信能力还有持久化、堆积、消费组管理、重试机制等能力。但代价也不小需要单独部署集群、监控消费积压、处理消息积压后的问题。我见过不少项目明明只有两个服务要做异步通知也硬上一个中间件集群最后复杂程度完全失控。架构选型不是越重越好能用简单方案解决的问题就不要先引入重型组件。3.2 进程内的最小可用事件总线这一节我先给一个不依赖任何中间件的最小实现。它不能替代真正的消息中间件但能把 Pub/Sub 这个模式的核心流程跑出来也适合理解机制。我用 Python 写一个非常简单的内存版发布订阅总线from typing import Callable, Dict, List class EventBus: def __init__(self) - None: self._subscribers: Dict[str, List[Callable]] {} def subscribe(self, topic: str, callback: Callable) - None: if topic not in self._subscribers: self._subscribers[topic] [] self._subscribers[topic].append(callback) def publish(self, topic: str, payload: object) - None: for callback in self._subscribers.get(topic, []): callback(payload) # 使用示例 bus EventBus() bus.subscribe(trade.order.paid, lambda event: print(库存扣减, event)) bus.subscribe(trade.order.paid, lambda event: print(积分增加, event)) bus.publish(trade.order.paid, {order_id: 10086, user_id: 123})这段代码的核心逻辑很短维护一个从主题到回调列表的字典发布时遍历回调逐个执行。一次发布多个订阅者都能收到这就是 Pub/Sub 最基本的语义。但这种实现有很多局限。发布和订阅在同一个线程里回调阻塞时会拖慢发布者没有消息持久化进程一重启就什么都没了没有消费确认回调抛异常了消息也就丢了。这些局限正是后面要上独立消息中间件的原因。不过作为理解模式的脚手架它非常合适。3.3 引入消息中间件后的完整链路换成独立消息中间件之后同样的流程就变成了四个步骤。我以常见分布式消息引擎为例写一套伪代码风格的接入步骤。第一步是初始化一个客户端连接。发布端和订阅端各自连到同一个消息服务地址上连接参数包括服务地址、连接超时、重试次数等。对于生产环境我建议把超时时间调大一些连接异常时至少重试两到三次避免网络抖动导致发布失败。第二步是创建主题或声明主题。多数消息中间件支持自动创建主题但生产环境我强烈建议关闭自动创建改为在管理端人工审批。原因很简单自动创建容易让那些随意拼写的主题混进系统后期维护成本极高。第三步是消费者客户端订阅主题指定消费组和消息处理回调。这里有一个高频参数需要关注是从最新的消息开始消费还是从最早的位置开始消费。默认值通常是“从最新开始”但如果消费者进程挂了很久重启后可能漏掉这期间的很多消息。业务允许补数据的时候订阅时最好显式设置为从最早开始。第四步是消息确认。拉取模式下的流程是拉消息、执行业务逻辑、提交位移。提交位移再细分为自动和手动。自动提交省事但存在风险消息被拉回来后进程还没处理完就崩溃重启后这条消费位置可能已经被提交了消息就永久丢掉了。我习惯用手动提交处理成功后先记录业务本身的状态再提交消费位移这样至少在业务层不会丢关键数据。3.4 一个模拟订单事件的项目实践某次我在做一个订单事件同步的模拟项目场景不复杂订单服务每完成一笔订单就发布一个order.completed事件库存服务、积分服务、短信服务都要对这个事件做出反应。我先在消息服务里创建了主题order.completed然后让三个服务分别注册了三个不同的消费组。这里必须强调三个服务不能共用同一个消费组否则库存、积分、短信会互相竞争同一笔事件只会被其中一个服务处理掉。正确做法是三个消费组独立订阅消息中间件复制出三份分别投递。发布端代码核心如下只负责发送事件不关心下游谁处理def complete_order(order_id: str, user_id: int) - None: # 业务逻辑更新订单状态 event { event_type: order.completed, order_id: order_id, user_id: user_id, occurred_at: int(time.time()), } producer.send(order.completed, event)消费端代码则集中在处理逻辑和提交确认上def handle_order_completed(message): event message.event if event[event_type] ! order.completed: message.nack() return # 各自的业务处理扣库存、加积分、发短信 process(event) # 处理成功后手动提交 message.ack()整个流程跑下来印象最深的是消费组语义。只要消费组设置正确一个发布事件就能顺畅地被多个服务独立消费而各服务之间完全不需要感知对方存在。这也是解耦在实际系统里最直观的体现。4. 事中排查Pub/Sub 落地最常见的四个坑4.1 订阅者收不到消息初用消息中间件的人最容易遇到这个问题。订阅方已经启动日志里也显示成功订阅了主题但发布方发消息以后订阅方就是收不到。排查这个问题的顺序我来回用了几十次之后已经固定下来。第一步检查主题是不是同一个。发布方写的是order.completed订阅方写的是order.Completed大小写一不一致中间件不会报错因为很多中间件主题名就是普通字符串但这显然导致一条消息发到了两个完全不同的地方。第二步检查消费组和订阅关系是否建立正确。有些中间件允许同一个消费者服务注册多个订阅订阅关系和消费组之间如果搞混也会出现“订阅成功但收不到”的现象。第三步看消费者启动位置也就是消费起始位移。消费者启动时默认“从最新开始”而发布方是在它启动之前发的消息那自然一条都不消费。第四步看消费者所在集群有没有和消息中间件的网络隔离跨网络区域经常出现白名单没开的情况。这四个步骤里最容易踩的是第三步也是最容易被人忽略的。线上排查时我一般先看消费组管理页面的堆积量如果堆积量一直是 0说明消息根本没到消费者所在的消费组或者消费者没有正确建立订阅关系如果堆积量在涨那就是消费者处理速度太慢跟“收不到”是两码事。4.2 消息重复与幂等消息中间件为了可靠投递通常会采用“至少一次”的投递语义。这意味着正常情况下消息不丢但可能会在某些异常时刻重复投递。比如消费者处理完消息之后还没来得及提交消费位移进程就崩了重启后中间件就会把这条消息重新投递一遍。这是可靠性和重复之间的必然权衡所有分布式消息系统都绕不开。应对方案只有一个核心原则消费逻辑必须幂等。幂等不是说让中间件保证消息不重复而是让同一条消息被重复处理多次时最终产生的业务效果与只处理一次相同。例如扣减库存的操作不要设计成无条件的库存减一而是先判断这个订单号对应的扣减记录是否已经存在如果存在就直接跳过或者用消息里的唯一事件 ID 建一张去重表已经处理过的 ID 不再处理第二次。我之前遇到过一个实际事故短信服务收到重复事件后给用户连发了两条相同短信。后来在消费侧增加了一个事件 ID 缓存处理前先查缓存判断虽然极端情况下还是有窗口期但至少把重复概率压到了很低。注意不要指望消息中间件的配置能完全抹掉重复稳妥的做法永远是在自己的业务逻辑里做兜底。4.3 消息堆积和消费延迟堆积是 Pub/Sub 系统里最常见也最难缠的问题。生产者的写入速度远大于消费者的处理速度时消息就会在中间件上越积越多消费延迟越来越大。排查这类问题第一件事不是调参数而是看消费端究竟慢在哪里。常见原因有三个。第一个是消费回调里做了太多重操作比如同步调数据库、调用外部接口等了很久第二个是单个消费组只有一个消费者实例消息只能串行处理处理能力天然受限第三个是批量拉取的消息条数设置得太小消费者每次只能拉很少的数据导致频繁发起网络请求吞吐上不去。定位之后再做针对性处理。如果消费逻辑本身很重就要考虑拆分流程把发放短信、推送通知这类耗时操作继续丢给更下游的组件如果单个消费者实例太少就调整消费组内的实例数量让消息被分摊处理如果批量参数有问题就适当放宽单次拉取的条数和长轮询等待时间。每次调整参数后要持续观察消费堆积量有没有下降不能只凭感觉调。4.4 顺序性无法满足业务要求在某些场景下消息顺序非常重要例如同一个订单的状态变化不能乱必须是待支付、已支付、已发货这个顺序。消息中间件默认是无法保证全局顺序的多个消费者并行消费时先后发布的消息可能被不同的消费者实例同时处理顺序自然就无法保证。解决思路通常是分区有序。在向主题发送消息时把订单 ID 作为分区键中间件对分区键做哈希同一个订单 ID 的消息始终进入同一个分区。消费者在同一方向上尽量只让一个实例处理同一个分区分区内就保证了顺序。需要注意分区有序不等于全局有序。如果业务要求全主题所有消息都严格按顺序处理那就意味着只能有一个消费者实例串行消费吞吐量会很低。这种需求本身就值得怀疑绝大多数业务其实只需要按某个业务键局部有序。设计阶段就要想清楚是哪种有序否则后期再改消息模型非常痛苦。下面把常见的排查点整理成速查表现象可能原因排查/处理建议订阅者收不到消息主题名不一致、消费起始位移不对、订阅关系错误先核对主题再查消费组堆积量消息频繁重复处理消费后未及时提交位移导致重新投递消费侧做幂等按事件 ID 去重消息堆积不降消费逻辑重、实例数少、拉取批量太小先定位消费耗时再扩容实例或调参消息顺序错乱多个消费者并行消费、分区键缺失设置业务分区键按分区消费5. 选型思考项目中该不该上真正的消息中间件5.1 内存事件总线能解决的问题在给项目选型之前先确认内存事件总线是不是已经够用。如果所有组件都部署在同一个进程内服务之间不需要跨网络通信那完全没必要引入独立中间件。内存总线的好处很明显零网络开销、零序列化开销、调试起来直观。断点一打发布端和订阅端都在同一个进程里问题定位非常快。典型场景是单体应用内部的业务解耦。比如下单接口完成后需要记录操作日志、发送站内信、更新统计指标这些动作都可以通过事件发布出去让对应的处理器异步执行。哪怕这个异步只是换一个线程去处理对接口响应时间的改善也是立竿见影的。内存总线也有明显的天花板。重启丢消息、不支持多语言跨进程通信、消费者必须和生产者同生命周期这些限制决定了它只能作为模块内部的协作工具一旦涉及多个独立部署的服务就必须升级方案。5.2 出现这些信号就果断上独立中间件什么信号出现时该上独立消息中间件我盘点几个实际经历过的判断标准。跨服务事件广播是最直接的信号。订单服务想通知库存、积分、通知三个不同服务这几个服务部署在不同机器上内存总线做不到只能通过网络发消息。第二个信号是消息需要持久化。业务要求消费者宕机一段时间后还能恢复消费消息不能丢这时就需要中间件保存消息副本。第三个信号是削峰填谷。系统要应对瞬时高流量下游处理能力不足必须靠中间件缓冲流量。还有一个容易被忽视的信号需要消费组管理。如果你希望同一条主题消息被多个消费组各拿一份还要分别记录每个组消费到哪了那必须靠独立中间件的消费位点管理机制这已经是分布式协作的范畴了。5.3 不适合用 Pub/Sub 的场景不是所有异步都适合用 Pub/Sub。强一致性的业务不要用。比如转账场景必须同步确认扣款和入账都成功才能给用户反馈如果用发布订阅模式把扣款动作广播出去用户收到“转账成功”时钱可能还没真正到账这个体验和账务风险都很难接受。这类操作应该走事务型接口确保实时结果。实时性要求极高的指令操作也不适合。如果一条指令必须在几十毫秒内得到响应Pub/Sub 的异步处理和排队过程反而会拖慢整体时延。调度型场景更适合点对点队列比如把任务分发到空闲机器上执行而不是让每台机器都收到任务。我对 Pub/Sub 的整体体会是它是一个强大的消息传递模式但强在“解耦”和“分发”并不擅长“强一致”和“强实时”。用之前先把业务目标想清楚你会发现很多问题其实根本不该用 Pub/Sub 解决。这篇文章就到这里。最后再分享一个我自己的小经验每次引入新的消息订阅不要只写代码先给消息定义一个完整的事件结构把事件 ID、业务域、发生时间都放进去。这个习惯救了我很多次排查重复消息和顺序错乱的时候光靠这个结构就能省掉一大半的时间。