背景介绍
1,最近有一个大数据量插入的操作入库的业务场景,需要先做一些其他修改操作,然后在执行插入操作,由于插入数据可能会很多,用到多线程去拆分数据并行处理来提高响应时间,如果有一个线程执行失败,则全部回滚。
2,在spring中可以使用@transactional注解去控制事务,使出现异常时会进行回滚,在多线程中,这个注解则不会生效,如果主线程需要先执行一些修改数据库的操作,当子线程在进行处理出现异常时,主线程修改的数据则不会回滚,导致数据错误。
3,下面用一个简单示例演示多线程事务。
公用的类和方法
/** * 平均拆分list方法. * @param source * @param n * @param * @return */public static list averageassign(list source,int n){ list result=new arraylist(); int remaider=source.size()%n; int number=source.size()/n; int offset=0;//偏移量 for(int i=0;i0){ value=source.sublist(i*number+offset, (i+1)*number+offset+1); remaider--; offset++; }else{ value=source.sublist(i*number+offset, (i+1)*number+offset); } result.add(value); } return result;}/** 线程池配置 * @version v1.0 */public class executorconfig { private static int maxpoolsize = runtime.getruntime().availableprocessors(); private volatile static executorservice executorservice; public static executorservice getthreadpool() { if (executorservice == null){ synchronized (executorconfig.class){ if (executorservice == null){ executorservice = newthreadpool(); } } } return executorservice; } private static executorservice newthreadpool(){ int queuesize = 500; int corepool = math.min(5, maxpoolsize); return new threadpoolexecutor(corepool, maxpoolsize, 10000l, timeunit.milliseconds, new linkedblockingqueue(queuesize),new threadpoolexecutor.abortpolicy()); } private executorconfig(){}}/** 获取sqlsession * @author 86182 * @version v1.0 */@componentpublic class sqlcontext { @resource private sqlsessiontemplate sqlsessiontemplate; public sqlsession getsqlsession(){ sqlsessionfactory sqlsessionfactory = sqlsessiontemplate.getsqlsessionfactory(); return sqlsessionfactory.opensession(); }} 示例事务不成功操作
/** * 测试多线程事务. * @param employeedolist */@override@transactionalpublic void savethread(list employeedolist) { try { //先做删除操作,如果子线程出现异常,此操作不会回滚 this.getbasemapper().delete(null); //获取线程池 executorservice service = executorconfig.getthreadpool(); //拆分数据,拆分5份 list lists=averageassign(employeedolist, 5); //执行的线程 thread []threadarray = new thread[lists.size()]; //监控子线程执行完毕,再执行主线程,要不然会导致主线程关闭,子线程也会随着关闭 countdownlatch countdownlatch = new countdownlatch(lists.size()); atomicboolean atomicboolean = new atomicboolean(true); for (int i =0;i { try { //最后一个线程抛出异常 if (!atomicboolean.get()){ throw new serviceexception(001,出现异常); } //批量添加,mybatisplus中自带的batch方法 this.savebatch(list); }finally { countdownlatch.countdown(); } }); } for (int i = 0; i
//测试用例@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 employeedolist = new arraylist(size); for (int i = 0; i可以发现子线程组执行时,有一个线程执行失败,其他线程也会抛出异常,但是主线程中执行的删除操作,没有回滚,@transactional注解没有生效。
使用sqlsession控制手动提交事务
@resource sqlcontext sqlcontext; /** * 测试多线程事务. * @param employeedolist */@overridepublic void savethread(list employeedolist) throws sqlexception { // 获取数据库连接,获取会话(内部自有事务) sqlsession sqlsession = sqlcontext.getsqlsession(); connection connection = sqlsession.getconnection(); try { // 设置手动提交 connection.setautocommit(false); //获取mapper employeemapper employeemapper = sqlsession.getmapper(employeemapper.class); //先做删除操作 employeemapper.delete(null); //获取执行器 executorservice service = executorconfig.getthreadpool(); list callablelist = new arraylist(); //拆分list list lists=averageassign(employeedolist, 5); atomicboolean atomicboolean = new atomicboolean(true); for (int i =0;i { //让最后一个线程抛出异常 if (!atomicboolean.get()){ throw new serviceexception(001,出现异常); } return employeemapper.savebatch(list); }; callablelist.add(callable); } //执行子线程 list futures = service.invokeall(callablelist); for (future future:futures) { //如果有一个执行不成功,则全部回滚 if (future.get()<=0){ connection.rollback(); return; } } connection.commit(); system.out.println(添加完毕); }catch (exception e){ connection.rollback(); log.info(error,e); throw new serviceexception(002,出现异常); }finally { connection.close(); }}// sql insert into employee (employee_id,age,employee_name,birth_date,gender,id_number,creat_time,update_time,status) values ( #{item.employeeid}, #{item.age}, #{item.employeename}, #{item.birthdate}, #{item.gender}, #{item.idnumber}, #{item.creattime}, #{item.updatetime}, #{item.status} ) 数据库中一条数据:
测试结果:抛出异常,
删除操作的数据回滚了,数据库中的数据依旧存在,说明事务成功了。
成功操作示例:
@resourcesqlcontext sqlcontext;/** * 测试多线程事务. * @param employeedolist */@overridepublic void savethread(list employeedolist) throws sqlexception { // 获取数据库连接,获取会话(内部自有事务) sqlsession sqlsession = sqlcontext.getsqlsession(); connection connection = sqlsession.getconnection(); try { // 设置手动提交 connection.setautocommit(false); employeemapper employeemapper = sqlsession.getmapper(employeemapper.class); //先做删除操作 employeemapper.delete(null); executorservice service = executorconfig.getthreadpool(); list callablelist = new arraylist(); list lists=averageassign(employeedolist, 5); for (int i =0;i employeemapper.savebatch(list); callablelist.add(callable); } //执行子线程 list futures = service.invokeall(callablelist); for (future future:futures) { if (future.get()<=0){ connection.rollback(); return; } } connection.commit(); system.out.println(添加完毕); }catch (exception e){ connection.rollback(); log.info(error,e); throw new serviceexception(002,出现异常); // throw new serviceexception(exceptioncodeenum.employee_save_or_update_error); }} 测试结果:
数据库中数据:
删除的删除了,添加的添加成功了,测试成功。
Firefly-RK3128主板编译固件介绍
罗马仕充电宝 2020年改变你的充电体验
ST推出首款LED照明驱动器HVLED805
2SA2151和2SC6100设计的分立元件功放电路
自动铺床的被子 懒人的福利
多线程事务怎么回滚?一个简单示例演示多线程事务
性能更高更稳定!爱普特携手平头哥推进基于RISC-V的MCU生态发展
工业机器人的开发有什么是需要注意的
安富利推出Xilinx Virtex-6 FPGA DSP开
Xilinx Zynq-7000系列安全配置策略
整流器稳流和稳压的区别
奋达科技智能音箱业务出彩 成为国内领先的智能音箱ODM龙头
衡量电气绝缘性能的电气强度测试
今年市场的一大亮点:科技行业IPO IPO市场并购活动热度将保持两至三年
浅谈PLC控制系统设计要点及其在使用中的问题
飞利浦发布一款超宽曲面屏 CA屏+100Hz刷新率售价约合人民币4950元
华为P10最新消息:采用正面指纹方案 是无缺旗舰!
智源联合清华发布首个支持PyTorch框架的高性能MoE系统
MAXQ微控制器上的多路复用JTAG接口引脚
3000多台一次消谐器用于南部电网农网改造工程