如何在多线程场景下结合使用@Async与@Transactional实现事务回滚
问题根因分析
- 多线程未生效、启动仅创建1个线程的原因:
你把配置类注解@Configuration和业务类注解混用,业务方法上加了@Bean注解,Spring启动时会主动调用标注了@Bean的方法注册容器组件,你看到的单个线程就是启动阶段的调用。另外如果@Async标注的方法在本类直接调用,不会走Spring AOP代理,异步逻辑不生效,自然跑不起来多线程。 - 事务回滚失效的原因:
Spring原生事务上下文是绑定在当前线程的ThreadLocal中的,多线程场景下每个线程的事务完全独立,互相没有感知,单个线程抛出异常只会触发自身事务回滚,不会通知其他线程做回滚操作。
修复步骤
1. 修正基础注解错误,恢复多线程能力
首先拆分配置类和业务类,移除业务方法上的@Bean注解,保证异步方法被外部类调用触发AOP代理。
线程池配置类(单独抽离)
package com.demo.multithread.config; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.scheduling.annotation.EnableAsync; import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; import java.util.concurrent.ExecutorService; @Configuration @EnableAsync public class AsyncConfig { @Bean(name = "threadPoolTaskExecutor") public ExecutorService executorService() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(10); executor.setMaxPoolSize(10); executor.setQueueCapacity(10); executor.setThreadNamePrefix("myThread-"); executor.initialize(); return executor.getThreadPoolExecutor(); } }
DemoService修正后代码
package com.demo.multithread.service; import org.springframework.stereotype.Repository; import org.springframework.transaction.annotation.Transactional; import javax.persistence.EntityManager; import javax.persistence.PersistenceContext; import javax.persistence.StoredProcedureQuery; @Repository public class DemoService { @PersistenceContext private EntityManager entityManager; @Transactional(rollbackFor = {Exception.class, Error.class}) public void callProcedure() { StoredProcedureQuery query = entityManager.createNamedStoredProcedureQuery("MyProcedure"); query.execute(); } }
异步调用类修正后代码
package com.demo.multithread.service; import org.springframework.scheduling.annotation.Async; import org.springframework.stereotype.Repository; import javax.annotation.Resource; @Repository public class DemoDAO { @Resource private DemoService demoService; @Async(value = "threadPoolTaskExecutor") public void asyncMethodWithConfiguredExecutor() { System.out.println("Thread ID~" + Thread.currentThread().getId() + " 运行中,线程名~" + Thread.currentThread().getName()); demoService.callProcedure(); } }
注意:异步方法必须从其他类注入后调用,不能在本类直接调用,否则AOP不生效,多线程不会触发。
2. 实现多线程全局事务统一回滚
Spring原生的@Transactional不支持跨线程全局事务控制,需要手动实现两阶段提交逻辑,用闭锁协调所有线程的提交/回滚动作,参考实现如下:
事务回调接口定义
@FunctionalInterface public interface TransactionCallback { void execute() throws Exception; }
多线程事务协调工具类
@Component public class MultiThreadTransactionManager { @Resource private PlatformTransactionManager transactionManager; public void execute(List<TransactionCallback> tasks, ExecutorService executor) throws InterruptedException { int taskCount = tasks.size(); // 协调器等待所有任务执行完成的闭锁 CountDownLatch allTaskFinishedLatch = new CountDownLatch(taskCount); // 所有任务等待协调器发送提交/回滚信号的闭锁 CountDownLatch commitSignalLatch = new CountDownLatch(1); // 标记是否有任务执行异常 AtomicBoolean hasException = new AtomicBoolean(false); for (TransactionCallback task : tasks) { executor.submit(() -> { TransactionStatus status = transactionManager.getTransaction(new DefaultTransactionDefinition()); try { // 执行业务逻辑 task.execute(); } catch (Exception e) { hasException.set(true); e.printStackTrace(); } finally { allTaskFinishedLatch.countDown(); } try { // 等待协调器通知 commitSignalLatch.await(); if (hasException.get()) { transactionManager.rollback(status); } else { transactionManager.commit(status); } } catch (InterruptedException e) { transactionManager.rollback(status); Thread.currentThread().interrupt(); } }); } // 等待所有任务执行完成,统计执行结果 allTaskFinishedLatch.await(); // 发送统一提交/回滚信号 commitSignalLatch.countDown(); } }
使用时将需要异步执行的业务逻辑封装为TransactionCallback列表,传入工具类执行即可,只要有一个任务执行异常,所有线程的事务都会统一回滚。
内容的提问来源于stack exchange,提问作者user2488578
相关产品推荐
相关产品推荐

