来源:https://urlify.cn/zYfEVb 1,最近有一个大数据量插入的宝多操作入库的业务场景,需要先做一些其他修改操作,线程然后在执行插入操作,事务说用由于插入数据可能会很多,回滚回去用到多线程去拆分数据并行处理来提高响应时间,等通如果有一个线程执行失败,支付知则全部回滚。宝多 2,线程在spring中可以使用@Transactional注解去控制事务,事务说用使出现异常时会进行回滚,回滚回去在多线程中,等通这个注解则不会生效,支付知如果主线程需要先执行一些修改数据库的站群服务器宝多操作,当子线程在进行处理出现异常时,线程主线程修改的数据则不会回滚,导致数据错误。 3,下面用一个简单示例演示多线程事务。 / * 平均拆分list方法. * @param n * @param */ ,int n){ List int remaider=source.size()%n; int number=source.size()/n; int offset=0;//偏移量 (int i=0;i List (remaider>0){ value=source.subList(i*number+offset, (i+1)*number+offset+1); remaider--; offset++; { value=source.subList(i*number+offset, (i+1)*number+offset); } result.add(value); } result; } /** 线程池配置 * @version V1.0 */ public class ExecutorConfig { private static int maxPoolSize = Runtime.getRuntime().availableProcessors(); private volatile static ExecutorService executorService; () { (executorService == null){ synchronized (ExecutorConfig.class){ (executorService == null){ executorService = newThreadPool(); } } } executorService; } (){ int queueSize = 500; int corePool = Math.min(5, maxPoolSize); new ThreadPoolExecutor(corePool, maxPoolSize, 10000L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue<>(queueSize),new ThreadPoolExecutor.AbortPolicy()); } (){ } } /** 获取sqlSession * @author 86182 * @version V1.0 */ @Component public class SqlContext { @Resource private SqlSessionTemplate sqlSessionTemplate; (){ SqlSessionFactory sqlSessionFactory = sqlSessionTemplate.getSqlSessionFactory(); sqlSessionFactory.openSession(); } } / * 测试多线程事务. * @param employeeDOList */ @Override @Transactional public void saveThread(List try { //先做删除操作,如果子线程出现异常,此操作不会回滚 this.getBaseMapper().delete(null); //获取线程池 ExecutorService service = ExecutorConfig.getThreadPool(); //拆分数据,拆分5份 List //执行的线程 Thread []threadArray = new Thread[lists.size()]; //监控子线程执行完毕,再执行主线程,要不然会导致主线程关闭,子线程也会随着关闭 CountDownLatch countDownLatch = new CountDownLatch(lists.size()); ); (int i =0;i (i==lists.size()-1){ ); } List threadArray[i] = new Thread(() -> { try { //最后一个线程抛出异常 (!atomicBoolean.get()){ ); } //批量添加,mybatisPlus中自带的batch方法 this.saveBatch(list); }finally { countDownLatch.countDown(); } }); } (int i = 0; i service.execute(threadArray[i]); } //当子线程执行完毕时,主线程再往下执行 countDownLatch.await(); ); }catch (Exception e){ ,e); ); }finally { connection.close(); } } 数据库中存在一条数据: //测试用例 @RunWith(SpringRunner.class) @SpringBootTest(classes = { ThreadTest01.class, MainApplication.class}) public class ThreadTest01 { @Resource private EmployeeBO employeeBO; / * 测试多线程事务. * @throws InterruptedException */ @Test public void MoreThreadTest2() throws InterruptedException { int size = 10; List (int i = 0; i EmployeeDO employeeDO = new EmployeeDO(); +i); employeeDO.setAge(18); employeeDO.setGender(1); ); employeeDO.setCreatTime(Calendar.getInstance().getTime()); employeeDOList.add(employeeDO); } try { employeeBO.saveThread(employeeDOList); ); }catch (Exception e){ e.printStackTrace(); } } } 测试结果: 可以发现子线程组执行时,有一个线程执行失败,其他线程也会抛出异常,但是主线程中执行的删除操作,服务器租用没有回滚,@Transactional注解没有生效。 使用sqlSession控制手动提交事务 @Resource SqlContext sqlContext; / * 测试多线程事务. * @param employeeDOList */ @Override public void saveThread(List // 获取数据库连接,获取会话(内部自有事务) SqlSession sqlSession = sqlContext.getSqlSession(); Connection connection = sqlSession.getConnection(); try { // 设置手动提交 ); //获取mapper EmployeeMapper employeeMapper = sqlSession.getMapper(EmployeeMapper.class); //先做删除操作 employeeMapper.delete(null); //获取执行器 ExecutorService service = ExecutorConfig.getThreadPool(); List //拆分list List ); (int i =0;i (i==lists.size()-1){ ); } List //使用返回结果的callable去执行, Callable //让最后一个线程抛出异常 (!atomicBoolean.get()){ ); } employeeMapper.saveBatch(list); }; callableList.add(callable); } //执行子线程 List (Future //如果有一个执行不成功,则全部回滚 (future.get()<=0){ connection.rollback(); ; } } connection.commit(); ); }catch (Exception e){ connection.rollback(); ,e); ); }finally { connection.close(); } } // sql > INSERT INTO employee (employee_id,age,employee_name,birth_date,gender,id_number,creat_time,update_time,status) values > ( ) 数据库中一条数据: 测试结果:抛出异常, 删除操作的数据回滚了,数据库中的数据依旧存在,说明事务成功了。 成功操作示例: @Resource SqlContext sqlContext; / * 测试多线程事务. * @param employeeDOList */ @Override public void saveThread(List // 获取数据库连接,获取会话(内部自有事务) SqlSession sqlSession = sqlContext.getSqlSession(); Connection connection = sqlSession.getConnection(); try { // 设置手动提交 ); EmployeeMapper employeeMapper = sqlSession.getMapper(EmployeeMapper.class); //先做删除操作 employeeMapper.delete(null); ExecutorService service = ExecutorConfig.getThreadPool(); List List (int i =0;i List Callable callableList.add(callable); } //执行子线程 List (Future (future.get()<=0){ connection.rollback(); ; } } connection.commit(); ); }catch (Exception e){ connection.rollback(); ,e); ); // throw new ServiceException(ExceptionCodeEnum.EMPLOYEE_SAVE_OR_UPDATE_ERROR); } } 测试结果: 数据库中数据: 删除的删除了,添加的添加成功了,测试成功。背景介绍
公用的类和方法
> result=new ArrayList
>();
示例事务不成功操作
> lists=averageAssign(employeeDOList, 5);
> lists=averageAssign(employeeDOList, 5);
> lists=averageAssign(employeeDOList, 5);