CyclicBarrier之多線程中的循環(huán)柵欄詳解
1. CyclicBarrier簡介
現實生活中我們經常會遇到這樣的情景,在進行某個活動前需要等待人全部都齊了才開始。例如吃飯時要等全家人都上座了才動筷子,旅游時要等全部人都到齊了才出發(fā),比賽時要等運動員都上場后才開始。
在JUC包中為我們提供了一個同步工具類能夠很好的模擬這類場景,它就是CyclicBarrier類。利用CyclicBarrier類可以實現一組線程相互等待,當所有線程都到達某個屏障點(柵欄)后再進行后續(xù)的操作。下圖演示了這一過程:

CyclicBarrier可以使一定數量的線程反復地在柵欄(不同輪次或不同代)位置處匯集。
- CyclicBarrier字面意思是“可重復使用的柵欄”,CyclicBarrier 和 CountDownLatch 很像,只是 CyclicBarrier 可以有不止一個柵欄,因為它的柵欄(Barrier)可以重復使用(Cyclic)。
- 當線程到達柵欄位置時將調用await()方法,這個方法將阻塞(當前線程)直到所有線程都到達柵欄位置。如果所有線程都到達柵欄位置,那么柵欄將打開,此時所有的線程都將被釋放,而柵欄將被重置以便下次使用。
2.CyclicBarrier的使用
2.1 常用方法
//參數parties:表示要到達屏障 (柵欄)的線程數量
//參數Runnable: 最后一個線程到達屏障之后要做的任務
?
//構造方法1
public CyclicBarrier(int parties)
//構造方法2
public CyclicBarrier(int parties, Runnable barrierAction)
?
//線程調用await()方法表示當前線程已經到達柵欄,然后會被阻塞
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));
}2.2 使用舉例
適用場景:可用于需要多個線程均到達某一步之后才能繼續(xù)往下執(zhí)行的場景
//循環(huán)柵欄-可多次循環(huán)使用
CyclicBarrier cyclicBarrier = new CyclicBarrier(5,()->{
System.out.println(Thread.currentThread().getName()+" 完成最后任務!");
});
?
IntStream.range(1,6).forEach(i->{
new Thread(()->{
try {
Thread.sleep(new Double(Math.random()*3000).longValue());
System.out.println(Thread.currentThread().getName()+" 到達柵欄A");
cyclicBarrier.await();//屏障點A,當前線程會阻塞至此,等待計數器=0
System.out.println(Thread.currentThread().getName()+" 沖破柵欄A");
}catch (Exception e){
e.printStackTrace();
}
}).start();
});3.CyclicBarrier原理
CyclicBarrier是一道屏障,調用await()方法后,當前線程進入阻塞,當parties數量的線程調用await方法后,所有的await方法會返回并繼續(xù)往下執(zhí)行。
3.1 成員變量
/** CyclicBarrier使用的排他鎖*/
private final ReentrantLock lock = new ReentrantLock();
/** barrier被沖破前,線程等待的condition*/
private final Condition trip = lock.newCondition();
/** barrier被沖破時,需要滿足的參與線程數。*/
private final int parties;
/* barrier被沖破后執(zhí)行的方法。*/
private final Runnable barrierCommand;
/** 當其輪次 */
private Generation generation = new Generation();
?
/**
*目前等待剩余的參與者數量。從 parties倒數到0。每個輪次該值會被重置回parties
*/
private int count;(1)CyclicBarrier內部是通過條件隊列trip來對線程進行阻塞的,并且其內部維護了兩個int型的變量parties和count。
- parties表示每次攔截的線程數,該值在構造時進行賦值。
- count是內部計數器,它的初始值和parties相同,以后隨著每次await方法的調用而減1,直到減為0就將所有線程喚醒。
(2)CyclicBarrier有一個靜態(tài)內部類Generation,該類的對象代表柵欄的當前代,就像玩游戲時代表的本局游戲,利用它可以實現循環(huán)等待
(3)barrierCommand表示換代前執(zhí)行的任務,當count減為0時表示本局游戲結束,需要轉到下一局。在轉到下一局游戲之前會將所有阻塞的線程喚醒,在喚醒所有線程之前你可以通過指定barrierCommand來執(zhí)行自己的任務。
3.2 構造器
//構造器1:指定本局要攔截的線程數parties 及 本局結束時要執(zhí)行的任務
public CyclicBarrier(int parties, Runnable barrierAction) {
if (parties <= 0) throw new IllegalArgumentException();
this.parties = parties;
this.count = parties;
this.barrierCommand = barrierAction;
}
?
//構造器2
public CyclicBarrier(int parties) {
this(parties, null);
}3.3 等待的方法
CyclicBarrier類最主要的功能就是使先到達屏障點的線程阻塞并等待后面的線程,其中它提供了兩種等待的方法,分別是定時等待和非定時等待。源代碼中await()方法最終調用的是dowait()方法:
private int dowait(boolean timed, long nanos) throws InterruptedException, BrokenBarrierException,TimeoutException {
// 獲取獨占鎖
final ReentrantLock lock = this.lock;
lock.lock();//對共享資源count,generation操作前,需先上鎖保證線程安全
try {
// 當前代--當前輪次對象的引用
final Generation g = generation;
// 如果這輪次broken了,拋出異常
if (g.broken)
throw new BrokenBarrierException();
// 如果線程中斷了,拋出異常
if (Thread.interrupted()) {
breakBarrier();//如果被打斷,通過此方法設置當前輪次為broken狀態(tài),通知當前輪次所有等待的線程
throw new InterruptedException();//拋出異常
}
//自旋前
//1、count值-1
int index = --count;
// 2、判斷是否到0,若是,則沖破屏障點(說明最后一個線程已經到達)
if (index == 0) {
boolean ranAction = false;
try {
final Runnable command = barrierCommand;
// 3、執(zhí)行柵欄任務(若CyclicBarrier構造時傳入了Runnable,則調用)
if (command != null)
command.run();
ranAction = true;
// 4、更新一輪次,將count重置,將generation重置,喚醒之前等待的線程
nextGeneration();
return 0;
} finally {
// 如果執(zhí)行柵欄任務(command)的時候出現了異常,那么就認為本輪次破環(huán)了
if (!ranAction)
breakBarrier();
}
}
//計數器沒有到0 =》開始自旋,直到屏障被沖破,或者interrupted或者超時
for (;;) {
try {
// 開始等待;如果沒有時間限制,則直接等待,直到被喚醒(讓其他線程進入到lock的代碼塊執(zhí)行以上邏輯)
if (!timed)
trip.await();//阻塞,此時會釋放鎖,以讓其他線程進入await方法中,等待屏障被沖破后,向后執(zhí)行
// 如果有時間限制,則等待指定時間
else if (nanos > 0L)
nanos = trip.awaitNanos(nanos);
} catch (InterruptedException ie) {
// 如果當前線程阻塞被interrupt了,并且本輪次還沒有被break,那么修改本輪次狀態(tài)為broken
if (g == generation && ! g.broken) {
// 讓柵欄失效
breakBarrier();
throw ie;
} else {
// 上面條件不滿足,說明這個線程不是這輪次的,就不會影響當前這代柵欄的執(zhí)行,所以,就打個中斷標記
Thread.currentThread().interrupt();
}
}
// 當有任何一個線程中斷了,就會調用breakBarrier方法,就會喚醒其他的線程,其他線程醒來后,也要拋出異常
if (g.broken)
throw new BrokenBarrierException();
// g != generation表示正常換輪次了,返回當前線程所在柵欄的下標
// 如果 g == generation,說明還沒有換,那為什么會醒了?
// 因為一個線程可以使用多個柵欄,當別的柵欄喚醒了這個線程,就會走到這里,所以需要判斷是否是當前代。
// 正是因為這個原因,才需要generation來保證正確。
if (g != generation)
return index;
// 如果有時間限制,且時間小于等于0,銷毀柵欄并拋出異常
if (timed && nanos <= 0L) {
breakBarrier();
throw new TimeoutException();
}
}//自旋
} finally {
// 釋放獨占鎖
lock.unlock();
}
}更新本輪次的方法:nextGeneration()
private void nextGeneration() {
trip.signalAll();//喚醒本輪次等待的線程
count = parties;//重置count值為初始值,為下一輪次(代)使用
generation = new Generation();//更新本輪次對象,這樣自旋中的線程才會跳出自旋。
}
?
private static class Generation {
boolean broken = false;
}
?
private void breakBarrier() {
generation.broken = true;//設置標識
count = parties;//重置count值為初始值
trip.signalAll();//喚醒所有等待線程
}總結
以上為個人經驗,希望能給大家一個參考,也希望大家多多支持腳本之家。
相關文章
Spring?cloud?Hystrix注解初始化源碼過程解讀
這篇文章主要為大家介紹了Hystrix初始化部分,我們從源碼的角度分析一下@EnableCircuitBreaker以及@HystrixCommand注解的初始化過程,有需要的朋友可以借鑒參考下,希望能夠有所幫助2023-12-12
SpringBoot中使用JeecgBoot的Autopoi導出Excel的方法步驟
這篇文章主要介紹了SpringBoot中使用JeecgBoot的Autopoi導出Excel的方法步驟,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧2020-09-09
使用SpringBoot?+?Vue?+?Redis實現驗證碼登錄功能全過程
在現代web應用中,用戶驗證是非常重要的一部分,這篇文章主要介紹了使用SpringBoot?+?Vue?+?Redis實現驗證碼登錄功能的相關資料,文中通過代碼介紹的非常詳細,需要的朋友可以參考下2025-10-10

