Skip to content

一次状态变化怎样通知多方:观察者、事件与消息的边界 ​

「发一个事件,谁关心谁订阅」听起来让发布者和订阅者彻底解耦。但每把事件往外推一层:从直接调用到同步监听器,到异步监听器,到提交后监听器,再到跨进程的消息,就要多回答一个关于失败的问题。

报名成功后要做四件事:记积分、写审计、发邮件、通知组织者。用 Spring 的事件机制实现,四个订阅者分别是两个同步监听器、一个 @Async 监听器和一个 @TransactionalEventListener。实测里有两个结果和直觉不一样:

text
异步线程池(1 个线程、队列 2)连续收到 10 个事件:
  3 个被接受,7 个在发布者线程抛出 TaskRejectedException
  10 条报名全部已经保存 → 7 条报名永远不会有邮件

提交后监听器(通知组织者)抛出异常:
  报名已经提交;发布者没有收到任何异常
  日志里只有一行 SEVERE: TransactionSynchronization.afterCompletion threw exception

本文用 Spring Framework 7.0.9 的实测,把事件机制拆成四个维度,说明每一层解决了什么、又引入了什么。

一、先说结论 ​

  • 先拆开四个维度:调用关系(直接调用或发布订阅)、线程(同步或异步)、进程(进程内或跨进程)、事务(事务内、提交后或 Outbox)。「用事件」只回答了第一个。
  • 同步监听器和发布者同生共死。 实测审计监听器抛出异常后,发布者收到异常,报名回滚。这有时正是想要的:审计写不进去就不允许报名。
  • 异步不会让工作消失,只是换了个地方失败。 实测 @Async 监听器的异常不会回到发布者;线程池满时,拒绝发生在发布者线程上,数据已经保存而事件丢了。
  • 提交后监听器解决了「回滚了却发了通知」,但不保证一定执行。 实测没有事务时它不执行,执行失败时异常被吞掉;进程在提交后崩溃,它也不会执行。
  • 事件不能丢时,把它和数据写进同一个事务。 这就是 Outbox:由投递器读取并重试,订阅方做幂等,见 Kafka 不丢、不重与 Exactly Once。

二、观察者的四个维度 ​

同步监听器同一线程、同一事务实测:抛异常 → 报名回滚@Async 监听器线程池执行实测:异常不回发布者提交后监听器AFTER_COMMIT实测:无事务不执行,异常被吞Outbox + 消息跨进程与业务同一事务写入要回答:订阅者失败该不该让发布者失败?要回答:排队、背压、失败后谁来补?要回答:提交后进程崩溃或监听器失败怎么办?要回答:重复、乱序、schema 演进?从左到右,发布者与订阅者的耦合越来越松,需要处理的失败越来越多只有一个订阅者、而且发布者本来就知道它时,直接调用最清楚
图 1 · 同步监听器与发布者同生共死;异步监听器把失败和排队挪到线程池里;提交后监听器解决回滚后的副作用,却不保证一定执行;跨进程时还要加上持久化与投递语义
维度选项必须回答的问题
调用关系直接调用 / 发布订阅发布者真的不该知道订阅者吗?只有一个订阅者时,直接调用更清楚
线程同步 / 异步订阅者的耗时和失败要不要算在发布者头上?排队满了怎么办?
进程进程内 / 跨进程进程崩溃时事件还在吗?重复、乱序怎么处理?
事务事务内 / 提交后 / Outbox业务回滚时副作用要不要撤回?提交后失败谁来补?

经典的观察者模式只回答了第一个维度。JDK 早期的 java.util.Observable 与 Observer 自 JDK 9 起已被标记为废弃,官方建议改用 java.beans 包的事件模型或并发包里的类;业务系统里更常见的是框架的事件机制与消息中间件。

三、同步监听器:顺序、线程与异常 ​

实测一次报名(没有失败时)的调用轨迹:

text
points@publisher  audit@publisher  email@async-1  afterCommit@publisher
  • 两个同步监听器按 @Order 顺序在发布者线程执行,它们的耗时直接加在报名请求上;
  • 审计监听器抛出异常时,发布者收到「审计写入失败」,报名回滚(已提交 0 条),提交后监听器也不再执行。

这正是同步监听器的语义:订阅者属于同一个工作单元。适合「订阅者失败就不允许继续」的场景;如果订阅者失败不应该影响主流程,就不要用同步监听器,或者在监听器里自己捕获异常并记录。

四、异步监听器:失败与积压换了地方 ​

加上 @Async 后,发布者在邮件监听器完成之前就返回了。实测两件事:

异常不会回到发布者。 邮件监听器抛出的异常只进入 AsyncUncaughtExceptionHandler。如果这个处理器只打一行日志,邮件就静悄悄地没发出去。

