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

淺析Java中并發(fā)工具類(lèi)的使用

 更新時(shí)間:2022年12月07日 10:49:48   作者:初念初戀  
在JDK的并發(fā)包里提供了幾個(gè)非常有用的并發(fā)工具類(lèi)。CountDownLatch、CyclicBarrier和Semaphore工具類(lèi)提供了一種并發(fā)流程控制的手段,Exchanger工具類(lèi)提供了在線(xiàn)程間交換數(shù)據(jù)的一種方法。本文主要介紹了它們的使用,需要的可以參考一下

在JDK的并發(fā)包里提供了幾個(gè)非常有用的并發(fā)工具類(lèi)。CountDownLatch、CyclicBarrier和Semaphore工具類(lèi)提供了一種并發(fā)流程控制的手段,Exchanger工具類(lèi)提供了在線(xiàn)程間交換數(shù)據(jù)的一種方法。

它們都在java.util.concurrent包下。先總體概括一下都有哪些工具類(lèi),它們有什么作用,然后再分別介紹它們的主要使用方法和原理。

類(lèi)作用
CountDownLatch線(xiàn)程等待直到計(jì)數(shù)器減為0時(shí)開(kāi)始工作
CyclicBarrier作用跟CountDownLatch類(lèi)似,但是可以重復(fù)使用
Semaphore限制線(xiàn)程的數(shù)量
Exchanger兩個(gè)線(xiàn)程交換數(shù)據(jù)

下面分別介紹這幾個(gè)類(lèi)。

CountDownLatch

概述

CountDownLatch可以使一個(gè)或多個(gè)線(xiàn)程等待其他線(xiàn)程各自執(zhí)行完畢后再執(zhí)行。

CountDownLatch定義了一個(gè)計(jì)數(shù)器,和一個(gè)阻塞隊(duì)列, 當(dāng)計(jì)數(shù)器的值遞減為0之前,阻塞隊(duì)列里面的線(xiàn)程處于掛起狀態(tài),當(dāng)計(jì)數(shù)器遞減到0時(shí)會(huì)喚醒阻塞隊(duì)列所有線(xiàn)程,這里的計(jì)數(shù)器是一個(gè)標(biāo)志,可以表示一個(gè)任務(wù)一個(gè)線(xiàn)程,也可以表示一個(gè)倒計(jì)時(shí)器。

案例

玩吃雞游戲的時(shí)候,正式開(kāi)始游戲之前,肯定會(huì)加載一些前置場(chǎng)景,例如:“加載地圖”、“加載人物模型”、“加載背景音樂(lè)”等。

public class CountDownLatchDemo {
    // 定義前置任務(wù)線(xiàn)程
    static class PreTaskThread implements Runnable {

        private String task;
        private CountDownLatch countDownLatch;

        public PreTaskThread(String task, CountDownLatch countDownLatch) {
            this.task = task;
            this.countDownLatch = countDownLatch;
        }

        @Override
        public void run() {
            try {
                Random random = new Random();
                Thread.sleep(random.nextInt(1000));
                System.out.println(task + " - 任務(wù)完成");
                countDownLatch.countDown();
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }
    }

