最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

SpringBoot用多線程批量導(dǎo)入數(shù)據(jù)庫實(shí)現(xiàn)方法

 更新時(shí)間:2023年02月03日 10:24:29   作者:愿做無知一猿  
這篇文章主要介紹了SpringBoot用多線程批量導(dǎo)入數(shù)據(jù)庫實(shí)現(xiàn)方法,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)吧

環(huán)境

springboot、mybatisPlus、mysql8

mysql8(部署在1核2G的服務(wù)器上,很卡,所以下面的數(shù)據(jù)條數(shù)用5000,太大怕不是要等到花兒都謝了 0.0)

原始的for循環(huán)入庫

@Service
@Slf4j
public class MoreTestServiceImpl extends ServiceImpl<MoreTestMapper, MoreTestEntity> implements MoreTestService {
    @Override
    @Transactional(rollbackFor = Exception.class)
    public Object doTest() {
        long start = System.currentTimeMillis();
        List<MoreTestEntity> entityList = new ArrayList<>();
        for (int i = 0; i < 5000; i++) {
            MoreTestEntity entity = new MoreTestEntity();
            entity.setId((long) i);
            entity.setA(UUID.randomUUID().toString());
            entity.setB(UUID.randomUUID().toString());
            entity.setC(UUID.randomUUID().toString());
            entity.setD(UUID.randomUUID().toString());
            entity.setE(UUID.randomUUID().toString());
            entity.setF(UUID.randomUUID().toString());
            entity.setG(UUID.randomUUID().toString());
            entity.setH(UUID.randomUUID().toString());
            entity.setI(UUID.randomUUID().toString());
            entity.setJ(UUID.randomUUID().toString());
            entity.setK(UUID.randomUUID().toString());
            entityList.add(entity);
            //在循環(huán)中入庫
            baseMapper.insert(entity);
        }
        long end = System.currentTimeMillis();
        System.err.println(end - start);
        return end - start;
    }
}

共耗時(shí):180121 ms

批量保存操作

@Service
@Slf4j
public class MoreTestServiceImpl extends ServiceImpl<MoreTestMapper, MoreTestEntity> implements MoreTestService {
    @Override
    @Transactional(rollbackFor = Exception.class)
    public Object doTest() {
        long start = System.currentTimeMillis();
        List<MoreTestEntity> entityList = new ArrayList<>();
        for (int i = 0; i < 5000; i++) {
            MoreTestEntity entity = new MoreTestEntity();
            entity.setId((long) i);
            entity.setA(UUID.randomUUID().toString());
            entity.setB(UUID.randomUUID().toString());
            entity.setC(UUID.randomUUID().toString());
            entity.setD(UUID.randomUUID().toString());
            entity.setE(UUID.randomUUID().toString());
            entity.setF(UUID.randomUUID().toString());
            entity.setG(UUID.randomUUID().toString());
            entity.setH(UUID.randomUUID().toString());
            entity.setI(UUID.randomUUID().toString());
            entity.setJ(UUID.randomUUID().toString());
            entity.setK(UUID.randomUUID().toString());
            entityList.add(entity);
        }
      	//mybatisPlus提供的批量保存方法,數(shù)字代表每幾條數(shù)據(jù)提交一次事務(wù),默認(rèn)1000
        saveBatch(entityList, 1000);
        long end = System.currentTimeMillis();
        System.err.println(end - start);
        return end - start;
    }
}

耗時(shí)時(shí)間:87217ms

在批量插入的基礎(chǔ)上使用多線程

