Skip to content
Go back

观察者模式——Spring事件驱动架构的基石

观察者模式: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. 数据埋点
}

这段代码有三个致命问题:

  1. 强耦合placeOrder() 需要知道所有后续操作的存在。新增一个”发 Push 通知” → 必须改这段代码。
  2. 事务放大:扣库存和发短信在同一个事务里——发短信失败导致回滚库存?不合理。
  3. 响应时间膨胀:等 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
典型实现@EventListenerKafka、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:multicastEventsupportsAsyncExecution() 分支在做什么? 结论先行:它决定了监听器是丢线程池异步执行、还是在调用线程里同步执行。因为配了 taskExecutor 且监听器声明支持异步时,才 executor.execute(...) 异步跑;否则直接 invokeListener 同步回调——这正是 @Async 的底层来源。

Q4:异步监听有什么隐患? 结论先行:脱离了事务和调用线程,数据可见性与异常处理都要重做。因为异步线程里事务未必已提交、数据未必可见,且异常不会回滚主流程;所以异步监听必须幂等、失败要重试或落到兜底,否则消息会丢或重复。

Q5:新增一个通知渠道怎么做到主流程零修改? 结论先行:加一个 @EventListener 方法即可,因为主题只依赖观察者接口、不关心具体监听器。发布方只 publishEvent,新增的监听器自己注册自己,placeOrder() 一行不用改——这正是把「硬编码调用链」解耦成「可插拔监听网络」的价值。


Share this post on:

Previous Post
责任链模式——从Servlet FilterChain到Spring Interceptor
Next Post
装饰器模式——Java IO为什么是一层套一层的洋葱?