为何@Transactional在Flux doFinally调用的服务方法中不生效?
问题分析与解答
代码重现
@Service @Transactional public class MyService{ public Flux<Entity1> method1(){ // some non flux code return Flux.interval(Duration.ofSeconds(1)) .filter(...) .doFinally(signalType -> updateDataSource(entity1)); } public void updateDataSource(Entity1 entity1){ myRepo.save(entity1); // code that can throw some exceptions myRepo.save(entity2); } }
事务不生效的核心原因
1. 内部方法调用绕过Spring代理
Spring的@Transactional依赖动态代理实现事务拦截,只有通过代理对象调用方法时,事务逻辑才会触发。但method1里直接调用this.updateDataSource(...)属于类内部调用,完全绕过了代理对象,事务拦截器根本无法介入,自然不会对updateDataSource的操作做事务管理。
2. 反应式流的懒加载导致事务上下文丢失
Flux是懒加载模型,method1返回Flux时,流的执行逻辑(包括doFinally里的代码)并未启动,真正的执行要等到订阅阶段才触发。而类上的@Transactional是同步事务注解,它的事务上下文会在method1执行完毕(即返回Flux的瞬间)就结束。等到doFinally执行updateDataSource时,已经没有有效的事务上下文,所有数据库操作都是无事务的独立提交。
3. 同步事务注解不匹配反应式场景
你的代码使用了反应式的Flux,但默认的@Transactional是为同步方法设计的。反应式场景需要搭配ReactiveTransactionManager,并使用@Transactional(reactive = true),或者直接用编程式反应式事务API(如TransactionalOperator)来对齐事务与反应式流的生命周期。
修复方案
方案1:解决内部调用代理问题
通过AopContext获取自身代理对象,确保调用走代理触发事务:
// 先在配置类开启代理暴露 @EnableAspectJAutoProxy(exposeProxy = true) // 修改MyService的method1和updateDataSource public Flux<Entity1> method1(){ // some non flux code MyService proxyService = AopContext.currentProxy(); return Flux.interval(Duration.ofSeconds(1)) .filter(...) .doFinally(signalType -> proxyService.updateDataSource(entity1)); } // 给updateDataSource单独标注事务 @Transactional public void updateDataSource(Entity1 entity1){ myRepo.save(entity1); // code that can throw some exceptions myRepo.save(entity2); }
方案2:适配反应式事务场景
用编程式反应式事务API,确保事务与反应式流生命周期对齐:
@Service public class MyService{ private final TransactionalOperator transactionalOperator; private final MyRepo myRepo; // 构造注入依赖 public MyService(TransactionalOperator transactionalOperator, MyRepo myRepo) { this.transactionalOperator = transactionalOperator; this.myRepo = myRepo; } public Flux<Entity1> method1(){ // some non flux code return Flux.interval(Duration.ofSeconds(1)) .filter(...) .doFinally(signalType -> { // 用事务包裹更新逻辑 Mono.fromRunnable(() -> { myRepo.save(entity1); // code that can throw some exceptions myRepo.save(entity2); }).as(transactionalOperator::transactional) .block(); // 非反应式上下文用block,反应式上下文可直接返回Mono并订阅 }); } }
内容的提问来源于stack exchange,提问作者Gestalt
相关产品推荐
相关产品推荐