    public static void main(String[] args) {
        // 假設(shè)有三個(gè)模塊需要加載
        CountDownLatch countDownLatch = new CountDownLatch(3);

        // 主任務(wù)
        new Thread(() -> {
            try {
                System.out.println("等待數(shù)據(jù)加載...");
                System.out.println(String.format("還有%d個(gè)前置任務(wù)", countDownLatch.getCount()));
                countDownLatch.await();
                System.out.println("數(shù)據(jù)加載完成,正式開(kāi)始游戲!");
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }).start();

        // 前置任務(wù)
        new Thread(new PreTaskThread("加載地圖數(shù)據(jù)", countDownLatch)).start();
        new Thread(new PreTaskThread("加載人物模型", countDownLatch)).start();
        new Thread(new PreTaskThread("加載背景音樂(lè)", countDownLatch)).start();
    }
}

輸出:

等待數(shù)據(jù)加載...
還有3個(gè)前置任務(wù)
加載地圖數(shù)據(jù) - 任務(wù)完成
加載人物模型 - 任務(wù)完成
加載背景音樂(lè) - 任務(wù)完成
數(shù)據(jù)加載完成,正式開(kāi)始游戲!

原理

CountDownLatch的方法很簡(jiǎn)單,如下:

// 構(gòu)造方法:
public CountDownLatch(int count)

public void await() // 等待
public boolean await(long timeout, TimeUnit unit) // 超時(shí)等待
public void countDown() // count - 1
public long getCount() // 獲取當(dāng)前還有多少count

CountDownLatch構(gòu)造器中的計(jì)數(shù)值(count)實(shí)際上就是閉鎖需要等待的線(xiàn)程數(shù)量。這個(gè)值只能被設(shè)置一次,而且CountDownLatch沒(méi)有提供任何機(jī)制去重新設(shè)置這個(gè)計(jì)數(shù)值。

與CountDownLatch的第一次交互是主線(xiàn)程等待其他線(xiàn)程。主線(xiàn)程必須在啟動(dòng)其他線(xiàn)程后立即調(diào)用CountDownLatch.await()方法。這樣主線(xiàn)程的操作就會(huì)在這個(gè)方法上阻塞,直到其他線(xiàn)程完成各自的任務(wù)。

其他N 個(gè)線(xiàn)程必須引用閉鎖對(duì)象,因?yàn)樗麄冃枰ㄖ狢ountDownLatch對(duì)象,他們已經(jīng)完成了各自的任務(wù)。這種通知機(jī)制是通過(guò)CountDownLatch.countDown()方法來(lái)完成的;每調(diào)用一次這個(gè)方法,在構(gòu)造函數(shù)中初始化的count值就減1。所以當(dāng)N個(gè)線(xiàn)程都調(diào) 用了這個(gè)方法,count的值等于0,然后主線(xiàn)程就能通過(guò)await()方法,恢復(fù)執(zhí)行自己的任務(wù)。

源碼分析

CountDownLatch有一個(gè)內(nèi)部類(lèi)叫做Sync,它繼承了AbstractQueuedSynchronizer類(lèi),其中維護(hù)了一個(gè)整數(shù)state,并且保證了修改state的可見(jiàn)性和原子性,源碼如下:

private static final class Sync extends AbstractQueuedSynchronizer {
        private static final long serialVersionUID = 4982264981922014374L;

        Sync(int count) {
            setState(count);
        }

        int getCount() {
            return getState();
        }

        protected int tryAcquireShared(int acquires) {
            return (getState() == 0) ? 1 : -1;
        }

