观察者模式:Spring 事件驱动架构的基石
一句话结论(30s)
观察者模式的本质是让「状态变化方」只负责通知、不关心谁在听,因为主题只调用观察者接口方法,双方互不知晓对方内部实现,从而把硬编码调用链解耦成可插拔的监听网络。权衡:通知顺序与失败处理需额外设计、同步通知会放大事务边界,但新增通知渠道只需加一个监听方法、主流程零修改。Spring 的 ApplicationEventPublisher + @EventListener 就是标准实现。
核心原理(2min)
主流程:主题状态变化 → 发布事件 → 遍历通知所有已注册观察者 → 观察者各自响应。关键机制在 Spring 三件套:ApplicationEventPublisher.publishEvent() 发布事件、@EventListener 注册监听方法,二者默认同线程同事务同步执行;@TransactionalEventListener(phase = AFTER_COMMIT) 把监听绑定到事务提交/回滚边界,解决「发短信失败导致回滚库存」的事务放大问题;@Async 让监听异步执行、主线程立即返回,把响应时间从「扣库存+发短信+加积分」缩短为「仅扣库存」。源码落地:下单后扣库存、发短信、加积分、核销优惠券等横切动作各自独立成监听器,新增渠道零改主流程。
底层深入(5-10min)
一个下订单后的”通知链条”
用户下了一个订单,需要做多少事?
下单成功 → 扣减库存 → 发送短信 → 积分累计 → 优惠券核销 → 打标签 → 数据埋点
如果全部写在一个方法里:
public void placeOrder(Order order) {
orderDao.insert(order); // 1. 存订单
inventoryService.deduct(order); // 2. 扣库存
smsService.send(order); // 3. 发短信
pointService.add(order); // 4. 加积分
couponService.use(order); // 5. 核销优惠券
tagService.mark(order); // 6. 用户打标
analyticsService.track(order); // 7. 数据埋点
}
这段代码有三个致命问题:
- 强耦合:
placeOrder()需要知道所有后续操作的存在。新增一个”发 Push 通知” → 必须改这段代码。 - 事务放大:扣库存和发短信在同一个事务里——发短信失败导致回滚库存?不合理。
- 响应时间膨胀:等 7 个操作全部完成才返回用户 → 用户盯着”加载中”圈圈,体验极差。
思考穿插:为什么「扣库存 + 发短信」放在一个事务里是错的? 事务的边界应该包住「必须同生共死」的数据操作,而不是包住所有副作用。发短信、加积分这类动作失败了不该回滚库存——库存是核心数据,短信是外围通知,把外围失败放大成核心回滚,就是「事务边界放大」。这也引出后面 @TransactionalEventListener(AFTER_COMMIT) 的解法:等事务确定提交了,再去做这些外围动作。
观察者模式的结构
观察者模式(Observer Pattern)解决的就是”一个对象状态变化时,自动通知所有依赖它的对象,且双方互不知晓彼此的存在。”
Subject(主题/被观察者)
↓ notify
Observer1 Observer2 Observer3 (观察者们)
核心原则:主题只负责”通知事件发生了”,不关心谁在听、听者做什么。观察者自己决定如何响应。
手写版实现
// 观察者接口
public interface OrderListener {
void onOrderPlaced(OrderEvent event);
}
// 主题(被观察者)
public class OrderService {
private List<OrderListener> listeners = new ArrayList<>();
public void register(OrderListener listener) {
listeners.add(listener);
}
public void placeOrder(Order order) {
// 核心逻辑
OrderEvent event = new OrderEvent(order);
// 通知所有观察者
for (OrderListener listener : listeners) {
listener.onOrderPlaced(event);
}
}
}
Spring 的观察者模式实现
Spring 内置了一套完整的事件发布/订阅机制,本质就是观察者模式:
定义事件
// 继承 ApplicationEvent
public class OrderPlacedEvent extends ApplicationEvent {
private final Order order;
public OrderPlacedEvent(Object source, Order order) {
super(source);
this.order = order;
}
// getter...
}
发布事件
@Service
public class OrderService {
@Autowired
private ApplicationEventPublisher publisher;
@Transactional
public void placeOrder(Order order) {
orderDao.insert(order);
// 发布事件(这行代码瞬间返回,不等待监听器执行完)
publisher.publishEvent(new OrderPlacedEvent(this, order));
}
}
监听事件
@Component
public class SmsListener {
@EventListener
public void handleOrderPlaced(OrderPlacedEvent event) {
smsService.send(event.getOrder().getPhone(), "下单成功!");
}
}
@Component
public class PointListener {
@EventListener
public void handleOrderPlaced(OrderPlacedEvent event) {
pointService.addPoints(event.getOrder().getUserId(), 10);
}
}
新增一个通知渠道?加一个 @EventListener 方法即可,OrderService 一行不改。 这就是松耦合。
真实源码:multicastEvent 如何遍历通知
ApplicationEventPublisher.publishEvent() 最终会落到 SimpleApplicationEventMulticaster.multicastEvent()——它就是观察者模式里那个「遍历所有观察者逐个通知」的 Subject 核心。以下是 Spring Framework 真实源码:
@Override
public void multicastEvent(ApplicationEvent event, @Nullable ResolvableType eventType) {
ResolvableType type = (eventType != null ? eventType : ResolvableType.forInstance(event));
Executor executor = getTaskExecutor();
for (ApplicationListener<?> listener : getApplicationListeners(event, type)) {
if (executor != null && listener.supportsAsyncExecution()) {
try {
executor.execute(() -> invokeListener(listener, event));
}
catch (RejectedExecutionException ex) {
// Probably on shutdown -> invoke listener locally instead
invokeListener(listener, event);
}
}
else {
invokeListener(listener, event);
}
}
}
而 invokeListener 层层透传,最终在 doInvokeListener 里回调观察者自己的 onApplicationEvent:
@SuppressWarnings({"rawtypes", "unchecked"})
private void doInvokeListener(ApplicationListener listener, ApplicationEvent event) {
try {
listener.onApplicationEvent(event);
}
catch (ClassCastException ex) {
String msg = ex.getMessage();
if (msg == null || matchesClassCastMessage(msg, event.getClass()) ||
(event instanceof PayloadApplicationEvent payloadEvent &&
matchesClassCastMessage(msg, payloadEvent.getPayload().getClass()))) {
// Possibly a lambda-defined listener which we could not resolve the generic event type for
// -> let's suppress the exception.
Log loggerToUse = this.lazyLogger;
if (loggerToUse == null) {
loggerToUse = LogFactory.getLog(getClass());
this.lazyLogger = loggerToUse;
}
if (loggerToUse.isTraceEnabled()) {
loggerToUse.trace("Non-matching event type for listener: " + listener, ex);
}
}
else {
throw ex;
}
}
}
这段源码印证了观察者模式的三个关键点:getApplicationListeners(event, type) 先根据事件类型筛出所有匹配的监听器,for 循环再逐个回调 listener.onApplicationEvent(event),而主题本身完全不关心监听器内部做了什么。executor != null && listener.supportsAsyncExecution() 这个分支正是前面 @Async 异步执行的底层来源——配了线程池且监听器声明支持异步时就丢给线程池,否则就在调用线程里同步执行,和文章前面「默认同线程同步执行」的结论完全吻合。
@TransactionalEventListener:事务提交后再通知
上面的代码有个隐蔽的坑:@EventListener 默认在发布事件的同一个线程、同一个事务中同步执行。如果发短信抛异常 → 事务回滚 → 库存已经扣了但订单没存上。
Spring 4.2+ 提供了 @TransactionalEventListener:
@Component
public class SmsListener {
@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
public void handleOrderPlaced(OrderPlacedEvent event) {
// 只在事务成功提交后才执行
smsService.send(event.getOrder().getPhone(), "下单成功!");
}
}
phase 有四个选项:
| Phase | 时机 |
|---|---|
| BEFORE_COMMIT | 事务提交前 |
| AFTER_COMMIT | 事务提交后(默认) |
| AFTER_ROLLBACK | 事务回滚后 |
| AFTER_COMPLETION | 事务完成(提交或回滚)后 |
异步事件:不阻塞主线程
@Component
@EnableAsync
public class PointListener {
@Async // ← 在新线程中异步执行
@EventListener
public void handle(OrderPlacedEvent event) {
pointService.add(event.getOrder().getUserId(), 10);
}
}
主线程 placeOrder() 发布事件后立即返回——监听器的执行完全异步。用户的响应时间从”扣库存 + 发短信 + 加积分”缩短为”仅扣库存”。
观察者模式 vs 发布-订阅模式
很多人把它们混为一谈,但有细微区别:
| 观察者模式 | 发布-订阅 | |
|---|---|---|
| 耦合程度 | 观察者知道 Subject 存在 | 发布者和订阅者通过 Broker 通信,完全互不知晓 |
| 通信媒介 | 直接调用 | 消息队列 / EventBus |
| 典型实现 | @EventListener | Kafka、RabbitMQ、Redis Pub/Sub |
Spring 的 ApplicationEventPublisher 更接近观察者模式(在同一个 JVM 内直接方法调用),而 RocketMQ / Kafka 实现的是发布-订阅(跨进程、跨服务)。
思考穿插:观察者和发布-订阅到底差在哪? 一句话看「中间有没有 Broker、双方认不认识」:观察者是「观察者注册到主题上,主题直接遍历通知」,双方还是互相持有一层引用;发布-订阅是「发布者和订阅者都只认识 Broker,中间靠消息队列中转」,彼此完全不知道对方存在。所以 Spring 事件是观察者(同 JVM 直接调用),Kafka/RabbitMQ 是发布-订阅(跨进程、跨服务)。
思考穿插:@Async 异步监听有没有坑? 有——异步意味着脱离了事务和调用线程,事务内的数据此刻未必可见,且异常不会回滚主流程、只能靠日志兜底。所以异步监听里不能假设「数据库里已经有这条数据」,要幂等 + 失败重试,否则消息丢了或重复了都难查。
总结
观察者模式让代码从”硬编码调用链”变成”发布-订阅松耦合网络”。Spring 中三件套 ApplicationEvent + @EventListener + ApplicationEventPublisher 是最常用的实现,再配合 @TransactionalEventListener 绑定事务边界、@Async 异步执行,就能构建出灵活、可扩展的事件驱动架构。
章末提问
Q1:观察者模式和发布-订阅模式的区别? 结论先行:看「中间有没有 Broker、双方认不认识」——观察者由主题直接遍历通知、双方互相持有一层引用;发布-订阅靠消息队列中转、双方完全互不知晓。因为观察者要注册到 Subject 上、Subject 要遍历调用观察者接口,而同 JVM 的 Spring 事件正是如此;Kafka/RabbitMQ 则把耦合全部收敛到 Broker,跨进程解耦。
Q2:@EventListener 和 @TransactionalEventListener 的区别?
结论先行:前者默认同线程、同事务、同步执行,后者把监听绑定到事务提交/回滚边界(默认 AFTER_COMMIT)。因为同步监听里一旦发短信抛异常,会让整个事务回滚、库存白扣;绑定到 AFTER_COMMIT 后,监听只在事务确定提交后才触发,外围副作用不会反过来回滚核心事务。
Q3:multicastEvent 里 supportsAsyncExecution() 分支在做什么?
结论先行:它决定了监听器是丢线程池异步执行、还是在调用线程里同步执行。因为配了 taskExecutor 且监听器声明支持异步时,才 executor.execute(...) 异步跑;否则直接 invokeListener 同步回调——这正是 @Async 的底层来源。
Q4:异步监听有什么隐患? 结论先行:脱离了事务和调用线程,数据可见性与异常处理都要重做。因为异步线程里事务未必已提交、数据未必可见,且异常不会回滚主流程;所以异步监听必须幂等、失败要重试或落到兜底,否则消息会丢或重复。
Q5:新增一个通知渠道怎么做到主流程零修改?
结论先行:加一个 @EventListener 方法即可,因为主题只依赖观察者接口、不关心具体监听器。发布方只 publishEvent,新增的监听器自己注册自己,placeOrder() 一行不用改——这正是把「硬编码调用链」解耦成「可插拔监听网络」的价值。