多线程事务怎么回滚?一个简单示例演示多线程事务

背景介绍
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多台一次消谐器用于南部电网农网改造工程