        protected boolean tryReleaseShared(int releases) {
            // Decrement count; signal when transition to zero
            for (;;) {
                int c = getState();
                if (c == 0)
                    return false;
                int nextc = c - 1;
                if (compareAndSetState(c, nextc))
                    return nextc == 0;
            }
        }
    }

創(chuàng)建CountDownLatch實(shí)例時(shí),也會(huì)創(chuàng)建一個(gè)Sync的實(shí)例,同時(shí)把計(jì)數(shù)器的值傳給Sync實(shí)例,源碼如下:

public CountDownLatch(int count) {
  if (count < 0) throw new IllegalArgumentException("count < 0");
  this.sync = new Sync(count);
}

countDown方法中,只調(diào)用了Sync實(shí)例的releaseShared方法,源碼如下:

public void countDown() {
    sync.releaseShared(1);
}

其中的releaseShared方法,先對(duì)計(jì)數(shù)器進(jìn)行減1操作,如果減1后的計(jì)數(shù)器為0,喚醒被await方法阻塞的所有線(xiàn)程,源碼如下:

public final boolean releaseShared(int arg) {
    if (tryReleaseShared(arg)) { //對(duì)計(jì)數(shù)器進(jìn)行減一操作
        doReleaseShared();//如果計(jì)數(shù)器為0,喚醒被await方法阻塞的所有線(xiàn)程
        return true;
    }
    return false;
}

其中的tryReleaseShared方法,先獲取當(dāng)前計(jì)數(shù)器的值,如果計(jì)數(shù)器為0時(shí),就直接返回;如果不為0時(shí),使用CAS方法對(duì)計(jì)數(shù)器進(jìn)行減1操作,源碼如下:

protected boolean tryReleaseShared(int releases) {
    for (;;) {//死循環(huán),如果CAS操作失敗就會(huì)不斷繼續(xù)嘗試。
        int c = getState();//獲取當(dāng)前計(jì)數(shù)器的值。
        if (c == 0)// 計(jì)數(shù)器為0時(shí),就直接返回。
            return false;
        int nextc = c-1;
        if (compareAndSetState(c, nextc))// 使用CAS方法對(duì)計(jì)數(shù)器進(jìn)行減1操作
            return nextc == 0;//如果操作成功,返回計(jì)數(shù)器是否為0
    }
}

await方法中,只調(diào)用了Sync實(shí)例的acquireSharedInterruptibly方法,源碼如下:

public void await() throws InterruptedException {
    sync.acquireSharedInterruptibly(1);
}

其中acquireSharedInterruptibly方法,判斷計(jì)數(shù)器是否為0,如果不為0則阻塞當(dāng)前線(xiàn)程,源碼如下:

public final void acquireSharedInterruptibly(int arg)
        throws InterruptedException {
    if (Thread.interrupted())
        throw new InterruptedException();
    if (tryAcquireShared(arg) < 0)//判斷計(jì)數(shù)器是否為0
        doAcquireSharedInterruptibly(arg);//如果不為0則阻塞當(dāng)前線(xiàn)程
}

其中tryAcquireShared方法,是AbstractQueuedSynchronizer中的一個(gè)模板方法,其具體實(shí)現(xiàn)在Sync類(lèi)中,其主要是判斷計(jì)數(shù)器是否為零,如果為零則返回1,如果不為零則返回-1,源碼如下:

protected int tryAcquireShared(int acquires) {
    return (getState() == 0) ? 1 : -1;
}

CyclicBarrier

概述

CyclicBarrier 翻譯為中文是循環(huán)(Cyclic)柵欄(Barrier)的意思,它的大概含義是實(shí)現(xiàn)一個(gè)可循環(huán)利用的屏障。

CyclicBarrier 作用是讓一組線(xiàn)程相互等待,當(dāng)達(dá)到一個(gè)共同點(diǎn)時(shí),所有之前等待的線(xiàn)程再繼續(xù)執(zhí)行,且 CyclicBarrier 功能可重復(fù)使用,使用reset()方法重置屏障。

案例

同樣用玩游戲的例子。如果玩一個(gè)游戲有多個(gè)“關(guān)卡”,那使用CountDownLatch顯然不太合適,那需要為每個(gè)關(guān)卡都創(chuàng)建一個(gè)實(shí)例。那我們可以使用CyclicBarrier來(lái)實(shí)現(xiàn)每個(gè)關(guān)卡的數(shù)據(jù)加載等待功能。

public class CyclicBarrierDemo {
    static class PreTaskThread implements Runnable {

        private String task;
        private CyclicBarrier cyclicBarrier;

        public PreTaskThread(String task, CyclicBarrier cyclicBarrier) {
            this.task = task;
            this.cyclicBarrier = cyclicBarrier;
        }

        @Override
        public void run() {
            // 假設(shè)總共三個(gè)關(guān)卡
            for (int i = 1; i < 4; i++) {
                try {
                    Random random = new Random();
                    Thread.sleep(random.nextInt(1000));
                    System.out.println(String.format("關(guān)卡%d的任務(wù)%s完成", i, task));
                    cyclicBarrier.await();
                } catch (InterruptedException | BrokenBarrierException e) {
                    e.printStackTrace();
                }
            }
        }
    }

    public static void main(String[] args) {
        CyclicBarrier cyclicBarrier = new CyclicBarrier(3, () -> {
            System.out.println("本關(guān)卡所有前置任務(wù)完成,開(kāi)始游戲...");
        });

        new Thread(new PreTaskThread("加載地圖數(shù)據(jù)", cyclicBarrier)).start();
        new Thread(new PreTaskThread("加載人物模型", cyclicBarrier)).start();
        new Thread(new PreTaskThread("加載背景音樂(lè)", cyclicBarrier)).start();
    }
}

輸出:

關(guān)卡1的任務(wù)加載背景音樂(lè)完成
關(guān)卡1的任務(wù)加載地圖數(shù)據(jù)完成
關(guān)卡1的任務(wù)加載人物模型完成
本關(guān)卡所有前置任務(wù)完成,開(kāi)始游戲...
關(guān)卡2的任務(wù)加載人物模型完成
關(guān)卡2的任務(wù)加載背景音樂(lè)完成
關(guān)卡2的任務(wù)加載地圖數(shù)據(jù)完成
本關(guān)卡所有前置任務(wù)完成,開(kāi)始游戲...
關(guān)卡3的任務(wù)加載背景音樂(lè)完成
關(guān)卡3的任務(wù)加載地圖數(shù)據(jù)完成
關(guān)卡3的任務(wù)加載人物模型完成
本關(guān)卡所有前置任務(wù)完成,開(kāi)始游戲...

與CountDownLatch有一些不同。CyclicBarrier沒(méi)有分為await()countDown(),而是只有單獨(dú)的一個(gè)await()方法。

一旦調(diào)用await()方法的線(xiàn)程數(shù)量等于構(gòu)造方法中傳入的任務(wù)總量,就代表達(dá)到屏障了。CyclicBarrier允許我們?cè)谶_(dá)到屏障的時(shí)候可以執(zhí)行一個(gè)任務(wù),可以在構(gòu)造方法傳入一個(gè)Runnable類(lèi)型的對(duì)象。