@Service
@Slf4j
public class MoreTestServiceImpl extends ServiceImpl<MoreTestMapper, MoreTestEntity> implements MoreTestService {
    @Override
    @Transactional(rollbackFor = Exception.class)
    public Object doTest() throws InterruptedException {
        long start = System.currentTimeMillis();
        //手動(dòng)創(chuàng)建線程池,注意你 數(shù)據(jù)庫連接池的 允許連接數(shù)量,別超過了就行。
        ThreadPoolExecutor poolExecutor = new ThreadPoolExecutor(
                5,
                5,
                30,
                TimeUnit.SECONDS,
                new LinkedBlockingDeque<>(10),
                //isDaemon 設(shè)置線程是否是守護(hù)線程,true的話,主線程結(jié)束,new的線程就不會(huì)繼續(xù)工作
                new NamedThreadFactory("執(zhí)行線程", false),
                (r, executor) -> System.out.println("拒絕" + r));
        List<MoreTestEntity> entityList = new ArrayList<>();
        for (int i = 0; i < 5000; i++) {
            MoreTestEntity entity = new MoreTestEntity();
            entity.setId((long) i);
            entity.setA(UUID.randomUUID().toString());
            entity.setB(UUID.randomUUID().toString());
            entity.setC(UUID.randomUUID().toString());
            entity.setD(UUID.randomUUID().toString());
            entity.setE(UUID.randomUUID().toString());
            entity.setF(UUID.randomUUID().toString());
            entity.setG(UUID.randomUUID().toString());
            entity.setH(UUID.randomUUID().toString());
            entity.setI(UUID.randomUUID().toString());
            entity.setJ(UUID.randomUUID().toString());
            entity.setK(UUID.randomUUID().toString());
            entityList.add(entity);
        }
        //拆分list,將其拆分成5份,然后上面線程池創(chuàng)建也是5個(gè)核心線程,剛好執(zhí)行
        List<List<MoreTestEntity>> partition = ListUtils.partition(entityList, 1000);
        //使用CountDownLatch保證所有線程都執(zhí)行完成
        CountDownLatch latch = new CountDownLatch(5);
        partition.forEach(item -> {
            poolExecutor.execute(() -> {
                saveBatch(item, 1000);
                latch.countDown();
            });
        });
        latch.await();
        // 也可以這么寫,設(shè)定超時(shí)時(shí)間
        //latch.await(100,TimeUnit.SECONDS);
        long end = System.currentTimeMillis();
        System.err.println(end - start);
        //關(guān)閉線程池
        poolExecutor.shutdown();
        return end - start;
    }
}

耗時(shí)時(shí)間: 28235

可見時(shí)間從180秒,縮短到了28秒,但是@Transactional對(duì)于多線程是控制不了所有的事務(wù)的。

Spring實(shí)現(xiàn)事務(wù)的原理是通過ThreadLocal把數(shù)據(jù)庫連接綁定到當(dāng)前線程中,同一個(gè)事務(wù)中數(shù)據(jù)庫操作使用同一個(gè)jdbc connection,新開啟的線程獲取不到當(dāng)前jdbc connection。

如下代碼:

partition.forEach(item -> {
            poolExecutor.execute(() -> {
                saveBatch(item, 1000);
                latch.countDown();
                //讓每個(gè)都報(bào)錯(cuò)
                int i = 1/0;
            });
        });

控制臺(tái)打?。?/p>

Exception in thread "執(zhí)行線程5" java.lang.ArithmeticException: / by zero
    at com.kusch.ares.service.impl.MoreTestServiceImpl.lambda$null$1(MoreTestServiceImpl.java:68)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:750)
Exception in thread "執(zhí)行線程2" java.lang.ArithmeticException: / by zero
    at com.kusch.ares.service.impl.MoreTestServiceImpl.lambda$null$1(MoreTestServiceImpl.java:68)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:750)
Exception in thread "執(zhí)行線程4" java.lang.ArithmeticException: / by zero
    at com.kusch.ares.service.impl.MoreTestServiceImpl.lambda$null$1(MoreTestServiceImpl.java:68)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:750)
Exception in thread "執(zhí)行線程1" java.lang.ArithmeticException: / by zero
    at com.kusch.ares.service.impl.MoreTestServiceImpl.lambda$null$1(MoreTestServiceImpl.java:68)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:750)
Exception in thread "執(zhí)行線程3" 30179
java.lang.ArithmeticException: / by zero
    at com.kusch.ares.service.impl.MoreTestServiceImpl.lambda$null$1(MoreTestServiceImpl.java:68)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:750)

