AOP 实战场景 - @Cacheable/@Transactional/@Async/@EventListener 原理
概述
Spring Framework 的 AOP(Aspect-Oriented Programming)能力是许多声明式功能的基础。@Cacheable、@Transactional、@Async 和 @EventListener 这些注解虽然面向不同的问题域,但底层都依赖 Spring AOP 的拦截机制。本文将从源码层面深入分析这四个核心注解的实现原理,揭示它们如何通过 AOP 拦截器链完成各自的功能。
1. 声明式缓存:@Cacheable / @CachePut / @CacheEvict / @Caching
1.1 核心组件架构
Spring 缓存抽象的核心拦截器是 CacheInterceptor,它继承了 CacheAspectSupport,形成了如下拦截器链:
MethodInvocation
↓
CacheInterceptor.invoke(invocation)
↓
CacheAspectSupport.execute(cacheOperationInvoker, method, args, target, targetClass)
↓
解析 @Cacheable/@CachePut/@CacheEvict/@Caching 注解 → CacheOperation 集合
↓
按操作类型执行缓存逻辑1.2 启用方式
@Configuration
@EnableCaching
public class CacheConfig implements CachingConfigurer {
@Bean
@Override
public CacheManager cacheManager() {
return new ConcurrentMapCacheManager("users");
}
}@EnableCaching 通过 @Import(CachingConfigurationSelector.class) 注册了两个基础设施 Bean:
CacheInterceptor:AOP 拦截器,拦截所有带有缓存注解的方法BeanFactoryCacheOperationSourceAdvisor:切面通知器,匹配带有缓存注解的 Bean
1.3 CacheInterceptor 源码分析
// CacheInterceptor.java (Spring 5.3.x)
public class CacheInterceptor extends CacheAspectSupport
implements MethodInterceptor, Serializable {
@Override
@Nullable
public Object invoke(final MethodInvocation invocation) throws Throwable {
Method method = invocation.getMethod();
CacheOperationInvoker aopAllianceInvoker = () -> {
try {
return invocation.proceed();
} catch (Throwable ex) {
throw new CacheOperationInvoker.ThrowableWrapper(ex);
}
};
try {
// 委托给父类 CacheAspectSupport.execute() 执行缓存逻辑
return execute(aopAllianceInvoker, method, invocation.getThis(),
invocation.getArguments());
} catch (CacheOperationInvoker.ThrowableWrapper th) {
throw th.getOriginal();
}
}
}1.4 CacheAspectSupport.execute() 核心逻辑
// CacheAspectSupport.java (Spring 5.3.x) — 核心方法
@Nullable
protected Object execute(CacheOperationInvoker invoker, Method method, Object target,
Object[] args) {
// 1. 获取目标方法上的所有缓存操作(@Cacheable / @CachePut / @CacheEvict 等)
Collection<CacheOperation> operations = getCacheOperationSource()
.getCacheOperations(method, targetClass);
if (operations == null) {
// 没有缓存配置,直接执行原方法
return invoker.invoke();
}
// 2. 为每个操作创建 CacheOperationContext(绑定缓存名称、条件、键等)
CacheOperationContexts contexts = createOperationContexts(operations, method, args, target, targetClass);
// 3. 按操作类型分类执行
Object result = execute(invoker, method, contexts);
return result;
}1.5 @Cacheable 执行流程
// 示例用法
@Service
public class UserService {
@Cacheable(value = "users", key = "#id", unless = "#result == null")
public User getUserById(Long id) {
// 模拟耗时数据库查询
return userRepository.findById(id).orElse(null);
}
}CacheAspectSupport 中对 @Cacheable 的处理流程:
// 伪代码 — CacheAspectSupport 内部逻辑
private Object execute(CacheOperationInvoker invoker, Method method,
CacheOperationContexts contexts) {
// 1. 先执行所有 @CacheEvict(beforeInvocation=true)
processCacheEvicts(contexts.get(CacheEvictOperation.class), true);
// 2. 检查是否命中缓存(仅针对 @Cacheable)
Cache.ValueWrapper cacheHit = findCachedItem(contexts.get(CacheableOperation.class));
if (cacheHit != null) {
// 缓存命中,直接返回缓存值
return cacheHit.get();
}
// 3. 将 @CachePut 和 @Cacheable 的操作集合起来,用于后续写入
Set<CacheOperation> cachePutRequests = new LinkedHashSet<>();
cachePutRequests.addAll(contexts.get(CacheableOperation.class));
Object result;
try {
// 4. 执行目标方法(原始业务逻辑)
result = invoker.invoke();
} catch (Throwable ex) {
// 5. 异常时执行 @CacheEvict(..., beforeInvocation=false)
processCacheEvicts(contexts.get(CacheEvictOperation.class), false);
throw ex;
}
// 6. 方法正常返回后,将结果写入缓存
result = processCachePut(contexts, result, cachePutRequests);
// 7. 执行 @CacheEvict(beforeInvocation=false)
processCacheEvicts(contexts.get(CacheEvictOperation.class), false);
return result;
}1.6 @CachePut / @CacheEvict / @Caching
@Service
public class UserService {
// @CachePut 总会执行方法并将结果写入缓存
@CachePut(value = "users", key = "#user.id")
public User updateUser(User user) {
return userRepository.save(user);
}
// @CacheEvict 清除缓存
@CacheEvict(value = "users", key = "#id")
public void deleteUser(Long id) {
userRepository.deleteById(id);
}
// @CacheEvict(allEntries=true) 清除整个缓存区域
@CacheEvict(value = "users", allEntries = true)
public void clearAllUsers() {
// 清空缓存
}
// @Caching 组合多个缓存操作
@Caching(
cacheable = @Cacheable(value = "users", key = "#username"),
put = @CachePut(value = "users", key = "#result.id"),
evict = @CacheEvict(value = "userIndex", key = "#username")
)
public User findByUsername(String username) {
return userRepository.findByUsername(username).orElse(null);
}
}1.7 缓存条件与 SpEL 支持
@Cacheable(
value = "orders",
key = "#orderId",
condition = "#orderId != null", // 满足条件才进行缓存
unless = "#result.status == 'ERROR'" // 满足条件不缓存结果
)
public Order getOrder(String orderId) {
return orderRepository.findById(orderId).orElse(null);
}1.8 拦截器链小结
调用方 → JdkDynamicAopProxy/CglibAopProxy
→ CacheInterceptor (执行缓存逻辑)
→ 原始方法2. 声明式事务:@Transactional 原理
2.1 启用方式
@Configuration
@EnableTransactionManagement
public class TransactionConfig implements TransactionManagementConfigurer {
@Bean
@Override
public PlatformTransactionManager annotationDrivenTransactionManager() {
return new DataSourceTransactionManager(dataSource());
}
}@EnableTransactionManagement 通过 @Import(TransactionManagementConfigurationSelector.class) 注册:
TransactionInterceptor:核心 AOP 拦截器BeanFactoryTransactionAttributeSourceAdvisor:切面通知器
2.2 TransactionInterceptor 源码
// TransactionInterceptor.java (Spring 5.3.x)
public class TransactionInterceptor extends TransactionAspectSupport
implements MethodInterceptor, Serializable {
public TransactionInterceptor() {}
@Override
@Nullable
public Object invoke(MethodInvocation invocation) throws Throwable {
Method method = invocation.getMethod();
Class<?> targetClass = invocation.getThis() != null
? AopUtils.getTargetClass(invocation.getThis())
: null;
// 适配成 TransactionAspectSupport 需要的回调
return invokeWithinTransaction(method, targetClass,
invocation::proceed, // 业务逻辑回调
invocation.getArguments());
}
}2.3 TransactionAspectSupport.invokeWithinTransaction() 源码
这是事务 AOP 的核心方法,完整分析了声明式事务的拦截链路。
// TransactionAspectSupport.java (Spring 5.3.x)
@Nullable
protected Object invokeWithinTransaction(Method method, @Nullable Class<?> targetClass,
final InvocationCallback invocation, @Nullable Object... args) throws Throwable {
// 1. 获取事务属性 — 解析 @Transactional 注解的配置
TransactionAttributeSource tas = getTransactionAttributeSource();
TransactionAttribute txAttr = (tas != null ? tas.getTransactionAttribute(method, targetClass) : null);
// 2. 确定事务管理器 — 如 @Transactional(transactionManager="xxx")
TransactionManager tm = determineTransactionManager(txAttr);
// 3. 根据事务管理器类型分派执行
if (this instanceof CallbackPreferringPlatformTransactionManager) {
// 回调优先模式 — 少见,略
...
}
// === 标准模式 ===
PlatformTransactionManager ptm = (PlatformTransactionManager) tm;
// 4. 构造事务名称(用于监控日志)
String joinpointIdentification = methodIdentification(method, targetClass, txAttr);
// 5. 开启事务 — 核心
TransactionInfo txInfo = createTransactionIfNecessary(ptm, txAttr, joinpointIdentification);
Object retVal;
try {
// 6. 执行业务逻辑(回调目标方法)
retVal = invocation.proceedWithInvocation();
}
catch (Throwable ex) {
// 7. 异常处理 — 判断是否回滚
completeTransactionAfterThrowing(txInfo, ex);
throw ex;
}
finally {
// 8. 清理事务信息(恢复挂起的事务)
cleanupTransactionInfo(txInfo);
}
// 9. 正常提交事务
commitTransactionAfterReturning(txInfo);
return retVal;
}2.4 开启事务:createTransactionIfNecessary()
// TransactionAspectSupport.java
protected TransactionInfo createTransactionIfNecessary(@Nullable PlatformTransactionManager tm,
@Nullable TransactionAttribute txAttr, final String joinpointIdentification) {
// 如果没有事务属性或事务管理器,返回空 TransactionInfo(不开启事务)
if (txAttr != null && tm != null) {
// 获取当前线程已有事务状态
TransactionStatus status = tm.getTransaction(txAttr);
// 构建事务信息并绑定到当前线程
return prepareTransactionInfo(tm, txAttr, joinpointIdentification, status);
}
return null;
}tm.getTransaction(txAttr) 内部处理了事务传播行为(Propagation):
| 传播行为 | 行为说明 |
|---|---|
REQUIRED(默认) | 当前存在事务则加入,否则新建 |
REQUIRES_NEW | 挂起当前事务,新建一个 |
NESTED | 使用 Savepoint 嵌套事务 |
MANDATORY | 必须在现有事务中运行,否则抛异常 |
SUPPORTS | 有事务则加入,没有也无所谓 |
NOT_SUPPORTED | 挂起当前事务,以非事务方式运行 |
NEVER | 必须在非事务环境中运行,否则抛异常 |
2.5 事务回滚:completeTransactionAfterThrowing()
// TransactionAspectSupport.java
protected void completeTransactionAfterThrowing(@Nullable TransactionInfo txInfo, Throwable ex) {
if (txInfo == null || txInfo.getTransactionStatus() == null) {
return;
}
// 判断异常类型是否需要回滚
// @Transactional(rollbackFor = Exception.class, noRollbackFor = SomeException.class)
if (txInfo.transactionAttribute != null
&& txInfo.transactionAttribute.rollbackOn(ex)) {
try {
// 执行事务回滚
txInfo.getTransactionManager().rollback(txInfo.getTransactionStatus());
}
catch (RuntimeException | Error ex2) {
// 回滚异常记录日志
throw ex2;
}
} else {
// 不需要回滚的异常 — 仍然提交(谨慎行为)
// 但会抛出原始异常
txInfo.getTransactionManager().commit(txInfo.getTransactionStatus());
}
}2.6 事务提交:commitTransactionAfterReturning()
// TransactionAspectSupport.java
protected void commitTransactionAfterReturning(@Nullable TransactionInfo txInfo) {
if (txInfo == null || txInfo.getTransactionStatus() == null) {
return;
}
// 提交事务
txInfo.getTransactionManager().commit(txInfo.getTransactionStatus());
}2.7 事务拦截链路全流程
调用方
↓
AOP 代理
↓
TransactionInterceptor.invoke()
↓
invokeWithinTransaction()
├── 1. 解析 @Transactional 属性
├── 2. 获取事务管理器
├── 3. createTransactionIfNecessary()
│ └── PlatformTransactionManager.getTransaction()
│ ├── 检查当前线程是否存在事务(TransactionSynchronizationManager)
│ ├── 根据传播行为处理(新建/挂起/加入)
│ ├── doBegin() 真正开启事务(如 DataSourceTransactionManager.doBegin())
│ └── 返回 TransactionStatus
├── 4. invocation.proceedWithInvocation() ← 执行业务方法
│ └── 业务代码执行成功/抛出异常
├── 5. 异常分支 → completeTransactionAfterThrowing()
│ └── rollbackOn(ex) ? rollback() : commit() + 抛异常
└── 6. 正常分支 → commitTransactionAfterReturning()
└── PlatformTransactionManager.commit()
├── 验证事务状态
└── doCommit() 真正提交事务Data source:
// DataSourceTransactionManager.doBegin() — 开启数据库事务
@Override
protected void doBegin(Object transaction, TransactionDefinition definition) {
DataSourceTransactionObject txObject = (DataSourceTransactionObject) transaction;
Connection con = txObject.getConnectionHolder().getConnection();
// 设置隔离级别、只读属性
Integer previousIsolationLevel = DataSourceUtils.prepareConnectionForTransaction(con, definition);
txObject.setPreviousIsolationLevel(previousIsolationLevel);
// 关闭自动提交,开启事务
con.setAutoCommit(false);
txObject.setMustRestoreAutoCommit(true);
}2.8 自调用问题
@Service
public class UserService {
@Transactional
public void outerMethod() {
// 调用本类的另一个 @Transactional 方法
this.innerMethod(); // ❌ 事务注解失效!AOP 代理无法拦截自调用
}
@Transactional(propagation = Propagation.REQUIRES_NEW)
public void innerMethod() {
// 这不会在单独的事务中执行
}
}解决方案:
@Service
public class UserService {
@Autowired
private UserService self; // 注入自身代理
@Transactional
public void outerMethod() {
self.innerMethod(); // ✅ 通过代理调用,事务生效
}
@Transactional(propagation = Propagation.REQUIRES_NEW)
public void innerMethod() {
// 正确在独立事务中执行
}
}3. @Async 异步执行原理
3.1 启用方式
@Configuration
@EnableAsync
public class AsyncConfig implements AsyncConfigurer {
@Override
public Executor getAsyncExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(10);
executor.setMaxPoolSize(50);
executor.setQueueCapacity(100);
executor.setThreadNamePrefix("async-");
executor.initialize();
return executor;
}
@Override
public AsyncUncaughtExceptionHandler getAsyncUncaughtExceptionHandler() {
return new SimpleAsyncUncaughtExceptionHandler();
}
}@EnableAsync 通过 @Import(AsyncConfigurationSelector.class) 注册:
AnnotationAsyncExecutionInterceptor(继承AsyncExecutionInterceptor):AOP 方法拦截器AsyncAnnotationAdvisor:切面通知器
3.2 AsyncExecutionInterceptor 源码
// AsyncExecutionInterceptor.java (Spring 5.3.x)
public class AsyncExecutionInterceptor extends AsyncExecutionAspectSupport
implements MethodInterceptor, Ordered {
public AsyncExecutionInterceptor(@Nullable Executor defaultExecutor) {
super(defaultExecutor);
}
@Override
@Nullable
public Object invoke(final MethodInvocation invocation) throws Throwable {
Class<?> targetClass = (invocation.getThis() != null
? AopUtils.getTargetClass(invocation.getThis())
: null);
Method specificMethod = ClassUtils.getMostSpecificMethod(invocation.getMethod(), targetClass);
final Method userDeclaredMethod = BridgeMethodResolver.findBridgedMethod(specificMethod);
// 1. 确定异步执行所需的 TaskExecutor
AsyncTaskExecutor executor = determineAsyncExecutor(userDeclaredMethod);
if (executor == null) {
// 没有找到 Executor,同步执行
return invocation.proceed();
}
// 2. 将方法调用封装为 Callable
Callable<Object> task = () -> {
try {
Object result = invocation.proceed();
// 如果返回值是 Future,等待其完成
if (result instanceof Future) {
((Future<?>) result).get();
}
} catch (ExecutionException ex) {
handleError(ex.getCause(), userDeclaredMethod, invocation.getArguments());
} catch (Throwable ex) {
handleError(ex, userDeclaredMethod, invocation.getArguments());
}
return null;
};
// 3. 提交到线程池异步执行
return doSubmit(task, executor, invocation.getMethod().getReturnType());
}
}3.3 TaskExecutor 选择策略
// AsyncExecutionAspectSupport.java
@Nullable
protected AsyncTaskExecutor determineAsyncExecutor(Method method) {
// 1. 先从缓存中查找
AsyncTaskExecutor executor = this.executors.get(method);
if (executor == null) {
// 2. 解析 @Async("beanName") 指定的 Executor 名称
String qualifier = getExecutorQualifier(method);
if (StringUtils.hasLength(qualifier)) {
// 按名称查找特定的 Executor Bean
executor = determineExecutorByQualifier(qualifier);
}
if (executor == null) {
// 3. 使用默认 Executor
executor = this.defaultExecutor;
if (executor == null) {
// 4. 没有配置 Executor → 使用 SimpleAsyncTaskExecutor
executor = new SimpleAsyncTaskExecutor();
}
}
// 缓存
this.executors.put(method, executor);
}
return executor;
}TaskExecutor 的选择优先级:
@Async("myExecutor")注解中指定的 qualifier- AsyncConfigurer.getAsyncExecutor() 返回的默认执行器
- 如果以上都未配置 →
SimpleAsyncTaskExecutor(每个任务新建一个线程,不推荐生产使用)
3.4 返回值处理 doSubmit()
// AsyncExecutionAspectSupport.java
@Nullable
protected Object doSubmit(Callable<Object> task, AsyncTaskExecutor executor, Class<?> returnType) {
// 根据方法返回类型决定提交方式
if (CompletableFuture.class.isAssignableFrom(returnType)) {
// 返回 CompletableFuture — 提交到线程池
return CompletableFuture.supplyAsync(() -> {
try {
return task.call();
} catch (Throwable ex) {
throw new CompletionException(ex);
}
}, executor);
} else if (ListenableFuture.class.isAssignableFrom(returnType)) {
// 返回 ListenableFuture(Spring 4.x 兼容)
return executor.submitListenable(task);
} else if (Future.class.isAssignableFrom(returnType)) {
// 返回 Future
return executor.submit(task);
} else {
// void 返回类型 — submit 后立即返回 null
executor.submit(task);
return null;
}
}3.5 异常处理
// 用法示例
@Service
public class NotificationService {
@Async
public void sendEmail(String to, String content) {
// 异步执行,异常不会传播到调用方
if (to == null) {
throw new IllegalArgumentException("收件人不能为空");
}
mailSender.send(to, content);
}
}
// @Async 方法有返回值时 — 异常通过 Future.get() 传播
@Async
public Future<String> generateReport(Long id) {
try {
return new AsyncResult<>(reportService.generate(id));
} catch (Exception e) {
return AsyncResult.forExecutionException(e);
}
}AsyncUncaughtExceptionHandler 处理 void 方法的异常:
public class SimpleAsyncUncaughtExceptionHandler implements AsyncUncaughtExceptionHandler {
@Override
public void handleUncaughtException(Throwable ex, Method method, Object... params) {
// 默认只打印日志
logger.error("未捕获的异步异常,方法: " + method, ex);
}
}3.6 @Async 的注意事项
// ❌ 错误:@Async 方法必须是 public 且不能是 static
@Async
private void privateMethod() { } // 不会异步执行!
// ❌ 错误:自调用
@Service
public class MyService {
public void doWork() {
this.asyncMethod(); // 不是通过代理调用,@Async 失效
}
@Async
public void asyncMethod() { }
}
// ✅ 正确
@Service
public class MyService {
@Async
public void asyncMethod() { }
}4. @EventListener 事件监听原理
4.1 事件模型架构
Spring 事件机制基于观察者模式,核心接口:
ApplicationEvent:事件基类ApplicationListener<E>:事件监听器接口ApplicationEventMulticaster:事件广播器ApplicationEventPublisher:事件发布器(注入到 Bean 中)
4.2 启用方式
// Spring Boot 自动启用,无需额外注解
// 纯 Spring 环境需要:
@Configuration
@ComponentScan
public class AppConfig {
// 自动装配 ApplicationEventPublisher
}4.3 EventListenerMethodProcessor — 注册监听器
@EventListener 注解的处理入口是 EventListenerMethodProcessor,它实现了 BeanFactoryPostProcessor,在 Bean 初始化阶段扫描所有带有 @EventListener 的方法。
// EventListenerMethodProcessor.java (Spring 5.3.x)
public class EventListenerMethodProcessor
implements SmartInitializingSingleton, ApplicationContextAware {
@Override
public void afterSingletonsInstantiated() {
// 容器中所有单例 Bean 实例化完成后执行
ConfigurableListableBeanFactory beanFactory = this.beanFactory;
String[] beanNames = beanFactory.getBeanNamesForType(Object.class, false, true);
for (String beanName : beanNames) {
Class<?> type = beanFactory.getType(beanName);
if (type != null && !shouldSkip(beanName, type, beanFactory)) {
// 检查每个 Bean 中是否有 @EventListener 方法
processBean(beanName, type, beanFactory);
}
}
}
private void processBean(String beanName, Class<?> targetType,
ConfigurableListableBeanFactory beanFactory) {
// 查找所有 @EventListener 方法
Map<Method, EventListener> annotatedMethods = MethodIntrospector.selectMethods(
targetType, (MethodIntrospector.MetadataLookup<EventListener>) method ->
AnnotatedElementUtils.findMergedAnnotation(method, EventListener.class));
if (annotatedMethods.isEmpty()) {
return;
}
// 为每个带注解的方法创建适配器并注册到事件广播器
for (Method method : annotatedMethods.keySet()) {
EventListener eventListener = annotatedMethods.get(method);
// 创建 ApplicationListenerMethodAdapter
for (Class<?> eventType : eventListener.classes()) {
ApplicationListenerMethodAdapter listener = new ApplicationListenerMethodAdapter(
beanName, targetType, method);
// 注册到 ApplicationEventMulticaster
beanFactory.getBean(APPLICATION_EVENT_MULTICASTER_BEAN_NAME,
ApplicationEventMulticaster.class).addApplicationListener(listener);
}
}
}
}4.4 ApplicationListenerMethodAdapter — 方法适配器
// ApplicationListenerMethodAdapter.java (Spring 5.3.x)
public class ApplicationListenerMethodAdapter
implements GenericApplicationListener {
private final String beanName;
private final Method method;
private final Method bridgedMethod;
private final ResolvableType declaredEventType;
private final List<ResolvableType> declaredEventTypes;
@Override
public void onApplicationEvent(ApplicationEvent event) {
// 执行监听方法
processEvent(event);
}
public void processEvent(ApplicationEvent event) {
// 提取事件对象 — 支持 ApplicationEvent 或任意 POJO
Object[] args = resolveArguments(event);
// 执行监听方法
Object result = doInvoke(args);
// 处理返回值 — 如果 @EventListener 方法有返回值,
// 且返回非 null,会发布为新事件
if (result != null) {
handleResult(result);
}
}
@Nullable
protected Object doInvoke(Object... args) {
// 通过反射调用目标 Bean 的监听方法
Object bean = getBean();
ReflectionUtils.makeAccessible(this.method);
return this.method.invoke(bean, args);
}
// 返回值发布为新事件 — 事件驱动链
protected void handleResult(Object result) {
if (result instanceof Iterable<?> iterable) {
// 返回集合时,每个元素作为独立事件发布
for (Object event : iterable) {
publishEvent(event);
}
} else {
publishEvent(result);
}
}
}4.5 @EventListener 基础用法
// 自定义事件
public class OrderCreatedEvent extends ApplicationEvent {
private final Long orderId;
private final String orderNo;
public OrderCreatedEvent(Object source, Long orderId, String orderNo) {
super(source);
this.orderId = orderId;
this.orderNo = orderNo;
}
// getters...
}
// 事件监听器
@Component
public class OrderEventListener {
private static final Logger log = LoggerFactory.getLogger(OrderEventListener.class);
@EventListener
public void handleOrderCreated(OrderCreatedEvent event) {
log.info("收到订单创建事件: orderId={}, orderNo={}",
event.getOrderId(), event.getOrderNo());
// 执行后续业务逻辑:发送通知、更新库存等
}
// 条件过滤 — SpEL 表达式
@EventListener(condition = "#event.orderNo.startsWith('ORD')")
public void handleImportantOrder(OrderCreatedEvent event) {
log.info("处理重要订单: {}", event.getOrderNo());
}
}
// 事件发布
@Component
public class OrderService {
@Autowired
private ApplicationEventPublisher publisher;
public void createOrder(Order order) {
// ... 创建订单逻辑 ...
// 发布事件
publisher.publishEvent(new OrderCreatedEvent(this, order.getId(), order.getOrderNo()));
}
}4.6 事件转发链
@EventListener 方法有返回值时,返回值会自动发布为新事件:
@Component
public class OrderEventChain {
@EventListener
public OrderPaidEvent handleOrderCreated(OrderCreatedEvent event) {
// 处理订单创建
// 返回一个新事件,Spring 自动发布
return new OrderPaidEvent(this, event.getOrderId());
}
@EventListener
public void handleOrderPaid(OrderPaidEvent event) {
// 支付事件处理
// 监听 OrderPaidEvent,形成事件链
}
}4.7 多事件类型与泛型事件
// 监听多个事件类型
@Component
public class MultiEventListener {
@EventListener({OrderCreatedEvent.class, OrderPaidEvent.class})
public void handleOrderEvents(Object event) {
if (event instanceof OrderCreatedEvent) {
// 处理创建
} else if (event instanceof OrderPaidEvent) {
// 处理支付
}
}
}
// 泛型事件(Spring 4.2+)
public class EntityCreatedEvent<T> extends ApplicationEvent {
private final T entity;
public EntityCreatedEvent(Object source, T entity) {
super(source);
this.entity = entity;
}
public T getEntity() { return entity; }
}
@Component
public class GenericEventListener {
@EventListener
public void handleUserCreated(EntityCreatedEvent<User> event) {
// ResolvableType 能正确解析泛型参数
User user = event.getEntity();
}
}5. @TransactionalEventListener 事务事件
5.1 事务阶段枚举
// TransactionPhase.java (Spring 5.3.x)
public enum TransactionPhase {
/** 事务提交成功后触发(默认) */
AFTER_COMMIT,
/** 事务回滚后触发 */
AFTER_ROLLBACK,
/** 事务完成后触发(无论是提交还是回滚) */
AFTER_COMPLETION,
/** 在事务执行前触发(不要求有活跃事务) */
BEFORE_COMMIT
}5.2 实现原理
@TransactionalEventListener 在 @EventListener 的基础上增加了事务阶段控制。其核心处理位于 TransactionalEventListenerFactory 和 ApplicationListenerMethodTransactionalAdapter。
// TransactionalEventListenerFactory.java (Spring 5.3.x)
public class TransactionalEventListenerFactory implements EventListenerFactory {
@Override
public boolean supportsMethod(Method method) {
// 只处理同时标注 @TransactionalEventListener 的方法
return AnnotatedElementUtils.hasAnnotation(method, TransactionalEventListener.class);
}
@Override
public ApplicationListener<?> createApplicationListener(String beanName, Class<?> type, Method method) {
// 创建事务感知的监听器适配器
return new ApplicationListenerMethodTransactionalAdapter(beanName, type, method);
}
}5.3 ApplicationListenerMethodTransactionalAdapter
// ApplicationListenerMethodTransactionalAdapter.java (Spring 5.3.x)
public class ApplicationListenerMethodTransactionalAdapter
extends ApplicationListenerMethodAdapter {
private final TransactionalEventListener annotation;
public ApplicationListenerMethodTransactionalAdapter(
String beanName, Class<?> targetClass, Method method) {
super(beanName, targetClass, method);
this.annotation = AnnotatedElementUtils
.findMergedAnnotation(method, TransactionalEventListener.class);
}
@Override
public void onApplicationEvent(ApplicationEvent event) {
// 如果当前没有活跃事务,直接处理事件
if (TransactionSynchronizationManager.isSynchronizationActive()) {
// 将事件处理注册为事务同步回调
TransactionSynchronizationManager.registerSynchronization(
new TransactionalSynchronization(event));
} else {
// 没有事务时,默认行为:日志警告(可配置 fallbackExecution 忽略)
if (this.annotation.fallbackExecution()) {
super.onApplicationEvent(event);
} else {
logger.warn("没有活跃事务,忽略 @TransactionalEventListener");
}
}
}
// 事务同步 — 控制事件在指定阶段触发
private class TransactionalSynchronization extends TransactionSynchronizationAdapter {
private final ApplicationEvent event;
@Override
public int getOrder() {
return this.annotation.phase() == TransactionPhase.BEFORE_COMMIT
? TransactionSynchronization.PRECEDENCE : TransactionSynchronization.SUPPORTS;
}
@Override
public void beforeCommit(boolean readOnly) {
if (this.annotation.phase() == TransactionPhase.BEFORE_COMMIT) {
processEvent(this.event);
}
}
@Override
public void afterCommit() {
if (this.annotation.phase() == TransactionPhase.AFTER_COMMIT) {
processEvent(this.event);
}
}
@Override
public void afterCompletion(int status) {
if (this.annotation.phase() == TransactionPhase.AFTER_ROLLBACK
&& status == STATUS_ROLLED_BACK) {
processEvent(this.event);
} else if (this.annotation.phase() == TransactionPhase.AFTER_COMPLETION) {
processEvent(this.event);
}
}
}
}5.4 使用示例
@Component
public class OrderEventListener {
private static final Logger log = LoggerFactory.getLogger(OrderEventListener.class);
// 默认 AFTER_COMMIT — 事务提交成功后执行
@TransactionalEventListener
public void handleOrderCreatedAfterCommit(OrderCreatedEvent event) {
// 事务已提交,数据已持久化,可以安全地:
// 1. 发送 MQ 消息
// 2. 发送邮件/短信
// 3. 调用外部 API
log.info("事务已提交,发送订单通知: {}", event.getOrderNo());
notificationService.sendOrderConfirmation(event.getOrderId());
}
// 事务回滚后执行
@TransactionalEventListener(phase = TransactionPhase.AFTER_ROLLBACK)
public void handleOrderRollback(OrderCreatedEvent event) {
log.warn("订单创建已回滚,清理临时资源: {}", event.getOrderNo());
}
// 事务提交前执行(仍在事务中)
@TransactionalEventListener(phase = TransactionPhase.BEFORE_COMMIT)
public void handleBeforeCommit(OrderCreatedEvent event) {
log.info("事务即将提交,执行预提交检查: {}", event.getOrderNo());
}
// 没有事务时也执行
@TransactionalEventListener(fallbackExecution = true)
public void handleAlways(OrderCreatedEvent event) {
// fallbackExecution=true,即使没有事务也执行
}
}5.5 与常规 @EventListener 的对比
| 特性 | @EventListener | @TransactionalEventListener |
|---|---|---|
| 触发时机 | 事件发布后立即执行 | 可控制在各事务阶段触发 |
| 事务关联 | 无关联 | 与事务绑定 |
| 默认阶段 | N/A | AFTER_COMMIT |
| 无事务时的行为 | 正常执行 | 默认忽略(可配 fallbackExecution) |
| 典型场景 | 内存内处理 | 需要确保数据已持久化后的处理 |
6. 总结:四种 AOP 场景对比
| 特性 | @Cacheable系列 | @Transactional | @Async | @EventListener |
|---|---|---|---|---|
| 核心拦截器 | CacheInterceptor | TransactionInterceptor | AsyncExecutionInterceptor | EventListenerMethodProcessor |
| 通知器 | BeanFactoryCacheOperationSourceAdvisor | BeanFactoryTransactionAttributeSourceAdvisor | AsyncAnnotationAdvisor | N/A(后处理器) |
| 切入点 | 方法上的缓存注解 | 方法上的 @Transactional | 方法上的 @Async | 方法上的 @EventListener |
| 通知类型 | Around | Around | Around | After(事件发布) |
| 线程模型 | 同步(调用方线程) | 同步(调用方线程) | 异步(线程池) | 同步(发布者线程) |
| 异常影响调用方 | 是 | 是 | 否(void 方法) | 是(默认) |
| 代理需求 | 外部调用 | 外部调用 | 外部调用 | 外部调用 |
所有四种机制都依赖于 Spring AOP 代理的核心能力——无论是通过 JdkDynamicAopProxy 还是 CglibAopProxy,拦截器拦截方法调用,插入横切逻辑,最终再通过反射或 CGLIB 调用原始目标方法,构成了 Spring 声明式编程模型的基石。