线程池满时,拒绝发生在发布者线程。 线程池配置为 1 个线程、队列 2,监听器每个耗时 500ms,连续发布 10 个事件:

保存报名10 条全部成功发布事件10 次线程池接受 3 个1 个执行 + 队列 2拒绝 7 个TaskRejectedException3 条有通知7 条没有通知数据已提交修复方向:事件与数据写进同一个事务(outbox),由投递器重试;或者让发布者感知拒绝并决定是否回滚不要:在无界队列上异步发布——拒绝消失了,积压转移到内存里
图 2 · 线程 1、队列 2,连续发布 10 个事件:3 个被接受,7 个在发布者线程抛出 TaskRejectedException;报名已经保存了 10 条,其中 7 条永远不会有对应的通知

3 个被接受,7 个在发布者线程抛出 TaskRejectedException;而发布发生在保存之后,10 条报名都已经保存。数据和事件在这里分叉了。换成无界队列,拒绝消失了,积压会转移到内存里,进程重启时一起丢失。

异步的价值在于缩短响应路径,代价是必须回答:排队的上限是多少、满了怎么办、失败了谁来补、重启时队列里的事件怎么办。进程内的异步机制本身都不持久化,Guava EventBus 也一样,见 Guava EventBus。

还有一个实现细节:带 @Async 方法的 Bean 在容器里是 CGLIB 代理。实测直接读代理对象的字段得到 null,状态只能通过方法访问,原因见 Spring AOP 为什么会失效。

五、提交后监听器:只保证「回滚时不执行」 ​

@TransactionalEventListener(phase = AFTER_COMMIT) 让监听器在事务提交之后才执行,解决了「数据回滚了,通知却发出去了」的问题。实测它还有三个边界:

情况结果
业务回滚不执行
发布时没有事务不执行(默认 fallbackExecution = false)
监听器抛出异常报名已提交;发布者没有收到异常,日志中一行 SEVERE
提交后、监听器执行前进程崩溃不执行(未实测,由机制决定:它只存在于内存中)

第二行最容易踩:同一个服务方法有时在事务里调用、有时不在,通知就时有时无。第三、四行说明它不是可靠投递机制:提交已经完成,失败无法回滚,也没有人重试。

六、事件不能丢:Outbox ​

当事件的丢失不可接受(给用户的通知、给下游系统的状态变更),就不能依赖进程内的任何机制。做法是把事件当作数据,和业务写进同一个事务:

sql
BEGIN;
INSERT INTO registrations (...) VALUES (...);
INSERT INTO outbox (event_id, type, payload, created_at) VALUES (...);
COMMIT;

独立的投递器读取 outbox,发给消息中间件或直接调用下游,成功后标记。它带来三件需要处理的事:

  • 至少一次:投递器可能在发送成功、标记之前崩溃,同一个事件会被再发一次,订阅方必须幂等;
  • 顺序:同一个聚合的事件按顺序投递,不同聚合之间通常不需要;
  • 事件是契约:发布者不知道订阅者,不代表可以随意修改事件结构,字段的增删要考虑已部署的订阅方。

状态机里「状态改了、通知失败」的对比实测,见 把状态分支变成可验证模型。

七、怎么选 ​

场景选择
只有一个订阅者,发布者本来就知道它直接调用
订阅者失败应该让主流程失败同步监听器
订阅者慢、失败可以容忍,丢了也能接受异步监听器,设置有界队列并处理拒绝
只在提交成功后执行,偶尔丢失可以接受提交后监听器
不能丢、要重试、要跨进程Outbox + 投递器 + 幂等订阅

八、常见误区 ​

  • 「用了事件就解耦了」:事件的结构、发布时机、失败语义都是契约,只是从编译期依赖变成了运行期约定。
  • 「异步更快」:它缩短的是响应路径,工作量和失败没有减少,实测满了会在发布者线程抛出拒绝。
  • 「@TransactionalEventListener 保证提交后一定执行」:没有事务时不执行,失败时异常被吞,崩溃时丢失。
  • 「进程内事件总线可以替代消息队列」:它没有持久化、重试和消费进度,重启即丢。

小结 ​

从直接调用到 Outbox,每走一步,发布者与订阅者的耦合更松,需要自己处理的失败也更多。选之前先回答四个问题:订阅者失败是否影响主流程、能否接受排队与丢失、是否需要在提交后执行、是否需要跨进程可靠投递。答案决定了用哪一层,而不是「用事件更优雅」。


配套实验

参考资料

文章以 CC BY-NC-SA 4.0 授权 · 代码片段以 MIT 授权