可見5個(gè)線程都報(bào)錯(cuò)了,但是去查詢數(shù)據(jù)庫,卻可以查詢到5000條數(shù)據(jù),這是不應(yīng)該出現(xiàn)的情況。

處理多線程入庫的事務(wù)問題

@Service
@Slf4j
public class MoreTestServiceImpl extends ServiceImpl<MoreTestMapper, MoreTestEntity> implements MoreTestService {
    @Resource
    private DataSourceTransactionManager dataSourceTransactionManager;
    @Resource
    private TransactionDefinition transactionDefinition;
    @Override
    //此處手動(dòng)管理事務(wù)的提交后,這個(gè)注解就可以去掉了
    //    @Transactional(rollbackFor = Exception.class)
    public Object doTest() {
        long start = System.currentTimeMillis();
        //手動(dòng)創(chuàng)建線程池,注意你 數(shù)據(jù)庫連接池的 允許連接數(shù)量,別超過了就行。
        ThreadPoolExecutor poolExecutor = new ThreadPoolExecutor(
                5,
                5,
                30,
                TimeUnit.SECONDS,
                new LinkedBlockingDeque<>(10),
                //isDaemon 設(shè)置線程是否是守護(hù)線程,true的話,主線程結(jié)束,new的線程就不會(huì)繼續(xù)工作
                new NamedThreadFactory("執(zhí)行線程", false),
                (r, executor) -> System.out.println("拒絕" + r));
        List<MoreTestEntity> entityList = new ArrayList<>();
        for (int i = 0; i < 50; i++) {
            MoreTestEntity entity = new MoreTestEntity();
            entity.setId((long) i);
            entity.setA(UUID.randomUUID().toString());
            entity.setB(UUID.randomUUID().toString());
            entity.setC(UUID.randomUUID().toString());
            entity.setD(UUID.randomUUID().toString());
            entity.setE(UUID.randomUUID().toString());
            entity.setF(UUID.randomUUID().toString());
            entity.setG(UUID.randomUUID().toString());
            entity.setH(UUID.randomUUID().toString());
            entity.setI(UUID.randomUUID().toString());
            entity.setJ(UUID.randomUUID().toString());
            entity.setK(UUID.randomUUID().toString());
            entityList.add(entity);
        }
        //拆分list,將其拆分成5份,然后上面線程池創(chuàng)建也是5個(gè)核心線程,剛好執(zhí)行
        List<List<MoreTestEntity>> partition = ListUtils.partition(entityList, 10);
        //使用CountDownLatch保證所有線程都執(zhí)行完成
        CountDownLatch sonLatch = new CountDownLatch(5);
        //主線程的 肯定為1
        CountDownLatch mainLatch = new CountDownLatch(1);
        AtomicBoolean hasError = new AtomicBoolean(false);
        partition.forEach(item -> {
            poolExecutor.execute(() -> {
                doSave(item, sonLatch, hasError, mainLatch);
            });
        });
        try {
            //此處應(yīng)該是用try catch 包裹著主線程的所有業(yè)務(wù)代碼,以此保證主線程中任何一處報(bào)錯(cuò)都可以通知子線程
            //這里加一個(gè)是為了調(diào)試主線程中的數(shù)據(jù)入庫操作
            MoreTestEntity entity = new MoreTestEntity();
            entity.setId((long) 99999);
            entity.setA(UUID.randomUUID().toString());
            entity.setB(UUID.randomUUID().toString());
            entity.setC(UUID.randomUUID().toString());
            entity.setD(UUID.randomUUID().toString());
            entity.setE(UUID.randomUUID().toString());
            entity.setF(UUID.randomUUID().toString());
            entity.setG(UUID.randomUUID().toString());
            entity.setH(UUID.randomUUID().toString());
            entity.setI(UUID.randomUUID().toString());
            entity.setJ(UUID.randomUUID().toString());
            entity.setK(UUID.randomUUID().toString());
            save(entity);
            //主線程報(bào)錯(cuò)
            int i = 10 / 0;
            sonLatch.await();
        } catch (InterruptedException e) {
            hasError.set(true);
            e.printStackTrace();
        }
        mainLatch.countDown();
        long end = System.currentTimeMillis();
        System.err.println(end - start);
        //關(guān)閉線程池
        if (!poolExecutor.isShutdown()) {
            poolExecutor.shutdown();
        }
        return end - start;
    }
    /**
     * 包裝后的子線程的保存代碼
     *
     * @param entityList 要保存的集合
     * @param sonLatch   子線程 CountDownLatch
     * @param hasError   是否發(fā)生錯(cuò)誤
     * @param mainLatch  主線程 CountDownLatch
     */
    private void doSave(List<MoreTestEntity> entityList,
                        CountDownLatch sonLatch,
                        AtomicBoolean hasError,
                        CountDownLatch mainLatch) {
        TransactionStatus transactionStatus = dataSourceTransactionManager.getTransaction(transactionDefinition);
        try {
            //            //子線程報(bào)錯(cuò)
            //            int i = 10 / 0;
            saveBatch(entityList);
        } catch (Throwable throwable) {
            throwable.printStackTrace();
            hasError.set(true);
        } finally {
            //這是必須的,每個(gè)子線程走完,要讓主線程繼續(xù)走,然后再回到子線程的每個(gè)任務(wù),決定是提交還是回滾
            sonLatch.countDown();
        }
        try {
            //等待主線程的執(zhí)行結(jié)束
            mainLatch.await();
        } catch (InterruptedException e) {
            e.printStackTrace();
            hasError.set(true);
        }
        //事務(wù)操作
        if (hasError.get()) {
            dataSourceTransactionManager.rollback(transactionStatus);
        } else {
            dataSourceTransactionManager.commit(transactionStatus);
        }
    }
}

