淺析Java中并發(fā)工具類(lèi)的使用
在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,在dowait里index==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)格式郵件的方法,實(shí)例分析了java發(fā)送超文本文件的相關(guān)技巧,需要的朋友可以參考下2015-05-05
Java編程guava RateLimiter實(shí)例解析
這篇文章主要介紹了Java編程guava RateLimiter實(shí)例解析,具有一定借鑒價(jià)值,需要的朋友可以參考下2018-01-01
springboot 整合 freemarker代碼實(shí)例
這篇文章主要介紹了springboot 整合 freemarker代碼實(shí)例,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2019-10-10
關(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)文件上傳功能解析
這篇文章主要為大家詳細(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é)
這幾天剛好在復(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 
這篇文章主要介紹了springboot解決Class path contains multiple SLF4J bindings問(wèn)題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2023-07-07
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