源碼分析

構(gòu)造函數(shù):

public CyclicBarrier(int parties, Runnable barrierAction) {
  if (parties <= 0) throw new IllegalArgumentException();
  this.parties = parties;
  this.count = parties;
  this.barrierCommand = barrierAction;
}

public CyclicBarrier(int parties) {
  this(parties, null);
}

默認(rèn)barrierAction是null,這個(gè)參數(shù)是Runnable參數(shù),當(dāng)最后線(xiàn)程達(dá)到的時(shí)候執(zhí)行的任務(wù),上述案例就是在達(dá)到屏障時(shí),輸出“本關(guān)卡所有前置任務(wù)完成,開(kāi)始游戲...”。parties 是參與的線(xiàn)程數(shù)。

接著看下await方法,有兩個(gè)重載,區(qū)別是是否有等待超時(shí),源碼如下:

 public int await() throws InterruptedException, BrokenBarrierException {
        try {
            return dowait(false, 0L);
        } catch (TimeoutException toe) {
            throw new Error(toe); // cannot happen
        }
    }

public int await(long timeout, TimeUnit unit)
    throws InterruptedException,
           BrokenBarrierException,
           TimeoutException {
    return dowait(true, unit.toNanos(timeout));
}

重點(diǎn)看下dowait(),核心邏輯就是這個(gè)方法,源碼如下:

private int dowait(boolean timed, long nanos)
        throws InterruptedException, BrokenBarrierException,
               TimeoutException {       
        final ReentrantLock lock = this.lock;
        lock.lock();
        try {
            // 每次使用屏障都會(huì)生成一個(gè)實(shí)例
            final Generation g = generation;

            // 如果被破壞了就拋異常
            if (g.broken)
                throw new BrokenBarrierException();

            // 線(xiàn)程中斷檢測(cè)
            if (Thread.interrupted()) {
                breakBarrier();
                throw new InterruptedException();
            }

            // 剩余的等待線(xiàn)程數(shù)
            int index = --count;
            // 最后線(xiàn)程到達(dá)時(shí) 
            if (index == 0) {  // tripped
                // 標(biāo)記任務(wù)是否被執(zhí)行(就是傳進(jìn)入的runable參數(shù))
                boolean ranAction = false;
                try {
                    final Runnable command = barrierCommand;
                    // 執(zhí)行任務(wù)
                    if (command != null)
                        command.run();
                    ranAction = true;
                    // 完成后 進(jìn)行下一組 初始化 generation 初始化 count 并喚醒所有等待的線(xiàn)程 
                    nextGeneration();
                    return 0;
                } finally {
                    if (!ranAction)
                        breakBarrier();
                }
            }

            // index 不為0時(shí) 進(jìn)入自旋
            for (;;) {
                try {
                    // 先判斷超時(shí) 沒(méi)超時(shí)就繼續(xù)等著
                    if (!timed)
                        trip.await();
                        // 如果超出指定時(shí)間 調(diào)用 awaitNanos 超時(shí)了釋放鎖
                    else if (nanos > 0L)
                        nanos = trip.awaitNanos(nanos);
                        // 中斷異常捕獲
                } catch (InterruptedException ie) {
                    // 判斷是否被破壞
                    if (g == generation && ! g.broken) {
                        breakBarrier();
                        throw ie;
                    } else {
                        // 否則的話(huà)中斷當(dāng)前線(xiàn)程
                        Thread.currentThread().interrupt();
                    }
                }

                // 被破壞拋異常
                if (g.broken)
                    throw new BrokenBarrierException();

                // 正常調(diào)用 就返回 
                if (g != generation)
                    return index;

                // 超時(shí)了而被喚醒的情況 調(diào)用 breakBarrier()
                if (timed && nanos <= 0L) {
                    breakBarrier();
                    throw new TimeoutException();
                }
            }
        } finally {
            lock.unlock();
        }
    }

