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

詳解Java使用雙異步后如何保證數(shù)據(jù)一致性

 更新時(shí)間:2024年01月22日 10:20:19   作者:哪吒編程  
這篇文章主要為大家詳細(xì)介紹了Java使用雙異步后如何保證數(shù)據(jù)一致性,文中的示例代碼講解詳細(xì),具有一定的借鑒價(jià)值,有需要的小伙伴可以了解下

一、前情提要

在上一篇文章中,我們通過(guò)雙異步的方式導(dǎo)入了10萬(wàn)行的Excel,有個(gè)小伙伴在評(píng)論區(qū)問(wèn)我,如何保證插入后數(shù)據(jù)的一致性呢?

很簡(jiǎn)單,通過(guò)對(duì)比Excel文件行數(shù)和入庫(kù)數(shù)量是否相等即可。

那么,如何獲取異步線程的返回值呢?

二、通過(guò)Future獲取異步返回值

我們可以通過(guò)給異步方法添加Future返回值的方式獲取結(jié)果。

FutureTask 除了實(shí)現(xiàn) Future 接口外,還實(shí)現(xiàn)了 Runnable 接口。因此,F(xiàn)utureTask 可以交給 Executor 執(zhí)行,也可以由調(diào)用線程直接執(zhí)行FutureTask.run()。

1、FutureTask 是基于 AbstractQueuedSynchronizer實(shí)現(xiàn)的

AbstractQueuedSynchronizer簡(jiǎn)稱AQS,它是一個(gè)同步框架,它提供通用機(jī)制來(lái)原子性管理同步狀態(tài)、阻塞和喚醒線程,以及 維護(hù)被阻塞線程的隊(duì)列。 基于 AQS 實(shí)現(xiàn)的同步器包括: ReentrantLock、Semaphore、ReentrantReadWriteLock、 CountDownLatch 和 FutureTask。

基于 AQS實(shí)現(xiàn)的同步器包含兩種操作:

  • acquire,阻塞調(diào)用線程,直到AQS的狀態(tài)允許這個(gè)線程繼續(xù)執(zhí)行,在FutureTask中,get()就是這個(gè)方法;
  • release,改變AQS的狀態(tài),使state變?yōu)榉亲枞麪顟B(tài),在FutureTask中,可以通過(guò)run()和cancel()實(shí)現(xiàn)。

2、FutureTask執(zhí)行流程

執(zhí)行@Async異步方法;

建立新線程async-executor-X,執(zhí)行Runnable的run()方法,(FutureTask實(shí)現(xiàn)RunnableFuture,RunnableFuture實(shí)現(xiàn)Runnable);

判斷狀態(tài)state;

  • 如果未新建或者不處于AQS,直接返回;
  • 否則進(jìn)入COMPLETING狀態(tài),執(zhí)行異步線程代碼;

如果執(zhí)行cancel()方法改變AQS的狀態(tài)時(shí),會(huì)喚醒AQS等待隊(duì)列中的第一個(gè)線程線程async-executor-1;

線程async-executor-1被喚醒后

  • 將自己從AQS隊(duì)列中移除;
  • 然后喚醒next線程async-executor-2;
  • 改變線程async-executor-1的state;
  • 等待get()線程取值。

next等待線程被喚醒后,循環(huán)線程async-executor-1的步驟

  • 被喚醒
  • 從AQS隊(duì)列中移除
  • 喚醒next線程
  • 改變異步線程狀態(tài)

新建線程async-executor-N,監(jiān)聽(tīng)異步方法的state

  • 如果處于EXCEPTIONAL以上狀態(tài),拋出異常;
  • 如果處于COMPLETING狀態(tài),加入AQS隊(duì)列等待;
  • 如果處于NORMAL狀態(tài),返回結(jié)果;

3、get()方法執(zhí)行流程

get()方法通過(guò)判斷狀態(tài)state觀測(cè)異步線程是否已結(jié)束,如果結(jié)束直接將結(jié)果返回,否則會(huì)將等待節(jié)點(diǎn)扔進(jìn)等待隊(duì)列自旋,阻塞住線程。