分別放開子線程報(bào)錯(cuò)和主線程報(bào)錯(cuò),會(huì)發(fā)現(xiàn)事務(wù)都可以正?;貪L,達(dá)到了預(yù)期的效果。

主要思路就是通過子線程CountDownLatch和主線程CountDownLatch,控制線程好代碼的執(zhí)行順序即可。

最后補(bǔ)充幾點(diǎn):

  • 上述代碼中的countDown()一旦出現(xiàn)不執(zhí)行的情況那會(huì)導(dǎo)致線程堵塞堆積,所以建議給await()增加超時(shí)時(shí)間
  • 這樣操作可能還會(huì)出現(xiàn)問題,比如主線程通知子線程可以進(jìn)行實(shí)務(wù)操作了,但是各個(gè)子線程之間非透明,所以還是有幾率存在某個(gè)子線程事務(wù)回滾失敗的情況。

到此這篇關(guān)于SpringBoot用多線程批量導(dǎo)入數(shù)據(jù)庫實(shí)現(xiàn)方法的文章就介紹到這了,更多相關(guān)SpringBoot多線程導(dǎo)入數(shù)據(jù)庫內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Java實(shí)現(xiàn)計(jì)算器的代碼

    Java實(shí)現(xiàn)計(jì)算器的代碼

    這篇文章主要為大家介紹了Java實(shí)現(xiàn)計(jì)算器的詳細(xì)代碼,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2018-06-06
  • java多線程編程之管道通信詳解

    java多線程編程之管道通信詳解

    這篇文章主要為大家詳細(xì)介紹了java多線程編程之線程間的通信,探討使用管道進(jìn)行通信,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2017-10-10
  • Springboot Thymeleaf字符串對(duì)象實(shí)例解析

    Springboot Thymeleaf字符串對(duì)象實(shí)例解析

    這篇文章主要介紹了Springboot Thymeleaf字符串對(duì)象實(shí)例解析,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2007-09-09
  • Java調(diào)用高德地圖API根據(jù)詳細(xì)地址獲取經(jīng)緯度詳細(xì)教程

    Java調(diào)用高德地圖API根據(jù)詳細(xì)地址獲取經(jīng)緯度詳細(xì)教程

    寫了一個(gè)經(jīng)緯度相關(guān)的工具,分享給有需求的小伙伴們,下面這篇文章主要給大家介紹了關(guān)于Java調(diào)用高德地圖API根據(jù)詳細(xì)地址獲取經(jīng)緯度,文中通過圖文以及代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2024-04-04
  • 基于Springboot的高校社團(tuán)管理系統(tǒng)的設(shè)計(jì)與實(shí)現(xiàn)

    基于Springboot的高校社團(tuán)管理系統(tǒng)的設(shè)計(jì)與實(shí)現(xiàn)

    本文將基于Springboot+Mybatis開發(fā)實(shí)現(xiàn)一個(gè)高校社團(tuán)管理系統(tǒng),系統(tǒng)包含三個(gè)角色:管理員、團(tuán)長(zhǎng)、會(huì)員。文中采用的技術(shù)有Springboot、Mybatis、Jquery、AjAX、JSP等,感興趣的可以了解一下
    2022-07-07
  • 利用反射實(shí)現(xiàn)Excel和CSV 轉(zhuǎn)換為Java對(duì)象功能

    利用反射實(shí)現(xiàn)Excel和CSV 轉(zhuǎn)換為Java對(duì)象功能

    將Excel或CSV文件轉(zhuǎn)換為Java對(duì)象(POJO)以及將Java對(duì)象轉(zhuǎn)換為Excel或CSV文件可能是一個(gè)復(fù)雜的過程,但如果使用正確的工具和技術(shù),這個(gè)過程就會(huì)變得十分簡(jiǎn)單,在本文中,我們將了解如何利用一個(gè)Java反射的庫來實(shí)現(xiàn)這個(gè)功能,需要的朋友可以參考下
    2023-11-11
  • JAVA環(huán)境搭建之MyEclipse10+jdk1.8+tomcat8環(huán)境搭建詳解

    JAVA環(huán)境搭建之MyEclipse10+jdk1.8+tomcat8環(huán)境搭建詳解

    本文詳細(xì)講解了MyEclipse10+jdk1.8+tomcat8的JAVA環(huán)境搭建方法,希望能幫助到大家
    2018-10-10
  • Java switch關(guān)鍵字原理及用法詳解

    Java switch關(guān)鍵字原理及用法詳解

    這篇文章主要介紹了Java中 switch關(guān)鍵原理及用法詳解,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2019-11-11
  • 利用Java編寫一個(gè)Java虛擬機(jī)

    利用Java編寫一個(gè)Java虛擬機(jī)

    這篇文章主要為大家詳細(xì)介紹了如何使用 Java17 編寫的 Java 虛擬機(jī),文中的示例代碼講解詳細(xì),具有一定的學(xué)習(xí)價(jià)值,感興趣的可以了解下
    2023-07-07
  • Java?JDBC使用入門講解

    Java?JDBC使用入門講解

    JDBC是指Java數(shù)據(jù)庫連接,是一種標(biāo)準(zhǔn)Java應(yīng)用編程接口(?JAVA?API),用來連接?Java?編程語言和廣泛的數(shù)據(jù)庫。從根本上來說,JDBC?是一種規(guī)范,它提供了一套完整的接口,允許便攜式訪問到底層數(shù)據(jù)庫,本篇文章我們來了解MySQL連接JDBC的流程方法
    2022-12-12

最新評(píng)論

丰镇市| 安康市| 双峰县| 广灵县| 杂多县| 浦县| 克什克腾旗| 夏河县| 平阴县| 克拉玛依市| 仪陇县| 镇坪县| 旬阳县| 炉霍县| 天峨县| 大竹县| 保山市| 精河县| 楚雄市| 靖远县| 缙云县| 荥阳市| 柳州市| 巴中市| 满洲里市| 大宁县| 文化| 正蓝旗| 眉山市| 芦溪县| 兴城市| 和政县| 靖远县| 进贤县| 综艺| 阳山县| 铁岭市| 五指山市| 阳城县| 东兰县| 新竹市|