總結(jié)下dowait()方法的邏輯:

  • 線(xiàn)程調(diào)用后,會(huì)檢查barrier的狀態(tài)、線(xiàn)程狀態(tài),異常狀態(tài)會(huì)中斷。
  • 在初始化CyclicBarrier時(shí),設(shè)置的資源值count,會(huì)進(jìn)行--count
  • 當(dāng)10個(gè)線(xiàn)程中前9個(gè)線(xiàn)程,執(zhí)行dowait()后,由于count!=0,因此會(huì)進(jìn)行for(;;),在內(nèi)部會(huì)執(zhí)行Condition的trip.await()方法,進(jìn)行阻塞。
  • 阻塞結(jié)束的條件有:超時(shí)、被喚醒、線(xiàn)程中斷。
  • 當(dāng)?shù)?0個(gè)線(xiàn)程執(zhí)行dowait()后,由于count==0,會(huì)先檢查并執(zhí)行command的內(nèi)容。
  • 最后執(zhí)行nextGeneration(),在內(nèi)部調(diào)用trip.signalAll()喚醒所有trip.await()的線(xiàn)程。

如果被破壞了怎么恢復(fù)呢?來(lái)看下reset()方法,源碼如下:

public void reset() {
    final ReentrantLock lock = this.lock;
    lock.lock();
    try {
        breakBarrier();   // break the current generation
        nextGeneration(); // start a new generation
    } finally {
        lock.unlock();
    }
}

源碼很簡(jiǎn)單,break之后重新生成新的實(shí)例,對(duì)應(yīng)的會(huì)重新初始化count,在dowaitindex==0也調(diào)用了nextGeneration,所以說(shuō)它是可以循環(huán)利用的。

與CountDonwLatch的區(qū)別

CountDownLatch減計(jì)數(shù),CyclicBarrier加計(jì)數(shù)。

CountDownLatch是一次性的,CyclicBarrier可以重用。

CountDownLatch和CyclicBarrier都有讓多個(gè)線(xiàn)程等待同步然后再開(kāi)始下一步動(dòng)作的意思,但是CountDownLatch的下一步的動(dòng)作實(shí)施者是主線(xiàn)程,具有不可重復(fù)性;而CyclicBarrier的下一步動(dòng)作實(shí)施者還是“其他線(xiàn)程”本身,具有往復(fù)多次實(shí)施動(dòng)作的特點(diǎn)。

Semaphore

概述

Semaphore 一般譯作 信號(hào)量,它也是一種線(xiàn)程同步工具,主要用于多個(gè)線(xiàn)程對(duì)共享資源進(jìn)行并行操作的一種工具類(lèi)。它代表了一種許可的概念,是否允許多線(xiàn)程對(duì)同一資源進(jìn)行操作的許可,使用 Semaphore 可以控制并發(fā)訪(fǎng)問(wèn)資源的線(xiàn)程個(gè)數(shù)。

使用場(chǎng)景

Semaphore 的使用場(chǎng)景主要用于流量控制。

比如數(shù)據(jù)庫(kù)連接,同時(shí)使用的數(shù)據(jù)庫(kù)連接會(huì)有數(shù)量限制,數(shù)據(jù)庫(kù)連接不能超過(guò)一定的數(shù)量,當(dāng)連接到達(dá)了限制數(shù)量后,后面的線(xiàn)程只能排隊(duì)等前面的線(xiàn)程釋放數(shù)據(jù)庫(kù)連接后才能獲得數(shù)據(jù)庫(kù)連接。

比如停車(chē)場(chǎng)的場(chǎng)景中,一個(gè)停車(chē)場(chǎng)有有限數(shù)量的車(chē)位,同時(shí)能夠容納多少臺(tái)車(chē),車(chē)位滿(mǎn)了之后只有等里面的車(chē)離開(kāi)停車(chē)場(chǎng)外面的車(chē)才可以進(jìn)入。

案例

模擬一下停車(chē)場(chǎng)的業(yè)務(wù)場(chǎng)景:

