基于 Spring 事务同步机制的事务后置动作收集器 Starter 实践
一、涉及的技术知识点
1.1 Spring 事务同步机制
| 知识点 | 说明 |
|---|---|
TransactionSynchronization | Spring 提供的事务同步回调接口,在事务提交/回滚后触发 |
TransactionSynchronizationAdapter | 适配器类,选择性覆写需要的回调方法 |
TransactionSynchronizationManager | 管理当前线程的事务同步器注册 |
afterCommit() | 事务提交成功后调用 |
afterCompletion(status) | 事务完成后调用(无论提交还是回滚),status 标识结果 |
| 事务状态常量 | STATUS_COMMITTED=0、STATUS_ROLLED_BACK=1、STATUS_UNKNOWN=2 |
1.2 设计模式
| 模式 | 应用 |
|---|---|
| 观察者模式 | 收集器作为事务的观察者,监听提交/回滚事件 |
| 命令模式 | AfterTransactionSyncAction接口封装延迟执行的操作 |
| 收集器模式 | 在事务内累积多个动作,事务结束后批量执行 |
| 策略模式 | 提交动作和回滚动作分开注册、分开执行 |
1.3 解决的核心问题
在@Transactional方法内,有些操作必须在事务提交后才能执行:
- 发 MQ 消息:事务未提交就发消息,消费者可能查到旧数据
- 释放分布式锁:事务未提交就解锁,其他线程读到未提交的中间状态
- 调外部接口:事务回滚了但外部接口已调用,无法撤销
- 缓存更新:事务未提交就更新缓存,缓存与数据库不一致
1.4 Java 核心
| 知识点 | 说明 |
|---|---|
| 函数式接口 | AfterTransactionSyncAction只有一个方法,支持 Lambda |
| LinkedList | 保持注册顺序执行 |
| 方法引用 | lock::unlock作为AfterTransactionSyncAction传入 |
注:
博客:
https://blog.csdn.net/badao_liumang_qizhi
二、包结构
xxx.xxx.lib.transaction ├── AfterTransactionSyncAction.java // 提交后动作接口 ├── AfterTransactionRollbackSyncAction.java // 回滚后动作接口 └── AfterTransactionActionCollector.java // 核心收集器(注册到事务同步器)注意:无spring.factories,不是自动配置库。使用时在业务代码中直接new并注册。
三、核心实现流程
@Transactional 方法内 │ ├─→ new AfterTransactionActionCollector() │ ├─→ collector.addCommitSyncAction(() -> sendMqMessage()) // 注册提交后动作 ├─→ collector.addCommitSyncAction(lock::unlock) // 注册释放锁动作 ├─→ collector.addRollbackSyncAction(() -> compensate()) // 注册回滚补偿动作 │ ├─→ TransactionSynchronizationManager │ .registerSynchronization(collector) // 注册到Spring事务管理器 │ ├─→ 业务逻辑执行(数据库操作等) │ └─→ 事务结束 │ ├── 提交成功 → afterCommit() │ └── 按顺序执行所有 commitActions:sendMqMessage() → lock.unlock() │ ├── 回滚 → afterCompletion(STATUS_ROLLED_BACK=1) │ └── 按顺序执行所有 rollbackActions:compensate() │ └── afterCompletion() 最终清理 └── commitActions.clear() + rollbackActions.clear()四、通用示例代码
4.1 AfterTransactionSyncAction(提交后动作接口)
packagecom.example.transaction;/** * 事务提交后执行的动作接口. * 函数式接口,支持Lambda表达式. */@FunctionalInterfacepublicinterfaceAfterTransactionSyncAction{voidexecute();}4.2 AfterTransactionRollbackSyncAction(回滚后动作接口)
packagecom.example.transaction;/** * 事务回滚后执行的动作接口. * 用于补偿操作,如撤销已发的外部通知等. */@FunctionalInterfacepublicinterfaceAfterTransactionRollbackSyncAction{voidexecute();}4.3 AfterTransactionActionCollector(核心收集器)
packagecom.example.transaction;importjava.util.LinkedList;importjava.util.List;importorg.slf4j.Logger;importorg.slf4j.LoggerFactory;importorg.springframework.transaction.support.TransactionSynchronizationAdapter;/** * 事务后置动作收集器. * * 在 @Transactional 方法内使用,收集需要在事务提交/回滚后执行的动作. * * 核心价值: * 1. 保证 MQ 消息在事务提交后才发送(消费者能查到已提交数据) * 2. 保证分布式锁在事务提交后才释放(其他线程读到一致数据) * 3. 支持事务回滚时执行补偿逻辑 * * 使用方式: * 1. new 一个 collector * 2. addCommitSyncAction() 注册提交后动作 * 3. addRollbackSyncAction() 注册回滚后动作 * 4. TransactionSynchronizationManager.registerSynchronization(collector) 注册 */publicclassAfterTransactionActionCollectorextendsTransactionSynchronizationAdapter{privatestaticfinalintSTATUS_COMMITTED=0;privatestaticfinalintSTATUS_ROLLED_BACK=1;privatestaticfinalintSTATUS_UNKNOWN=2;privatefinalList<AfterTransactionSyncAction>commitActions=newLinkedList<>();privatefinalList<AfterTransactionRollbackSyncAction>rollbackActions=newLinkedList<>();privatefinalLoggerlogger=LoggerFactory.getLogger(getClass());/** * 添加事务提交后执行的动作(等同于 addCommitSyncAction). */publicvoidaddSyncAction(AfterTransactionSyncActionaction){logger.debug("add action to thread local: {}",action);commitActions.add(action);}/** * 添加事务提交后执行的动作. */publicvoidaddCommitSyncAction(AfterTransactionSyncActionaction){logger.debug("add commit sync action: {}",action);commitActions.add(action);}/** * 添加事务回滚后执行的动作. */publicvoidaddRollbackSyncAction(AfterTransactionRollbackSyncActionaction){logger.debug("add rollback sync action: {}",action);rollbackActions.add(action);}/** * 事务提交成功后回调. * 按注册顺序依次执行所有 commitActions. */@OverridepublicvoidafterCommit(){if(!commitActions.isEmpty()){for(AfterTransactionSyncActionaction:commitActions){logger.debug("begin to execute after commit action: {}",action);action.execute();}}else{logger.debug("no commitActions to be executed");}}/** * 事务完成后回调(无论提交还是回滚). * 如果是回滚(status=1),执行所有 rollbackActions. * 最后清理所有已注册的动作. */@OverridepublicvoidafterCompletion(intstatus){logger.debug("after transaction with status {}, commit action size: {}, rollback action size: {}",status,commitActions.size(),rollbackActions.size());if(status==STATUS_ROLLED_BACK&&!rollbackActions.isEmpty()){logger.debug("the transaction rolled back! begin to execute roll back actions");for(AfterTransactionRollbackSyncActionaction:rollbackActions){logger.debug("begin to execute after rollback action: {}",action);action.execute();}}// 清理,防止内存泄漏commitActions.clear();rollbackActions.clear();}}4.4 pom.xml(Starter 侧)
<project><groupId>com.example</groupId><artifactId>example-transaction-action-starter</artifactId><version>1.0.0</version><packaging>jar</packaging><dependencies><dependency><groupId>org.springframework</groupId><artifactId>spring-tx</artifactId><scope>provided</scope></dependency><dependency><groupId>org.slf4j</groupId><artifactId>slf4j-api</artifactId><scope>provided</scope></dependency></dependencies></project>五、引入方使用
5.1 添加依赖
<dependency><groupId>com.example</groupId><artifactId>example-transaction-action-starter</artifactId><version>1.0.0</version></dependency>5.2 使用示例1:事务提交后发 MQ
@ServicepublicclassOrderService{@ResourceprivateOrderMqSenderorderMqSender;@Transactional(rollbackFor=Exception.class)publicvoidcreateOrder(OrderDtoorderDto){// 1. 创建收集器并注册到事务管理器AfterTransactionActionCollectorcollector=newAfterTransactionActionCollector();TransactionSynchronizationManager.registerSynchronization(collector);// 2. 执行数据库操作Orderorder=orderRepository.save(convertToEntity(orderDto));// 3. 注册事务提交后发MQ(Lambda 写法)collector.addCommitSyncAction(()->{orderMqSender.sendOrderCreatedMessage(order.getId());});// 4. 其他业务逻辑...// 事务提交后,MQ消息才会发出;如果事务回滚,消息不会发}}5.3 使用示例2:事务提交后释放分布式锁
@ServicepublicclassStockService{@ResourceprivateDistributedLockProviderdistributedLockProvider;@Transactional(rollbackFor=Exception.class)publicvoiddeductStock(IntegeritemId,Integerqty){// 1. 获取锁StringlockKey="stock_deduct_"+itemId;DistributedLocklock=distributedLockProvider.getLock(lockKey);lock.tryLock(TimeUnit.MINUTES,1);// 2. 注册事务提交后释放锁AfterTransactionActionCollectorcollector=newAfterTransactionActionCollector();TransactionSynchronizationManager.registerSynchronization(collector);collector.addCommitSyncAction(lock::unlock);// 方法引用// 3. 执行库存扣减stockRepository.deduct(itemId,qty);// 事务提交后锁才释放 → 其他线程读到的一定是已提交的数据// 事务回滚 → afterCompletion 中 clear 动作列表(但锁不会主动释放,需要等待过期)}}5.4 使用示例3:提交和回滚分别处理
@ServicepublicclassPaymentService{@Transactional(rollbackFor=Exception.class)publicvoidprocessPayment(PaymentDtodto){AfterTransactionActionCollectorcollector=newAfterTransactionActionCollector();TransactionSynchronizationManager.registerSynchronization(collector);// 调用第三方支付(已扣款)StringpaymentId=thirdPartyPayService.charge(dto.getAmount());// 保存支付记录paymentRepository.save(newPayment(paymentId,dto));// 事务提交后:通知下游发货collector.addCommitSyncAction(()->{deliveryService.triggerDelivery(dto.getOrderId());});// 事务回滚后:调第三方退款(补偿)collector.addRollbackSyncAction(()->{thirdPartyPayService.refund(paymentId);});}}5.5 使用示例4:结合分布式锁完整场景
@ServicepublicclassBatchDeliveryService{@Transactional(rollbackFor=Exception.class)publicvoidbatchConfirmDelivery(DeliveryParamsDtoparams){// 获取锁Stringkey="batch_delivery_"+params.getMemberId();DistributedLocklock=distributedLockProvider.getLock(key,TimeUnit.MINUTES,10);lock.tryLock(TimeUnit.MINUTES,8);// 注册事务后释放锁AfterTransactionActionCollectorcollector=newAfterTransactionActionCollector();TransactionSynchronizationManager.registerSynchronization(collector);collector.addCommitSyncAction(lock::unlock);// 执行批量发货业务// ... 写入发货单、写入扩展表(is_fixed_delivery_date)、写出库单 ...// 事务提交后发MQ通知collector.addCommitSyncAction(()->{deliveryMqSender.sendToYc(params.getDeliveryRecordCode());});}}六、关键设计总结
| 设计要点 | 实现方式 | 收益 |
|---|---|---|
| 事务与副作用解耦 | 副作用延迟到事务提交后执行 | MQ 消费者读到已提交数据,外部调用不会因回滚而产生脏数据 |
| 提交/回滚分离 | commitActions + rollbackActions 两个列表 | 不同事务结果执行不同逻辑,支持补偿 |
| 有序执行 | LinkedList 保持注册顺序 | 多个动作按业务预期顺序执行 |
| 自动清理 | afterCompletion 中 clear 两个列表 | 防止内存泄漏 |
| 零配置 | 无 spring.factories,直接 new 使用 | 极简,无侵入 |
| 函数式接口 | 支持 Lambda 和方法引用 | 代码简洁优雅 |
| 最小依赖 | 仅依赖 spring-tx + slf4j | 任何 Spring 事务项目都能用 |