观察者模式与 Spring 事件机制解耦核心业务通知
在大型企业级应用开发中,随着业务迭代推进,核心主链路代码往往会面临严重的“功能膨胀”与“强耦合”危机。以电商系统的“订单创建成功”或“支付成功”为例,最初的代码可能只有简洁的几行订单状态更新与库存扣减;但随着业务演进,短信通知、站内信推送、积分发放、分销佣金核算、仓储履约下发、风控特征打标等周边逻辑被不断堆砌进同一个业务方法中。
这种面条式代码不仅使得单一业务类突破数千行,更带来了致命的架构风险:
- 长事务蔓延:短信发送或第三方 RPC 耗时被拉入本地数据库事务中,导致数据库连接被长时间占用甚至引发连接池打满。
- 故障蔓延:非核心业务(如积分赠送服务抖动)抛出异常,导致主业务(用户支付)整体回滚失败。
- 维护成本剧增:新增一个下游监听业务需要反复修改核心业务代码,违反开闭原则(OCP)。
本文介绍如何利用经典的观察者模式(Observer Pattern),结合 Spring 的事件发布订阅机制(ApplicationEvent & @TransactionalEventListener),实现核心链路与旁路通知的高可靠解耦。
经典观察者模式与 Spring 事件演进
GoF 观察者模式定义了对象间一种一对多的依赖关系,当一个对象状态发生改变时,所有依赖于它的对象都会得到通知并自动更新。
[传统紧耦合链路] [OrderService] ──> [DB 写入] ──> [SmsService] ──> [PointsService] ──> [LogisticsService] (任何一个下游阻塞或报错,整个事务一同挂掉) [Spring 事务事件解耦链路] [OrderService] ──> [DB 写入 & 事务提交] │ ▼ (发布 OrderPaidEvent) [ApplicationEventMulticaster] │ ├─ [SmsEventListener] (异步线程池执行 / AFTER_COMMIT) ├─ [PointsEventListener] (异步线程池执行 / AFTER_COMMIT) └─ [LogisticsEventListener] (MQ 桥接投递 / AFTER_COMMIT)Spring 框架对观察者模式进行了高阶封装:
- Subject(目标发布者):通过
ApplicationEventPublisher统一发布事件。 - Observer(观察者):通过
@EventListener或@TransactionalEventListener注解实现无侵入监听。 - EventMulticaster(广播器):支持同步调度与异步线程池广播。
核心实现与实战代码
1. 领域事件定义
事件对象推荐定义为不可变 DTO,仅携带核心业务主键与状态快照:
package com.example.event.order; import java.io.Serializable; import java.math.BigDecimal; import java.time.LocalDateTime; public class OrderPaidEvent implements Serializable { private static final long serialVersionUID = 1L; private final Long orderId; private final Long userId; private final BigDecimal paidAmount; private final LocalDateTime paidTime; public OrderPaidEvent(Long orderId, Long userId, BigDecimal paidAmount, LocalDateTime paidTime) { this.orderId = orderId; this.userId = userId; this.paidAmount = paidAmount; this.paidTime = paidTime; } public Long getOrderId() { return orderId; } public Long getUserId() { return userId; } public BigDecimal getPaidAmount() { return paidAmount; } public LocalDateTime getPaidTime() { return paidTime; } }2. 核心主业务只负责事件发布与事务提交
业务服务注入ApplicationEventPublisher,在完成数据库核心状态变更后发布事件:
package com.example.event.service; import com.example.event.order.OrderPaidEvent; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.context.ApplicationEventPublisher; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; import java.math.BigDecimal; import java.time.LocalDateTime; @Service public class OrderCheckoutService { private static final Logger log = LoggerFactory.getLogger(OrderCheckoutService.class); private final ApplicationEventPublisher eventPublisher; private final OrderRepository orderRepository; public OrderCheckoutService(ApplicationEventPublisher eventPublisher, OrderRepository orderRepository) { this.eventPublisher = eventPublisher; this.orderRepository = orderRepository; } @Transactional(rollbackFor = Exception.class) public void markOrderPaid(Long orderId, Long userId, BigDecimal amount) { // 1. 执行核心表状态更新(本地事务) orderRepository.updateStatusToPaid(orderId, amount); log.info("订单 [{}] 状态成功流转为 PAID", orderId); // 2. 核心主链路仅发布领域事件,不直接调用任何旁路业务 OrderPaidEvent event = new OrderPaidEvent(orderId, userId, amount, LocalDateTime.now()); eventPublisher.publishEvent(event); } }3. 多维度监听器实现与事务阶段绑定
在旁路通知处理中,最关键的考量是事务生命周期绑定与异步线程隔离。使用@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)能够确保只有当数据库事务真正 Commit 成功后,才触发下游通知,彻底杜绝“主事务回滚但短信/积分已发出”的脏事件灾难。
package com.example.event.listener; import com.example.event.order.OrderPaidEvent; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.scheduling.annotation.Async; import org.springframework.stereotype.Component; import org.springframework.transaction.event.TransactionPhase; import org.springframework.transaction.event.TransactionalEventListener; @Component public class OrderNotificationListeners { private static final Logger log = LoggerFactory.getLogger(OrderNotificationListeners.class); /** * 短信通知监听:事务提交后异步执行,不占用主业务线程 */ @Async("businessEventExecutor") @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT) public void handleSmsNotification(OrderPaidEvent event) { try { log.info("【短信服务】正在为用户 [{}] 发送订单 [{}] 支付成功通知", event.getUserId(), event.getOrderId()); // 模拟调用第三方短信网关 RPC Thread.sleep(100); } catch (Exception e) { // 旁路异常仅记录报警日志,绝不影响主流程 log.error("【短信服务】发送失败,订单ID: {}", event.getOrderId(), e); } } /** * 积分发放监听:事务提交后异步执行 */ @Async("businessEventExecutor") @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT) public void handlePointsAward(OrderPaidEvent event) { try { int rewardPoints = event.getPaidAmount().intValue() * 10; log.info("【积分中心】为用户 [{}] 发放支付激励积分: {}", event.getUserId(), rewardPoints); // 执行积分变更逻辑 } catch (Exception e) { log.error("【积分中心】积分发放异常,记录补偿重试表", e); } } }4. 异步线程池配置
package com.example.event.config; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.scheduling.annotation.EnableAsync; import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; import java.util.concurrent.Executor; import java.util.concurrent.ThreadPoolExecutor; @Configuration @EnableAsync public class AsyncEventConfig { @Bean("businessEventExecutor") public Executor businessEventExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(8); executor.setMaxPoolSize(32); executor.setQueueCapacity(500); executor.setThreadNamePrefix("async-event-"); // 队列打满时由调用方线程兜底执行,防止丢弃事件 executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); executor.initialize(); return executor; } }生产避坑与架构准则
@Async与@TransactionalEventListener的配合规范:
在 Spring 默认机制下,若未开启@Async,@TransactionalEventListener仍会在主线程中以同步阻塞方式运行。若在AFTER_COMMIT阶段同步执行慢 RPC,依然会拖慢前端请求响应。因此,旁路通知监听器通常必须搭配@Async异步化。- 事务上下文与只读限制:
在AFTER_COMMIT阶段,主事务已经物理提交并关闭了数据库连接。如果在监听器内部再次执行更新数据库操作,必须在监听器方法上显式开启新事务(@Transactional(propagation = Propagation.REQUIRES_NEW)),否则可能因使用已关闭的 Connection 而抛出连接已关闭异常。 - 单机事件与分布式 MQ 的演进界限:
Spring 进程内事件机制非常适合单体服务内部组件解耦与轻量级微服务领域事件分发。但如果监听方属于跨系统服务(如物流系统、大数据数仓系统),或者要求极高的一致性(即使应用 Crash 也绝对不能丢失通知),则应在监听器中将事件桥接并投递至分布式消息中间件(如 RocketMQ/Kafka),实现进程外的可靠广播。