RxJava 1.1.5与Spring Boot事务不回滚问题求助
你猜得没错!这确实是RxJava的懒加载特性和Spring事务的线程绑定机制冲突导致的事务不生效。
Spring的@Transactional是基于ThreadLocal实现的——它会把事务上下文绑定到当前执行doProcess方法的线程上。但RxJava的Observable是懒执行的:你在doProcess里调用的repository1.insertRecord1等方法,并不会立刻执行数据库操作,只是创建了一个Observable任务;真正的数据库操作要等到这个Observable被订阅的时候才会触发。
而订阅动作往往发生在doProcess方法执行完成之后(比如调用方在其他线程订阅),或者RxJava默认切换到其他线程执行操作——这时候执行数据库操作的线程已经没有事务上下文了,每个Repository操作都在独立的非事务环境下运行,自然不会触发回滚。
针对你的场景(RxJava 1.1.5 + Spring Boot 1.3.5),有几种可行的修复方案:
方案1:强制所有Observable操作在事务线程执行
让RxJava不要切换线程,确保所有数据库操作都在doProcess方法所在的事务线程中执行。可以通过subscribeOn(Schedulers.immediate())强制每个Observable在当前线程执行:
@Transactional(rollbackFor=DataException.class) public Observable<Optional<Boolean>> doProcess(req1,req2,req3,req4){ return Observable.zip( // 给每个Repository调用添加subscribeOn,强制在当前线程执行 repository1.insertRecord1(req1.get(),useFallBack).subscribeOn(Schedulers.immediate()), repository1.insertRecord2(req2.get(),useFallBack).subscribeOn(Schedulers.immediate()), repository1.getJRecord(req3.getAccountId(),req4.getId(),useFallBack).subscribeOn(Schedulers.immediate()), repository1.getTaxRate(req3.getAccountId(),useFallBack).subscribeOn(Schedulers.immediate()), repository1.getCurrencyByCode(req3.getCurrencyCode(), useFallBack).subscribeOn(Schedulers.immediate()), (record1,record2,jRecord,taxRate, currencyOp)->{ BigDecimal amountWithoutTax = req4.getAmount().divide(taxRate.get(), currencyOp.get().getDecimalPlaces(), BigDecimal.ROUND_HALF_UP); Record3 record3= new Record3(); BigDecimal chargeJournalKey = req4.getId(); record3.setKey(record2.get()); record3.setCollectionKey(record1.get()); record3.setAmount(jRecord.get().getAmount().subtract(req4.getAmount())); record3.setAmount(amountWithoutTax); // 同样给insertRecord3添加线程指定 return repository1.insertRecord3(record3).subscribeOn(Schedulers.immediate()); }); }
注意:如果你的Repository方法内部自己用了异步线程(比如内部开了新线程执行数据库操作),这个方案就不生效了——必须确保所有数据库操作都在事务线程中执行。
方案2:手动传递事务上下文到RxJava线程
如果必须用异步线程执行,可以手动把Spring的事务上下文传递到RxJava的执行线程中。你可以自定义一个RxJava Scheduler,把当前线程的ThreadLocal事务上下文复制到新线程:
public class TransactionAwareScheduler extends Scheduler { private final Scheduler delegate; private final DataSource dataSource; public TransactionAwareScheduler(Scheduler delegate, DataSource dataSource) { this.delegate = delegate; this.dataSource = dataSource; } @Override public Worker createWorker() { // 获取当前线程的事务上下文 Object transactionContext = TransactionSynchronizationManager.getResource(dataSource); return new TransactionAwareWorker(delegate.createWorker(), transactionContext); } private class TransactionAwareWorker extends Worker { private final Worker delegate; private final Object transactionContext; public TransactionAwareWorker(Worker delegate, Object transactionContext) { this.delegate = delegate; this.transactionContext = transactionContext; } @Override public Subscription schedule(Action0 action) { return delegate.schedule(() -> { // 在新线程中设置事务上下文 if (transactionContext != null) { TransactionSynchronizationManager.bindResource(dataSource, transactionContext); } try { action.call(); } finally { // 清理上下文 if (transactionContext != null) { TransactionSynchronizationManager.unbindResource(dataSource); } } }); } @Override public void unsubscribe() { delegate.unsubscribe(); } @Override public boolean isUnsubscribed() { return delegate.isUnsubscribed(); } } }
然后在doProcess里使用这个自定义Scheduler:
@Autowired private DataSource dataSource; @Transactional(rollbackFor=DataException.class) public Observable<Optional<Boolean>> doProcess(req1,req2,req3,req4){ Scheduler transactionAwareScheduler = new TransactionAwareScheduler(Schedulers.io(), dataSource); return Observable.zip( repository1.insertRecord1(req1.get(),useFallBack).subscribeOn(transactionAwareScheduler), repository1.insertRecord2(req2.get(),useFallBack).subscribeOn(transactionAwareScheduler), // ... 其他Repository调用同理 (record1,record2,jRecord,taxRate, currencyOp)->{ // ... 业务逻辑 return repository1.insertRecord3(record3).subscribeOn(transactionAwareScheduler); }); }
这个方案适合需要异步执行的场景,但实现起来比较复杂,需要注意资源清理避免内存泄漏。
方案3:阻塞等待Observable结果(牺牲异步性)
如果你的业务场景可以接受同步执行,可以在doProcess方法内部阻塞等待Observable的执行结果,确保整个流程都在事务线程中完成:
@Transactional(rollbackFor=DataException.class) public Observable<Optional<Boolean>> doProcess(req1,req2,req3,req4){ // 阻塞执行所有操作,获取结果 Optional<Boolean> result = Observable.zip( repository1.insertRecord1(req1.get(),useFallBack), repository1.insertRecord2(req2.get(),useFallBack), repository1.getJRecord(req3.getAccountId(),req4.getId(),useFallBack), repository1.getTaxRate(req3.getAccountId(),useFallBack), repository1.getCurrencyByCode(req3.getCurrencyCode(), useFallBack), (record1,record2,jRecord,taxRate, currencyOp)->{ BigDecimal amountWithoutTax = req4.getAmount().divide(taxRate.get(), currencyOp.get().getDecimalPlaces(), BigDecimal.ROUND_HALF_UP); Record3 record3= new Record3(); BigDecimal chargeJournalKey = req4.getId(); record3.setKey(record2.get()); record3.setCollectionKey(record1.get()); record3.setAmount(jRecord.get().getAmount().subtract(req4.getAmount())); record3.setAmount(amountWithoutTax); return repository1.insertRecord3(record3).toBlocking().single(); }).toBlocking().single(); // 把结果包装成Observable返回 return Observable.just(result); }
这个方案最简单,但会失去RxJava的异步优势,适合对性能要求不高的场景。
核心原则就是:所有需要在同一事务中执行的数据库操作,必须在同一个持有Spring事务上下文的线程中运行。根据你的业务需求选择合适的方案即可。
内容的提问来源于stack exchange,提问作者ankit patel