在進(jìn)入停車(chē)場(chǎng)之前會(huì)有一個(gè)提示牌,上面顯示著停車(chē)位還有多少,當(dāng)車(chē)位為 0 時(shí),不能進(jìn)入停車(chē)場(chǎng),當(dāng)車(chē)位不為 0 時(shí),才會(huì)允許車(chē)輛進(jìn)入停車(chē)場(chǎng)。所以停車(chē)場(chǎng)有幾個(gè)關(guān)鍵因素:停車(chē)場(chǎng)車(chē)位的總?cè)萘浚?dāng)一輛車(chē)進(jìn)入時(shí),停車(chē)場(chǎng)車(chē)位的總?cè)萘?- 1,當(dāng)一輛車(chē)離開(kāi)時(shí),總?cè)萘?+ 1,停車(chē)場(chǎng)車(chē)位不足時(shí),車(chē)輛只能在停車(chē)場(chǎng)外等待。

public class SemaphoreDemo {

    private static Semaphore semaphore = new Semaphore(10);

    public static void main(String[] args) {
        for (int i = 0; i < 100; i++) {
            Thread thread = new Thread(new Runnable() {
                @Override
                public void run() {
                    System.out.println("歡迎 " + Thread.currentThread().getName() + " 來(lái)到停車(chē)場(chǎng)");
                    // 判斷是否允許停車(chē)
                    if (semaphore.availablePermits() == 0) {
                        System.out.println("車(chē)位不足,請(qǐng)耐心等待");
                    }
                    try {
                        // 嘗試獲取
                        semaphore.acquire();
                        System.out.println(Thread.currentThread().getName() + " 進(jìn)入停車(chē)場(chǎng)");
                        Thread.sleep(new Random().nextInt(10000));// 模擬車(chē)輛在停車(chē)場(chǎng)停留的時(shí)間
                        System.out.println(Thread.currentThread().getName() + " 駛出停車(chē)場(chǎng)");
                        semaphore.release();
                    } catch (InterruptedException e) {
                        e.printStackTrace();
                    }
                }
            }, i + "號(hào)車(chē)");
            thread.start();
        }
    }
}

Semaphore 的初始容量,也就是只有 10 個(gè)車(chē)位,我們用這 10 個(gè)車(chē)位來(lái)控制 100 輛車(chē)的流量,所以結(jié)果和我們預(yù)想的很相似,即大部分車(chē)都在等待狀態(tài)。但是同時(shí)仍允許一些車(chē)駛?cè)胪\?chē)場(chǎng),駛?cè)胪\?chē)場(chǎng)的車(chē)輛,就會(huì) semaphore.acquire 占用一個(gè)車(chē)位,駛出停車(chē)場(chǎng)時(shí),就會(huì) semaphore.release 讓出一個(gè)車(chē)位,讓后面的車(chē)再次駛?cè)搿?/p>

原理

Semaphore內(nèi)部有一個(gè)繼承了AQS的同步器Sync,重寫(xiě)了tryAcquireShared方法。在這個(gè)方法里,會(huì)去嘗試獲取資源。

如果獲取失?。ㄏ胍馁Y源數(shù)量小于目前已有的資源數(shù)量),就會(huì)返回一個(gè)負(fù)數(shù)(代表嘗試獲取資源失?。H缓螽?dāng)前線(xiàn)程就會(huì)進(jìn)入AQS的等待隊(duì)列。

Exchanger

概述

Exchanger類(lèi)用于兩個(gè)線(xiàn)程交換數(shù)據(jù)。它支持泛型,也就是說(shuō)你可以在兩個(gè)線(xiàn)程之間傳送任何數(shù)據(jù)。一個(gè)線(xiàn)程在完成一定的事務(wù)后想與另一個(gè)線(xiàn)程交換數(shù)據(jù),則第一個(gè)先拿出數(shù)據(jù)的線(xiàn)程會(huì)一直等待第二個(gè)線(xiàn)程,直到第二個(gè)線(xiàn)程拿著數(shù)據(jù)到來(lái)時(shí)才能彼此交換對(duì)應(yīng)數(shù)據(jù)。

案例

案例1:A同學(xué)和B同學(xué)交換各自收藏的大片。

public class ExchangerDemo {

