基于 Spring 事务同步机制的事务后置动作收集器 Starter 实践

基于 Spring 事务同步机制的事务后置动作收集器 Starter 实践


一、涉及的技术知识点

1.1 Spring 事务同步机制

知识点说明
TransactionSynchronizationSpring 提供的事务同步回调接口,在事务提交/回滚后触发
TransactionSynchronizationAdapter适配器类,选择性覆写需要的回调方法
TransactionSynchronizationManager管理当前线程的事务同步器注册
afterCommit()事务提交成功后调用
afterCompletion(status)事务完成后调用(无论提交还是回滚),status 标识结果
事务状态常量STATUS_COMMITTED=0STATUS_ROLLED_BACK=1STATUS_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 事务项目都能用