自旋直至異步線程執(zhí)行完畢,獲取另一邊的線程計(jì)算出結(jié)果或取消后,將等待隊(duì)列里的所有節(jié)點(diǎn)依次喚醒并移除隊(duì)列。

  • 如果state小于等于COMPLETING,表示任務(wù)還在執(zhí)行中;
    • 如果線程被中斷,從等待隊(duì)列中移除等待節(jié)點(diǎn)WaitNode,拋出中斷異常;
    • 如果state大于COMPLETING;
      • 如果已有等待節(jié)點(diǎn)WaitNode,將線程置空;
      • 返回當(dāng)前狀態(tài);
    • 如果任務(wù)正在執(zhí)行,讓出時(shí)間片;
    • 如果還未構(gòu)造等待節(jié)點(diǎn),則new一個(gè)新的等待節(jié)點(diǎn);
    • 如果未入隊(duì)列,CAS嘗試入隊(duì);
    • 如果有超時(shí)時(shí)間參數(shù);
      • 計(jì)算超時(shí)時(shí)間;
      • 如果超時(shí),則從等待隊(duì)列中移除等待節(jié)點(diǎn)WaitNode,返回當(dāng)前狀態(tài)state;
      • 阻塞隊(duì)列nanos毫秒。
    • 否則阻塞隊(duì)列;
  • 如果state大于COMPLETING;
    • 如果執(zhí)行完畢,返回結(jié)果;
    • 如果大于等于取消狀態(tài),則拋出異常。

很多小朋友對(duì)讀源碼,嗤之以鼻,工作3年、5年,還是沒(méi)認(rèn)真讀過(guò)任何源碼,覺(jué)得讀了也沒(méi)啥用,或者讀了也看不懂~

其實(shí),只要把源碼的執(zhí)行流程通過(guò)畫(huà)圖的形式呈現(xiàn)出來(lái),你就會(huì)幡然醒悟,原來(lái)是這樣的~

簡(jiǎn)而言之:

1. 如果異步線程還沒(méi)執(zhí)行完,則進(jìn)入CAS自旋; 2. 其它線程獲取結(jié)果或取消后,重新喚醒CAS隊(duì)列中等待的線程; 3. 再通過(guò)get()判斷狀態(tài)state; 4. 直至返回結(jié)果或(取消、超時(shí)、異常)為止。

三、FutureTask源碼具體分析

1、FutureTask源碼

通過(guò)定義整形狀態(tài)值,判斷state大小,這個(gè)思想很有意思,值得學(xué)習(xí)。

public interface RunnableFuture<V> extends Runnable, Future<V> {
    /**
     * Sets this Future to the result of its computation
     * unless it has been cancelled.
     */
    void run();
}
public class FutureTask<V> implements RunnableFuture<V> {

	// 最初始的狀態(tài)是new 新建狀態(tài)
	private volatile int state;
    private static final int NEW          = 0; // 新建狀態(tài)
    private static final int COMPLETING   = 1; // 完成中
    private static final int NORMAL       = 2; // 正常執(zhí)行完
    private static final int EXCEPTIONAL  = 3; // 異常
    private static final int CANCELLED    = 4; // 取消
    private static final int INTERRUPTING = 5; // 正在中斷
    private static final int INTERRUPTED  = 6; // 已中斷

	public V get() throws InterruptedException, ExecutionException {
	    int s = state;
	    // 任務(wù)還在執(zhí)行中
	    if (s <= COMPLETING)
	        s = awaitDone(false, 0L);
	    return report(s);
	}
	
	private int awaitDone(boolean timed, long nanos)
        throws InterruptedException {
        final long deadline = timed ? System.nanoTime() + nanos : 0L;
        WaitNode q = null;
        boolean queued = false;
        for (;;) {
        	// 線程被中斷,從等待隊(duì)列中移除等待節(jié)點(diǎn)WaitNode,拋出中斷異常
            if (Thread.interrupted()) {
                removeWaiter(q);
                throw new InterruptedException();
            }

            int s = state;
            // 任務(wù)已執(zhí)行完畢或取消
            if (s > COMPLETING) {
            	// 如果已有等待節(jié)點(diǎn)WaitNode,將線程置空
                if (q != null)
                    q.thread = null;
                return s;
            }
            // 任務(wù)正在執(zhí)行,讓出時(shí)間片
            else if (s == COMPLETING) // cannot time out yet
                Thread.yield();
            // 還未構(gòu)造等待節(jié)點(diǎn),則new一個(gè)新的等待節(jié)點(diǎn)
            else if (q == null)
                q = new WaitNode();
            // 未入隊(duì)列,CAS嘗試入隊(duì)
            else if (!queued)
                queued = UNSAFE.compareAndSwapObject(this, waitersOffset,
                                                     q.next = waiters, q);
            // 如果有超時(shí)時(shí)間參數(shù)
            else if (timed) {
            	// 計(jì)算超時(shí)時(shí)間
                nanos = deadline - System.nanoTime();
                // 如果超時(shí),則從等待隊(duì)列中移除等待節(jié)點(diǎn)WaitNode,返回當(dāng)前狀態(tài)state
                if (nanos <= 0L) {
                    removeWaiter(q);
                    return state;
                }
                // 阻塞隊(duì)列nanos毫秒
                LockSupport.parkNanos(this, nanos);
            }
            else
            	// 阻塞隊(duì)列
                LockSupport.park(this);
        }
    }
    
