You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

为何@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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.24 21:53:20