    public static void main(String[] args) throws InterruptedException {
        Exchanger<String> stringExchanger = new Exchanger<>();

        Thread studentA = new Thread(() -> {
            try {
                String dataA = "A同學(xué)收藏多年的大片";
                String dataB = stringExchanger.exchange(dataA);
                System.out.println("A同學(xué)得到了" + dataB);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        });

        System.out.println("這個(gè)時(shí)候A同學(xué)是阻塞的,在等待B同學(xué)的大片");
        Thread.sleep(1000);

        Thread studentB = new Thread(() -> {
            try {
                String dataB = "B同學(xué)收藏多年的大片";
                String dataA = stringExchanger.exchange(dataB);
                System.out.println("B同學(xué)得到了" + dataA);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        });

        studentA.start();
        studentB.start();
    }
}

輸出:

這個(gè)時(shí)候A同學(xué)是阻塞的,在等待B同學(xué)的大片
A同學(xué)得到了B同學(xué)收藏多年的大片
B同學(xué)得到了A同學(xué)收藏多年的大片

可以看到,當(dāng)一個(gè)線(xiàn)程調(diào)用exchange方法后,它是處于阻塞狀態(tài)的,只有當(dāng)另一個(gè)線(xiàn)程也調(diào)用了exchange方法,它才會(huì)繼續(xù)向下執(zhí)行。

Exchanger類(lèi)還有一個(gè)有超時(shí)參數(shù)的方法,如果在指定時(shí)間內(nèi)沒(méi)有另一個(gè)線(xiàn)程調(diào)用exchange,就會(huì)拋出一個(gè)超時(shí)異常。

public V exchange(V x, long timeout, TimeUnit unit)

案例2:A同學(xué)被放鴿子,交易失敗。

public class ExchangerDemo {

    public static void main(String[] args) {
        Exchanger<String> stringExchanger = new Exchanger<>();
        Thread studentA = new Thread(() -> {
            String dataB = null;
            try {
                String dataA = "A同學(xué)收藏多年的大片";
                dataB = stringExchanger.exchange(dataA,5, TimeUnit.SECONDS);
                System.out.println("A同學(xué)得到了" + dataB);
            } catch (InterruptedException e) {
                e.printStackTrace();
            } catch (TimeoutException e) {
                System.out.println("等待超時(shí)-TimeoutException");
            }
            System.out.println("A同學(xué)得到了:"+dataB);
        });

        studentA.start();
    }
}

輸出:

等待超時(shí)-TimeoutException
A同學(xué)得到了:null

原理

Exchanger類(lèi)底層關(guān)鍵的技術(shù)有:

  • 使用CAS自旋指令完成數(shù)據(jù)交換;
  • 使用LockSupport的park方法使交換線(xiàn)程進(jìn)入休眠等待,使用LockSupport的unpark方法喚醒等待線(xiàn)程。
  • 此外還聲明了一個(gè)Node對(duì)象用于存儲(chǔ)交換數(shù)據(jù)。

Exchanger一般用于兩個(gè)線(xiàn)程之間更方便地在內(nèi)存中交換數(shù)據(jù),因?yàn)槠渲С址盒停晕覀兛梢詡鬏斎魏蔚臄?shù)據(jù),比如IO流或者IO緩存。根據(jù)JDK里面的注釋的說(shuō)法,可以總結(jié)為一下特性:

  • 此類(lèi)提供對(duì)外的操作是同步的;
  • 用于成對(duì)出現(xiàn)的線(xiàn)程之間交換數(shù)據(jù);
  • 可以視作雙向的同步隊(duì)列;
  • 可應(yīng)用于遺傳算法、流水線(xiàn)設(shè)計(jì)等場(chǎng)景。

需要注意的是,exchange是可以重復(fù)使用的。也就是說(shuō),兩個(gè)線(xiàn)程可以使用Exchanger在內(nèi)存中不斷地再交換數(shù)據(jù)。

小結(jié)

本文配合一些應(yīng)用場(chǎng)景介紹了JDK中提供的幾個(gè)并發(fā)工具類(lèi),簡(jiǎn)單分析了一下使用原理及業(yè)務(wù)場(chǎng)景,工作中,一旦有對(duì)應(yīng)的業(yè)務(wù)場(chǎng)景,可以試試這些工具類(lèi)。