	private V report(int s) throws ExecutionException {
		// 獲取outcome中記錄的返回結(jié)果
        Object x = outcome;
        // 如果執(zhí)行完畢,返回結(jié)果
        if (s == NORMAL)
            return (V)x;
            // 如果大于等于取消狀態(tài),則拋出異常
        if (s >= CANCELLED)
            throw new CancellationException();
        throw new ExecutionException((Throwable)x);
    }
}

2、將異步方法的返回值改為Future<Integer>,將返回值放到new AsyncResult<>();中;

@Async("async-executor")
public void readXls(String filePath, String filename) {
    try {
    	// 此代碼為簡(jiǎn)化關(guān)鍵性代碼
        List<Future<Integer>> futureList = new ArrayList<>();
        for (int time = 0; time < times; time++) {
            Future<Integer> sumFuture = readExcelDataAsyncFutureService.readXlsCacheAsync();
            futureList.add(sumFuture);
        }
    }catch (Exception e){
        logger.error("readXlsCacheAsync---插入數(shù)據(jù)異常:",e);
    }
}
@Async("async-executor")
public Future<Integer> readXlsCacheAsync() {
    try {
        // 此代碼為簡(jiǎn)化關(guān)鍵性代碼
        return new AsyncResult<>(sum);
    }catch (Exception e){
        return new AsyncResult<>(0);
    }
}

3、通過(guò)Future<Integer>.get()獲取返回值:

public static boolean getFutureResult(List<Future<Integer>> futureList, int excelRow){
    int[] futureSumArr = new int[futureList.size()];
    for (int i = 0;i<futureList.size();i++) {
        try {
            Future<Integer> future = futureList.get(i);
            while (true) {
                if (future.isDone() && !future.isCancelled()) {
                    Integer futureSum = future.get();
                    logger.info("獲取Future返回值成功"+"----Future:" + future
                            + ",Result:" + futureSum);
                    futureSumArr[i] += futureSum;
                    break;
                } else {
                    logger.info("Future正在執(zhí)行---獲取Future返回值中---等待3秒");
                    Thread.sleep(3000);
                }
            }
        } catch (Exception e) {
            logger.error("獲取Future返回值異常: ", e);
        }
    }
    
    boolean insertFlag = getInsertSum(futureSumArr, excelRow);
    logger.info("獲取所有異步線程Future的返回值成功,Excel插入結(jié)果="+insertFlag);
    return insertFlag;
}

4、這里也可以通過(guò)新線程+Future獲取Future返回值

不過(guò)感覺(jué)多此一舉了,就當(dāng)練習(xí)Future異步取返回值了~

public static Future<Boolean> getFutureResultThreadFuture(List<Future<Integer>> futureList, int excelRow) {
    ExecutorService service = Executors.newSingleThreadExecutor();
    final boolean[] insertFlag = {false};
    service.execute(new Runnable() {
        public void run() {
            try {
                insertFlag[0] = getFutureResult(futureList, excelRow);
            } catch (Exception e) {
                logger.error("新線程+Future獲取Future返回值異常: ", e);
                insertFlag[0] = false;
            }
        }
    });
    service.shutdown();
    return new AsyncResult<>(insertFlag[0]);
}

獲取異步線程結(jié)果后,我們可以通過(guò)添加事務(wù)的方式,實(shí)現(xiàn)Excel入庫(kù)操作的數(shù)據(jù)一致性。

但Future會(huì)造成主線程的阻塞,這個(gè)就很不友好了,有沒(méi)有更優(yōu)解呢?