以上就是淺析Java中并發(fā)工具類(lèi)的使用的詳細(xì)內(nèi)容,更多關(guān)于Java并發(fā)工具類(lèi)的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • JavaMail實(shí)現(xiàn)發(fā)送超文本(html)格式郵件的方法

    JavaMail實(shí)現(xiàn)發(fā)送超文本(html)格式郵件的方法

    這篇文章主要介紹了JavaMail實(shí)現(xiàn)發(fā)送超文本(html)格式郵件的方法,實(shí)例分析了java發(fā)送超文本文件的相關(guān)技巧,需要的朋友可以參考下
    2015-05-05
  • SpringBoot JWT令牌的使用

    SpringBoot JWT令牌的使用

    JWT令牌中包含了一個(gè)用戶(hù)名和哈希值,這些都需要進(jìn)行驗(yàn)證,本文主要介紹了SpringBoot JWT令牌的使用,具有一定的參考價(jià)值,感興趣的可以了解一下
    2024-03-03
  • Java編程guava RateLimiter實(shí)例解析

    Java編程guava RateLimiter實(shí)例解析

    這篇文章主要介紹了Java編程guava RateLimiter實(shí)例解析,具有一定借鑒價(jià)值,需要的朋友可以參考下
    2018-01-01
  • springboot 整合 freemarker代碼實(shí)例

    springboot 整合 freemarker代碼實(shí)例

    這篇文章主要介紹了springboot 整合 freemarker代碼實(shí)例,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2019-10-10
  • J2SE 1.5版本的新特性一覽

    J2SE 1.5版本的新特性一覽

    J2SE 1.5版本的新特性一覽...
    2006-12-12
  • 關(guān)于Hadoop中Spark?Streaming的基本概念

    關(guān)于Hadoop中Spark?Streaming的基本概念

    這篇文章主要介紹了關(guān)于Hadoop中Spark?Streaming的基本概念,Spark?Streaming是構(gòu)建在Spark上的實(shí)時(shí)計(jì)算框架,它擴(kuò)展了Spark處理大規(guī)模流式數(shù)據(jù)的能力,Spark?Streaming可結(jié)合批處理和交互式查詢(xún),需要的朋友可以參考下
    2023-07-07
  • java基于servlet實(shí)現(xiàn)文件上傳功能解析

    java基于servlet實(shí)現(xiàn)文件上傳功能解析

    這篇文章主要為大家詳細(xì)介紹了java基于servlet實(shí)現(xiàn)上傳功能,后臺(tái)使用java實(shí)現(xiàn),前端主要是js的ajax實(shí)現(xiàn),感興趣的小伙伴們可以參考一下
    2016-05-05
  • Java設(shè)計(jì)模式之構(gòu)建者模式知識(shí)總結(jié)

    Java設(shè)計(jì)模式之構(gòu)建者模式知識(shí)總結(jié)

    這幾天剛好在復(fù)習(xí)Java的設(shè)計(jì)模式,今天就給小伙伴們?nèi)婵偨Y(jié)一下開(kāi)發(fā)中最常用的設(shè)計(jì)模式-建造者模式的相關(guān)知識(shí),里面有很詳細(xì)的代碼示例及注釋哦,需要的朋友可以參考下
    2021-05-05
  • springboot解決Class path contains multiple SLF4J bindings問(wèn)題

    springboot解決Class path contains multiple 

    這篇文章主要介紹了springboot解決Class path contains multiple SLF4J bindings問(wèn)題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2023-07-07
  • Java計(jì)算兩個(gè)程序運(yùn)行時(shí)間的實(shí)例

    Java計(jì)算兩個(gè)程序運(yùn)行時(shí)間的實(shí)例

    下面小編就為大家?guī)?lái)一篇Java計(jì)算兩個(gè)程序運(yùn)行時(shí)間的實(shí)例。小編覺(jué)得挺不錯(cuò)的,現(xiàn)在就分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧
    2017-04-04

最新評(píng)論

隆昌县| 五原县| 余庆县| 崇左市| 枞阳县| 陆丰市| 新泰市| 宜州市| 涡阳县| 阿克陶县| 从江县| 自贡市| 鄂托克前旗| 万荣县| 新巴尔虎右旗| 苍溪县| 迭部县| 宜章县| 屏南县| 增城市| 东莞市| 浙江省| 抚州市| 修武县| 镇沅| 辛集市| 榆树市| 会宁县| 永靖县| 桂东县| 营口市| 乡宁县| 砀山县| 长泰县| 滨州市| 涡阳县| 达州市| 克什克腾旗| 宁武县| 湘潭市| 呼和浩特市|