以上就是詳解Java使用雙異步后如何保證數(shù)據(jù)一致性的詳細(xì)內(nèi)容,更多關(guān)于Java雙異步如何保證數(shù)據(jù)一致性的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • Eclipse 2020-06 漢化包安裝步驟詳解(附漢化包+安裝教程)

    Eclipse 2020-06 漢化包安裝步驟詳解(附漢化包+安裝教程)

    這篇文章主要介紹了Eclipse 2020-06 漢化包安裝步驟(附漢化包+安裝教程),本文通過(guò)圖文并茂的形式給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2020-08-08
  • SpringBoot請(qǐng)求映射的五種優(yōu)化方式小結(jié)

    SpringBoot請(qǐng)求映射的五種優(yōu)化方式小結(jié)

    在Spring?Boot應(yīng)用開(kāi)發(fā)中,請(qǐng)求映射(Request?Mapping)是將HTTP請(qǐng)求路由到相應(yīng)控制器方法的核心機(jī)制,合理優(yōu)化請(qǐng)求映射不僅可以提升應(yīng)用性能,還能改善代碼結(jié)構(gòu),增強(qiáng)API的可維護(hù)性和可擴(kuò)展性,本文將介紹5種Spring?Boot請(qǐng)求映射優(yōu)化方式,需要的朋友可以參考下
    2025-06-06
  • 通過(guò)mybatis-plus進(jìn)行數(shù)據(jù)庫(kù)字段加解密方式

    通過(guò)mybatis-plus進(jìn)行數(shù)據(jù)庫(kù)字段加解密方式

    文章主要介紹了在Java開(kāi)發(fā)中,從編寫(xiě)處理程序(handler)到實(shí)現(xiàn)加解密工具(util),再到配置實(shí)體和字段,以及自定義MyBatis的mapper語(yǔ)句的全過(guò)程
    2026-01-01
  • Java實(shí)現(xiàn)遞歸刪除菜單和目錄及目錄下所有文件

    Java實(shí)現(xiàn)遞歸刪除菜單和目錄及目錄下所有文件

    這篇文章主要為大家詳細(xì)介紹了Java如何實(shí)現(xiàn)遞歸刪除菜單和刪除目錄及目錄下所有文件,文中的示例代碼講解詳細(xì),感興趣的小伙伴可以參考一下
    2025-03-03
  • springboot集成ftp實(shí)現(xiàn)文件上傳

    springboot集成ftp實(shí)現(xiàn)文件上傳

    這篇文章主要為大家詳細(xì)介紹了springboot集成ftp實(shí)現(xiàn)文件上傳,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2021-05-05
  • Java并發(fā) 結(jié)合源碼分析AQS原理

    Java并發(fā) 結(jié)合源碼分析AQS原理

    這篇文章主要介紹了Java并發(fā) 結(jié)合源碼分析AQS原理,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2019-10-10
  • maven打包時(shí)候修改包名稱帶上git版本號(hào)和打包時(shí)間方式

    maven打包時(shí)候修改包名稱帶上git版本號(hào)和打包時(shí)間方式

    這篇文章主要介紹了maven打包時(shí)候修改包名稱帶上git版本號(hào)和打包時(shí)間方式,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2023-04-04
  • 詳解Spring的@Value作用與使用場(chǎng)景

    詳解Spring的@Value作用與使用場(chǎng)景

    這篇文章主要介紹了詳解Spring的@Value作用與使用場(chǎng)景,Spring為大家提供許多開(kāi)箱即用的功能,@Value就是一個(gè)極其常用的功能,它能將配置信息注入到bean中去,需要的朋友可以參考下
    2023-05-05
  • 8張圖帶你全面了解Java?kafka的核心機(jī)制

    8張圖帶你全面了解Java?kafka的核心機(jī)制

    kafka是目前企業(yè)中很常用的消息隊(duì)列產(chǎn)品,可以用于削峰、解耦、異步通信,本文就通過(guò)幾張圖帶大家全面認(rèn)識(shí)一下kafka,現(xiàn)在我們不妨帶入kafka設(shè)計(jì)者的角度去思考該如何設(shè)計(jì),它的架構(gòu)是怎么樣的、都有哪些組件組成、如何進(jìn)行擴(kuò)展等等,需要的朋友可以參考下
    2023-05-05
  • Java實(shí)現(xiàn)快速排序算法可視化的示例代碼

    Java實(shí)現(xiàn)快速排序算法可視化的示例代碼

    快速排序算法通過(guò)多次比較和交換來(lái)實(shí)現(xiàn)排序,是對(duì)冒泡排序算法的一種改進(jìn)。本文將用Java語(yǔ)言實(shí)現(xiàn)快速排序算法并進(jìn)行可視化,感興趣的可以了解一下
    2022-08-08

最新評(píng)論

黑水县| 东台市| 新安县| 青冈县| 松潘县| 文山县| 阿城市| 大连市| 滁州市| 华蓥市| 苗栗县| 通江县| 荣成市| 梨树县| 白山市| 丹寨县| 绵阳市| 方正县| 岗巴县| 晴隆县| 邵阳县| 大余县| 邹城市| 青田县| 普兰店市| 台安县| 札达县| 绩溪县| 满洲里市| 西乌| 中西区| 岑溪市| 朝阳区| 金阳县| 马鞍山市| 类乌齐县| 镇江市| 独山县| 同心县| 合水县| 故城县|