--- title: AQS 详解 category: Java tag: - Javaå¹¶å --- ## AQS ä»ç» AQS çå ¨ç§°ä¸º `AbstractQueuedSynchronizer` ï¼ç¿»è¯è¿æ¥çææå°±æ¯æ½è±¡éå忥å¨ãè¿ä¸ªç±»å¨ `java.util.concurrent.locks` å ä¸é¢ã  AQS å°±æ¯ä¸ä¸ªæ½è±¡ç±»ï¼ä¸»è¦ç¨æ¥æå»ºéå忥å¨ã ```java public abstract class AbstractQueuedSynchronizer extends AbstractOwnableSynchronizer implements java.io.Serializable { } ``` AQS 为æå»ºéå忥卿ä¾äºä¸äºéç¨åè½çæ¯å®ç°ï¼å æ¤ï¼ä½¿ç¨ AQS è½ç®åä¸é«æå°æé åºåºç¨å¹¿æ³ç大éç忥å¨ï¼æ¯å¦æä»¬æå°ç `ReentrantLock`ï¼`Semaphore`ï¼å ¶ä»çè¯¸å¦ `ReentrantReadWriteLock`ï¼`SynchronousQueue`çççæ¯åºäº AQS çã ## AQS åç å¨é¢è¯ä¸è¢«é®å°å¹¶åç¥è¯çæ¶åï¼å¤§å¤é½ä¼è¢«é®å°âè¯·ä½ è¯´ä¸ä¸èªå·±å¯¹äº AQS åçççè§£âãä¸é¢ç»å¤§å®¶ä¸ä¸ªç¤ºä¾ä¾å¤§å®¶åèï¼é¢è¯ä¸æ¯èé¢ï¼å¤§å®¶ä¸å®è¦å å ¥èªå·±çææ³ï¼å³ä½¿å å ¥ä¸äºèªå·±çææ³ä¹è¦ä¿è¯èªå·±è½å¤éä¿çè®²åºæ¥è䏿¯èåºæ¥ã ### AQS æ ¸å¿ææ³ AQS æ ¸å¿ææ³æ¯ï¼å¦æè¢«è¯·æ±çå ±äº«èµæºç©ºé²ï¼åå°å½å请æ±èµæºç线ç¨è®¾ç½®ä¸ºææçå·¥ä½çº¿ç¨ï¼å¹¶ä¸å°å ±äº«èµæºè®¾ç½®ä¸ºéå®ç¶æãå¦æè¢«è¯·æ±çå ±äº«èµæºè¢«å ç¨ï¼é£ä¹å°±éè¦ä¸å¥çº¿ç¨é»å¡çå¾ ä»¥å被å¤éæ¶éåé çæºå¶ï¼è¿ä¸ªæºå¶ AQS æ¯åºäº **CLH é** ï¼Craig, Landin, and Hagersten locksï¼ å®ç°çã CLH éæ¯å¯¹èªæéçä¸ç§æ¹è¿ï¼æ¯ä¸ä¸ªèæçååéåï¼èæçååéåå³ä¸åå¨éåå®ä¾ï¼ä» åå¨ç»ç¹ä¹é´çå ³èå ³ç³»ï¼ï¼ææ¶è·åä¸å°éç线ç¨å°è¢«å å ¥å°è¯¥éåä¸ãAQS å°æ¯æ¡è¯·æ±å ±äº«èµæºç线ç¨å°è£ æä¸ä¸ª CLH éåéçä¸ä¸ªç»ç¹ï¼Nodeï¼æ¥å®ç°éçåé ãå¨ CLH éåéä¸ï¼ä¸ä¸ªèç¹è¡¨ç¤ºä¸ä¸ªçº¿ç¨ï¼å®ä¿åç线ç¨çå¼ç¨ï¼threadï¼ã å½åèç¹å¨éåä¸çç¶æï¼waitStatusï¼ãå驱èç¹ï¼prevï¼ãåç»§èç¹ï¼nextï¼ã CLH éåéç»æå¦ä¸å¾æç¤ºï¼  å ³äº AQS æ ¸å¿æ°æ®ç»æ-CLH éç详ç»è§£è¯»ï¼å¼ºçæ¨èé 读 [Java AQS æ ¸å¿æ°æ®ç»æ-CLH é - Qunar ææ¯æ²é¾](https://mp.weixin.qq.com/s/jEx-4XhNGOFdCo4Nou5tqg) è¿ç¯æç« ã AQS(`AbstractQueuedSynchronizer`)çæ ¸å¿åçå¾ï¼å¾æº[Java å¹¶åä¹ AQS 详解](https://www.cnblogs.com/waterystone/p/4920797.html)ï¼å¦ä¸ï¼  AQS ä½¿ç¨ **int æååé `state` è¡¨ç¤ºåæ¥ç¶æ**ï¼éè¿å ç½®ç **线ç¨çå¾ éå** æ¥å®æè·åèµæºçº¿ç¨çæéå·¥ä½ã `state` åéç± `volatile` 修饰ï¼ç¨äºå±ç¤ºå½å临çèµæºçè·éæ åµã ```java // å ±äº«åéï¼ä½¿ç¨volatile修饰ä¿è¯çº¿ç¨å¯è§æ§ private volatile int state; ``` å¦å¤ï¼ç¶æä¿¡æ¯ `state` å¯ä»¥éè¿ `protected` ç±»åç`getState()`ã`setState()`å`compareAndSetState()` è¿è¡æä½ãå¹¶ä¸ï¼è¿å ä¸ªæ¹æ³é½æ¯ `final` 修饰çï¼å¨åç±»ä¸æ æ³è¢«éåã ```java //è¿ååæ¥ç¶æçå½åå¼ protected final int getState() { return state; } // è®¾ç½®åæ¥ç¶æçå¼ protected final void setState(int newState) { state = newState; } //ååå°ï¼CASæä½ï¼å°åæ¥ç¶æå¼è®¾ç½®ä¸ºç»å®å¼update妿å½ååæ¥ç¶æçå¼çäºexpectï¼ææå¼ï¼ protected final boolean compareAndSetState(int expect, int update) { return unsafe.compareAndSwapInt(this, stateOffset, expect, update); } ``` 以 `ReentrantLock` 为ä¾ï¼`state` åå§å¼ä¸º 0ï¼è¡¨ç¤ºæªéå®ç¶æãA çº¿ç¨ `lock()` æ¶ï¼ä¼è°ç¨ `tryAcquire()` ç¬å 该éå¹¶å° `state+1` ãæ¤åï¼å ¶ä»çº¿ç¨å `tryAcquire()` æ¶å°±ä¼å¤±è´¥ï¼ç´å° A çº¿ç¨ `unlock()` å° `state=`0ï¼å³éæ¾éï¼ä¸ºæ¢ï¼å ¶å®çº¿ç¨æææºä¼è·å该éãå½ç¶ï¼éæ¾éä¹åï¼A 线ç¨èªå·±æ¯å¯ä»¥éå¤è·åæ¤éçï¼`state` ä¼ç´¯å ï¼ï¼è¿å°±æ¯å¯éå ¥çæ¦å¿µãä½è¦æ³¨æï¼è·åå¤å°æ¬¡å°±è¦éæ¾å¤å°æ¬¡ï¼è¿æ ·æè½ä¿è¯ state æ¯è½åå°é¶æçãç¸å ³é 读ï¼[ä» ReentrantLock çå®ç°ç AQS çåçååºç¨ - ç¾å¢ææ¯å¢é](./reentrantlock.md)ã å以 `CountDownLatch` 以ä¾ï¼ä»»å¡å为 N 个å线ç¨å»æ§è¡ï¼`state` ä¹åå§å为 Nï¼æ³¨æ N è¦ä¸çº¿ç¨ä¸ªæ°ä¸è´ï¼ãè¿ N 个åçº¿ç¨æ¯å¹¶è¡æ§è¡çï¼æ¯ä¸ªåçº¿ç¨æ§è¡å®å`countDown()` 䏿¬¡ï¼state ä¼ CAS(Compare and Swap) å 1ãçå°ææå线ç¨é½æ§è¡å®å(å³ `state=0` )ï¼ä¼ `unpark()` 主è°ç¨çº¿ç¨ï¼ç¶å主è°ç¨çº¿ç¨å°±ä¼ä» `await()` 彿°è¿åï¼ç»§ç»åä½å¨ä½ã ### AQS èµæºå ±äº«æ¹å¼ AQS å®ä¹ä¸¤ç§èµæºå ±äº«æ¹å¼ï¼`Exclusive`ï¼ç¬å ï¼åªæä¸ä¸ªçº¿ç¨è½æ§è¡ï¼å¦`ReentrantLock`ï¼å`Share`ï¼å ±äº«ï¼å¤ä¸ªçº¿ç¨å¯åæ¶æ§è¡ï¼å¦`Semaphore`/`CountDownLatch`ï¼ã ä¸è¬æ¥è¯´ï¼èªå®ä¹åæ¥å¨çå ±äº«æ¹å¼è¦ä¹æ¯ç¬å ï¼è¦ä¹æ¯å ±äº«ï¼ä»ä»¬ä¹åªéå®ç°`tryAcquire-tryRelease`ã`tryAcquireShared-tryReleaseShared`ä¸çä¸ç§å³å¯ãä½ AQS 乿¯æèªå®ä¹åæ¥å¨åæ¶å®ç°ç¬å åå ±äº«ä¸¤ç§æ¹å¼ï¼å¦`ReentrantReadWriteLock`ã ### èªå®ä¹åæ¥å¨ 忥å¨ç设计æ¯åºäºæ¨¡æ¿æ¹æ³æ¨¡å¼çï¼å¦æéè¦èªå®ä¹åæ¥å¨ä¸è¬çæ¹å¼æ¯è¿æ ·ï¼æ¨¡æ¿æ¹æ³æ¨¡å¼å¾ç»å ¸çä¸ä¸ªåºç¨ï¼ï¼ 1. 使ç¨è ç»§æ¿ `AbstractQueuedSynchronizer` å¹¶éåæå®çæ¹æ³ã 2. å° AQS ç»åå¨èªå®ä¹åæ¥ç»ä»¶çå®ç°ä¸ï¼å¹¶è°ç¨å ¶æ¨¡æ¿æ¹æ³ï¼èè¿äºæ¨¡æ¿æ¹æ³ä¼è°ç¨ä½¿ç¨è éåçæ¹æ³ã è¿åæä»¬ä»¥å¾éè¿å®ç°æ¥å£çæ¹å¼æå¾å¤§åºå«ï¼è¿æ¯æ¨¡æ¿æ¹æ³æ¨¡å¼å¾ç»å ¸çä¸ä¸ªè¿ç¨ã **AQS 使ç¨äºæ¨¡æ¿æ¹æ³æ¨¡å¼ï¼èªå®ä¹åæ¥å¨æ¶éè¦éåä¸é¢å 个 AQS æä¾çé©åæ¹æ³ï¼** ```java //ç¬å æ¹å¼ãå°è¯è·åèµæºï¼æååè¿åtrueï¼å¤±è´¥åè¿åfalseã protected boolean tryAcquire(int) //ç¬å æ¹å¼ãå°è¯éæ¾èµæºï¼æååè¿åtrueï¼å¤±è´¥åè¿åfalseã protected boolean tryRelease(int) //å ±äº«æ¹å¼ãå°è¯è·åèµæºãè´æ°è¡¨ç¤ºå¤±è´¥ï¼0表示æåï¼ä½æ²¡æå©ä½å¯ç¨èµæºï¼æ£æ°è¡¨ç¤ºæåï¼ä¸æå©ä½èµæºã protected int tryAcquireShared(int) //å ±äº«æ¹å¼ãå°è¯éæ¾èµæºï¼æååè¿åtrueï¼å¤±è´¥åè¿åfalseã protected boolean tryReleaseShared(int) //è¯¥çº¿ç¨æ¯å¦æ£å¨ç¬å èµæºãåªæç¨å°conditionæéè¦å»å®ç°å®ã protected boolean isHeldExclusively() ``` **ä»ä¹æ¯é©åæ¹æ³å¢ï¼** é©åæ¹æ³æ¯ä¸ç§è¢«å£°æå¨æ½è±¡ç±»ä¸çæ¹æ³ï¼ä¸è¬ä½¿ç¨ `protected` å ³é®å修饰ï¼å®å¯ä»¥æ¯ç©ºæ¹æ³ï¼ç±åç±»å®ç°ï¼ï¼ä¹å¯ä»¥æ¯é»è®¤å®ç°çæ¹æ³ã模æ¿è®¾è®¡æ¨¡å¼éè¿é©åæ¹æ³æ§å¶åºå®æ¥éª¤çå®ç°ã ç¯å¹ é®é¢ï¼è¿éå°±ä¸è¯¦ç»ä»ç»æ¨¡æ¿æ¹æ³æ¨¡å¼äºï¼ä¸å¤ªäºè§£çå°ä¼ä¼´å¯ä»¥ççè¿ç¯æç« ï¼[ç¨ Java8 æ¹é åçæ¨¡æ¿æ¹æ³æ¨¡å¼ççæ¯ yyds!](https://mp.weixin.qq.com/s/zpScSCktFpnSWHWIQem2jg)ã é¤äºä¸é¢æå°çé©åæ¹æ³ä¹å¤ï¼AQS ç±»ä¸çå ¶ä»æ¹æ³é½æ¯ `final` ï¼æä»¥æ æ³è¢«å ¶ä»ç±»éåã ## 常è§åæ¥å·¥å ·ç±» ä¸é¢ä»ç»å 个åºäº AQS ç常è§åæ¥å·¥å ·ç±»ã ### Semaphore(ä¿¡å·é) #### ä»ç» `synchronized` å `ReentrantLock` 齿¯ä¸æ¬¡åªå 许ä¸ä¸ªçº¿ç¨è®¿é®æä¸ªèµæºï¼è`Semaphore`(ä¿¡å·é)å¯ä»¥ç¨æ¥æ§å¶åæ¶è®¿é®ç¹å®èµæºççº¿ç¨æ°éã Semaphore ç使ç¨ç®åï¼æä»¬è¿éå设æ N(N>5) ä¸ªçº¿ç¨æ¥è·å `Semaphore` ä¸çå ±äº«èµæºï¼ä¸é¢ç代ç 表示å䏿¶å» N 个线ç¨ä¸åªæ 5 个线ç¨è½è·åå°å ±äº«èµæºï¼å ¶ä»çº¿ç¨é½ä¼é»å¡ï¼åªæè·åå°å ±äº«èµæºççº¿ç¨æè½æ§è¡ãçå°æçº¿ç¨éæ¾äºå ±äº«èµæºï¼å ¶ä»é»å¡ççº¿ç¨æè½è·åå°ã ```java // åå§å ±äº«èµæºæ°é final Semaphore semaphore = new Semaphore(5); // è·å1ä¸ªè®¸å¯ semaphore.acquire(); // éæ¾1ä¸ªè®¸å¯ semaphore.release(); ``` å½åå§çèµæºä¸ªæ°ä¸º 1 çæ¶åï¼`Semaphore` éå为æä»éã `Semaphore` æä¸¤ç§æ¨¡å¼ï¼ã - **å ¬å¹³æ¨¡å¼ï¼** è°ç¨ `acquire()` æ¹æ³ç顺åºå°±æ¯è·å许å¯è¯ç顺åºï¼éµå¾ª FIFOï¼ - **éå ¬å¹³æ¨¡å¼ï¼** æ¢å å¼çã `Semaphore` 对åºç两个æé æ¹æ³å¦ä¸ï¼ ```java public Semaphore(int permits) { sync = new NonfairSync(permits); } public Semaphore(int permits, boolean fair) { sync = fair ? new FairSync(permits) : new NonfairSync(permits); } ``` **è¿ä¸¤ä¸ªæé æ¹æ³ï¼é½å¿ é¡»æä¾è®¸å¯çæ°éï¼ç¬¬äºä¸ªæé æ¹æ³å¯ä»¥æå®æ¯å ¬å¹³æ¨¡å¼è¿æ¯éå ¬å¹³æ¨¡å¼ï¼é»è®¤éå ¬å¹³æ¨¡å¼ã** `Semaphore` é常ç¨äºé£äºèµæºææç¡®è®¿é®æ°ééå¶çåºæ¯æ¯å¦éæµï¼ä» éäºåæºæ¨¡å¼ï¼å®é 项ç®ä¸æ¨èä½¿ç¨ Redis +Lua æ¥åéæµï¼ã #### åç `Semaphore` æ¯å ±äº«éçä¸ç§å®ç°ï¼å®é»è®¤æé AQS ç `state` å¼ä¸º `permits`ï¼ä½ å¯ä»¥å° `permits` çå¼ç解为许å¯è¯çæ°éï¼åªææ¿å°è®¸å¯è¯ççº¿ç¨æè½æ§è¡ã è°ç¨`semaphore.acquire()` ï¼çº¿ç¨å°è¯è·å许å¯è¯ï¼å¦æ `state >= 0` çè¯ï¼å表示å¯ä»¥è·åæåã妿è·åæåçè¯ï¼ä½¿ç¨ CAS æä½å»ä¿®æ¹ `state` çå¼ `state=state-1`ã妿 `state<0` çè¯ï¼å表示许å¯è¯æ°éä¸è¶³ãæ¤æ¶ä¼å建ä¸ä¸ª Node èç¹å å ¥é»å¡éåï¼æèµ·å½å线ç¨ã ```java /** * è·å1个许å¯è¯ */ public void acquire() throws InterruptedException { sync.acquireSharedInterruptibly(1); } /** * å ±äº«æ¨¡å¼ä¸è·å许å¯è¯ï¼è·åæååè¿åï¼å¤±è´¥åå å ¥é»å¡éåï¼æèµ·çº¿ç¨ */ public final void acquireSharedInterruptibly(int arg) throws InterruptedException { if (Thread.interrupted()) throw new InterruptedException(); // å°è¯è·å许å¯è¯ï¼arg为è·å许å¯è¯ä¸ªæ°ï¼å½å¯ç¨è®¸å¯è¯æ°åå½åè·åç许å¯è¯æ°ç»æå°äº0,åå建ä¸ä¸ªèç¹å å ¥é»å¡éåï¼æèµ·å½å线ç¨ã if (tryAcquireShared(arg) < 0) doAcquireSharedInterruptibly(arg); } ``` è°ç¨`semaphore.release();` ï¼çº¿ç¨å°è¯éæ¾è®¸å¯è¯ï¼å¹¶ä½¿ç¨ CAS æä½å»ä¿®æ¹ `state` çå¼ `state=state+1`ãéæ¾è®¸å¯è¯æåä¹åï¼åæ¶ä¼å¤é忥éåä¸çä¸ä¸ªçº¿ç¨ã被å¤éç线ç¨ä¼éæ°å°è¯å»ä¿®æ¹ `state` çå¼ `state=state-1` ï¼å¦æ `state>=0` åè·å令çæåï¼å¦åéæ°è¿å ¥é»å¡éåï¼æèµ·çº¿ç¨ã ```java // éæ¾ä¸ä¸ªè®¸å¯è¯ public void release() { sync.releaseShared(1); } // éæ¾å ±äº«éï¼åæ¶ä¼å¤é忥éåä¸çä¸ä¸ªçº¿ç¨ã public final boolean releaseShared(int arg) { //éæ¾å ±äº«é if (tryReleaseShared(arg)) { //å¤é忥éåä¸çä¸ä¸ªçº¿ç¨ doReleaseShared(); return true; } return false; } ``` #### 宿 ```java /** * * @author Snailclimb * @date 2018å¹´9æ30æ¥ * @Description: éè¦ä¸æ¬¡æ§æ¿ä¸ä¸ªè®¸å¯çæ åµ */ public class SemaphoreExample1 { // 请æ±çæ°é private static final int threadCount = 550; public static void main(String[] args) throws InterruptedException { // å建ä¸ä¸ªå ·æåºå®çº¿ç¨æ°éççº¿ç¨æ± 对象ï¼å¦æè¿éçº¿ç¨æ± ççº¿ç¨æ°éç»å¤ªå°çè¯ä½ ä¼åç°æ§è¡ç徿 ¢ï¼ ExecutorService threadPool = Executors.newFixedThreadPool(300); // åå§è®¸å¯è¯æ°é final Semaphore semaphore = new Semaphore(20); for (int i = 0; i < threadCount; i++) { final int threadnum = i; threadPool.execute(() -> {// Lambda 表达å¼çè¿ç¨ try { semaphore.acquire();// è·åä¸ä¸ªè®¸å¯ï¼æä»¥å¯è¿è¡çº¿ç¨æ°é为20/1=20 test(threadnum); semaphore.release();// éæ¾ä¸ä¸ªè®¸å¯ } catch (InterruptedException e) { // TODO Auto-generated catch block e.printStackTrace(); } }); } threadPool.shutdown(); System.out.println("finish"); } public static void test(int threadnum) throws InterruptedException { Thread.sleep(1000);// 模æè¯·æ±çèæ¶æä½ System.out.println("threadnum:" + threadnum); Thread.sleep(1000);// 模æè¯·æ±çèæ¶æä½ } } ``` æ§è¡ `acquire()` æ¹æ³é»å¡ï¼ç´å°æä¸ä¸ªè®¸å¯è¯å¯ä»¥è·å¾ç¶åæ¿èµ°ä¸ä¸ªè®¸å¯è¯ï¼æ¯ä¸ª `release` æ¹æ³å¢å ä¸ä¸ªè®¸å¯è¯ï¼è¿å¯è½ä¼éæ¾ä¸ä¸ªé»å¡ç `acquire()` æ¹æ³ãç¶èï¼å ¶å®å¹¶æ²¡æå®é ç许å¯è¯è¿ä¸ªå¯¹è±¡ï¼`Semaphore` åªæ¯ç»´æäºä¸ä¸ªå¯è·å¾è®¸å¯è¯çæ°éã `Semaphore` ç»å¸¸ç¨äºéå¶è·åæç§èµæºççº¿ç¨æ°éã å½ç¶ä¸æ¬¡ä¹å¯ä»¥ä¸æ¬¡æ¿ååéæ¾å¤ä¸ªè®¸å¯ï¼ä¸è¿ä¸è¬æ²¡æå¿ è¦è¿æ ·åï¼ ```java semaphore.acquire(5);// è·å5个许å¯ï¼æä»¥å¯è¿è¡çº¿ç¨æ°é为20/5=4 test(threadnum); semaphore.release(5);// éæ¾5ä¸ªè®¸å¯ ``` é¤äº `acquire()` æ¹æ³ä¹å¤ï¼å¦ä¸ä¸ªæ¯è¾å¸¸ç¨çä¸ä¹å¯¹åºçæ¹æ³æ¯ `tryAcquire()` æ¹æ³ï¼è¯¥æ¹æ³å¦æè·åä¸å°è®¸å¯å°±ç«å³è¿å falseã [issue645 è¡¥å å 容](https://github.com/Snailclimb/JavaGuide/issues/645)ï¼ > `Semaphore` ä¸ `CountDownLatch` 䏿 ·ï¼ä¹æ¯å ±äº«éçä¸ç§å®ç°ãå®é»è®¤æé AQS ç `state` 为 `permits`ã彿§è¡ä»»å¡ççº¿ç¨æ°éè¶ åº `permits`ï¼é£ä¹å¤ä½ç线ç¨å°ä¼è¢«æ¾å ¥é»å¡éå `Park`,å¹¶èªæå¤æ `state` æ¯å¦å¤§äº 0ãåªæå½ `state` å¤§äº 0 çæ¶åï¼é»å¡ççº¿ç¨æè½ç»§ç»æ§è¡,æ¤æ¶å åæ§è¡ä»»å¡ç线ç¨ç»§ç»æ§è¡ `release()` æ¹æ³ï¼`release()` æ¹æ³ä½¿å¾ state çåéä¼å 1ï¼é£ä¹èªæç线ç¨ä¾¿ä¼å¤ææåã > 妿¤ï¼æ¯æ¬¡åªææå¤ä¸è¶ è¿ `permits` æ°éç线ç¨è½èªææåï¼ä¾¿éå¶äºæ§è¡ä»»å¡çº¿ç¨çæ°éã ### CountDownLatch ï¼å计æ¶å¨ï¼ #### ä»ç» `CountDownLatch` å 许 `count` 个线ç¨é»å¡å¨ä¸ä¸ªå°æ¹ï¼ç´è³ææçº¿ç¨çä»»å¡é½æ§è¡å®æ¯ã `CountDownLatch` æ¯ä¸æ¬¡æ§çï¼è®¡æ°å¨çå¼åªè½å¨æé æ¹æ³ä¸åå§å䏿¬¡ï¼ä¹å没æä»»ä½æºå¶åæ¬¡å¯¹å ¶è®¾ç½®å¼ï¼å½ `CountDownLatch` 使ç¨å®æ¯åï¼å®ä¸è½å次被使ç¨ã #### åç `CountDownLatch` æ¯å ±äº«éçä¸ç§å®ç°,å®é»è®¤æé AQS ç `state` å¼ä¸º `count`ãå½çº¿ç¨ä½¿ç¨ `countDown()` æ¹æ³æ¶,å ¶å®ä½¿ç¨äº`tryReleaseShared`æ¹æ³ä»¥ CAS çæä½æ¥åå° `state`,ç´è³ `state` 为 0 ãå½è°ç¨ `await()` æ¹æ³çæ¶åï¼å¦æ `state` ä¸ä¸º 0ï¼é£å°±è¯æä»»å¡è¿æ²¡ææ§è¡å®æ¯ï¼`await()` æ¹æ³å°±ä¼ä¸ç´é»å¡ï¼ä¹å°±æ¯è¯´ `await()` æ¹æ³ä¹åçè¯å¥ä¸ä¼è¢«æ§è¡ãç¶åï¼`CountDownLatch` ä¼èªæ CAS 夿 `state == 0`ï¼å¦æ `state == 0` çè¯ï¼å°±ä¼éæ¾ææçå¾ ç线ç¨ï¼`await()` æ¹æ³ä¹åçè¯å¥å¾å°æ§è¡ã #### 宿 **CountDownLatch ç两ç§å ¸åç¨æ³**ï¼ 1. æä¸çº¿ç¨å¨å¼å§è¿è¡åçå¾ n ä¸ªçº¿ç¨æ§è¡å®æ¯ : å° `CountDownLatch` ç计æ°å¨åå§å为 n ï¼`new CountDownLatch(n)`ï¼ï¼æ¯å½ä¸ä¸ªä»»å¡çº¿ç¨æ§è¡å®æ¯ï¼å°±å°è®¡æ°å¨å 1 ï¼`countdownlatch.countDown()`ï¼ï¼å½è®¡æ°å¨çå¼å为 0 æ¶ï¼å¨ `CountDownLatch ä¸ await()` ç线ç¨å°±ä¼è¢«å¤éãä¸ä¸ªå ¸ååºç¨åºæ¯å°±æ¯å¯å¨ä¸ä¸ªæå¡æ¶ï¼ä¸»çº¿ç¨éè¦çå¾ å¤ä¸ªç»ä»¶å è½½å®æ¯ï¼ä¹ååç»§ç»æ§è¡ã 2. å®ç°å¤ä¸ªçº¿ç¨å¼å§æ§è¡ä»»å¡çæå¤§å¹¶è¡æ§ï¼æ³¨ææ¯å¹¶è¡æ§ï¼ä¸æ¯å¹¶åï¼å¼ºè°çæ¯å¤ä¸ªçº¿ç¨å¨æä¸æ¶å»åæ¶å¼å§æ§è¡ã类似äºèµè·ï¼å°å¤ä¸ªçº¿ç¨æ¾å°èµ·ç¹ï¼çå¾ å令æªåï¼ç¶ååæ¶å¼è·ãåæ³æ¯åå§åä¸ä¸ªå ±äº«ç `CountDownLatch` 对象ï¼å°å ¶è®¡æ°å¨åå§å为 1 ï¼`new CountDownLatch(1)`ï¼ï¼å¤ä¸ªçº¿ç¨å¨å¼å§æ§è¡ä»»å¡åé¦å `coundownlatch.await()`ï¼å½ä¸»çº¿ç¨è°ç¨ `countDown()` æ¶ï¼è®¡æ°å¨å为 0ï¼å¤ä¸ªçº¿ç¨åæ¶è¢«å¤éã **CountDownLatch 代ç 示ä¾**ï¼ ```java /** * * @author SnailClimb * @date 2018å¹´10æ1æ¥ * @Description: CountDownLatch ä½¿ç¨æ¹æ³ç¤ºä¾ */ public class CountDownLatchExample1 { // 请æ±çæ°é private static final int threadCount = 550; public static void main(String[] args) throws InterruptedException { // å建ä¸ä¸ªå ·æåºå®çº¿ç¨æ°éççº¿ç¨æ± 对象ï¼å¦æè¿éçº¿ç¨æ± ççº¿ç¨æ°éç»å¤ªå°çè¯ä½ ä¼åç°æ§è¡ç徿 ¢ï¼ ExecutorService threadPool = Executors.newFixedThreadPool(300); final CountDownLatch countDownLatch = new CountDownLatch(threadCount); for (int i = 0; i < threadCount; i++) { final int threadnum = i; threadPool.execute(() -> {// Lambda 表达å¼çè¿ç¨ try { test(threadnum); } catch (InterruptedException e) { // TODO Auto-generated catch block e.printStackTrace(); } finally { countDownLatch.countDown();// 表示ä¸ä¸ªè¯·æ±å·²ç»è¢«å®æ } }); } countDownLatch.await(); threadPool.shutdown(); System.out.println("finish"); } public static void test(int threadnum) throws InterruptedException { Thread.sleep(1000);// 模æè¯·æ±çèæ¶æä½ System.out.println("threadnum:" + threadnum); Thread.sleep(1000);// 模æè¯·æ±çèæ¶æä½ } } ``` ä¸é¢ç代ç ä¸ï¼æä»¬å®ä¹äºè¯·æ±çæ°é为 550ï¼å½è¿ 550 个请æ±è¢«å¤ç宿ä¹åï¼æä¼æ§è¡`System.out.println("finish");`ã ä¸ `CountDownLatch` çç¬¬ä¸æ¬¡äº¤äºæ¯ä¸»çº¿ç¨çå¾ å ¶ä»çº¿ç¨ã主线ç¨å¿ é¡»å¨å¯å¨å ¶ä»çº¿ç¨åç«å³è°ç¨ `CountDownLatch.await()` æ¹æ³ãè¿æ ·ä¸»çº¿ç¨çæä½å°±ä¼å¨è¿ä¸ªæ¹æ³ä¸é»å¡ï¼ç´å°å ¶ä»çº¿ç¨å®æåèªçä»»å¡ã å ¶ä» N 个线ç¨å¿ é¡»å¼ç¨éé对象ï¼å 为ä»ä»¬éè¦éç¥ `CountDownLatch` 对象ï¼ä»ä»¬å·²ç»å®æäºåèªçä»»å¡ãè¿ç§éç¥æºå¶æ¯éè¿ `CountDownLatch.countDown()`æ¹æ³æ¥å®æçï¼æ¯è°ç¨ä¸æ¬¡è¿ä¸ªæ¹æ³ï¼å¨æé 彿°ä¸åå§åç count å¼å°±å 1ãæä»¥å½ N 个线ç¨é½è° ç¨äºè¿ä¸ªæ¹æ³ï¼count çå¼çäº 0ï¼ç¶å主线ç¨å°±è½éè¿ `await()`æ¹æ³ï¼æ¢å¤æ§è¡èªå·±çä»»å¡ã åæä¸å´ï¼`CountDownLatch` ç `await()` æ¹æ³ä½¿ç¨ä¸å½å¾å®¹æäº§çæ»éï¼æ¯å¦æä»¬ä¸é¢ä»£ç ä¸ç for å¾ªç¯æ¹ä¸ºï¼ ```java for (int i = 0; i < threadCount-1; i++) { ....... } ``` è¿æ ·å°±å¯¼è´ `count` ç弿²¡åæ³çäº 0ï¼ç¶åå°±ä¼å¯¼è´ä¸ç´çå¾ ã ### CyclicBarrier(å¾ªç¯æ æ ) #### ä»ç» `CyclicBarrier` å `CountDownLatch` é常类似ï¼å®ä¹å¯ä»¥å®ç°çº¿ç¨é´çææ¯çå¾ ï¼ä½æ¯å®çåè½æ¯ `CountDownLatch` æ´å 夿å强大ã主è¦åºç¨åºæ¯å `CountDownLatch` 类似ã > `CountDownLatch` çå®ç°æ¯åºäº AQS çï¼è `CycliBarrier` æ¯åºäº `ReentrantLock`(`ReentrantLock` ä¹å±äº AQS 忥å¨)å `Condition` çã `CyclicBarrier` çåé¢æææ¯å¯å¾ªç¯ä½¿ç¨ï¼Cyclicï¼çå±éï¼Barrierï¼ãå®è¦åçäºæ æ¯ï¼è®©ä¸ç»çº¿ç¨å°è¾¾ä¸ä¸ªå±éï¼ä¹å¯ä»¥å«åæ¥ç¹ï¼æ¶è¢«é»å¡ï¼ç´å°æåä¸ä¸ªçº¿ç¨å°è¾¾å±éæ¶ï¼å±éæä¼å¼é¨ï¼ææè¢«å±éæ¦æªççº¿ç¨æä¼ç»§ç»å¹²æ´»ã #### åç `CyclicBarrier` å é¨éè¿ä¸ä¸ª `count` åéä½ä¸ºè®¡æ°å¨ï¼`count` çåå§å¼ä¸º `parties` 屿§çåå§åå¼ï¼æ¯å½ä¸ä¸ªçº¿ç¨å°äºæ æ è¿éäºï¼é£ä¹å°±å°è®¡æ°å¨å 1ã妿 count å¼ä¸º 0 äºï¼è¡¨ç¤ºè¿æ¯è¿ä¸ä»£æåä¸ä¸ªçº¿ç¨å°è¾¾æ æ ï¼å°±å°è¯æ§è¡æä»¬æé æ¹æ³ä¸è¾å ¥çä»»å¡ã ```java //æ¯æ¬¡æ¦æªççº¿ç¨æ° private final int parties; //计æ°å¨ private int count; ``` ä¸é¢æä»¬ç»åæºç æ¥ç®åççã 1ã`CyclicBarrier` é»è®¤çæé æ¹æ³æ¯ `CyclicBarrier(int parties)`ï¼å ¶åæ°è¡¨ç¤ºå±éæ¦æªççº¿ç¨æ°éï¼æ¯ä¸ªçº¿ç¨è°ç¨ `await()` æ¹æ³åè¯ `CyclicBarrier` æå·²ç»å°è¾¾äºå±éï¼ç¶åå½å线ç¨è¢«é»å¡ã ```java public CyclicBarrier(int parties) { this(parties, null); } public CyclicBarrier(int parties, Runnable barrierAction) { if (parties <= 0) throw new IllegalArgumentException(); this.parties = parties; this.count = parties; this.barrierCommand = barrierAction; } ``` å ¶ä¸ï¼`parties` å°±ä»£è¡¨äºææ¦æªç线ç¨çæ°éï¼å½æ¦æªççº¿ç¨æ°éè¾¾å°è¿ä¸ªå¼çæ¶åå°±æå¼æ æ ï¼è®©ææçº¿ç¨éè¿ã 2ãå½è°ç¨ `CyclicBarrier` 对象è°ç¨ `await()` æ¹æ³æ¶ï¼å®é ä¸è°ç¨çæ¯ `dowait(false, 0L)`æ¹æ³ã `await()` æ¹æ³å°±åæ ç«èµ·ä¸ä¸ªæ æ çè¡ä¸ºä¸æ ·ï¼å°çº¿ç¨æ¡ä½äºï¼å½æ¦ä½ççº¿ç¨æ°éè¾¾å° `parties` ç弿¶ï¼æ æ æä¼æå¼ï¼çº¿ç¨æå¾ä»¥éè¿æ§è¡ã ```java public int await() throws InterruptedException, BrokenBarrierException { try { return dowait(false, 0L); } catch (TimeoutException toe) { throw new Error(toe); // cannot happen } } ``` `dowait(false, 0L)`æ¹æ³æºç åæå¦ä¸ï¼ ```java // å½çº¿ç¨æ°éæè è¯·æ±æ°éè¾¾å° count æ¶ await ä¹åçæ¹æ³æä¼è¢«æ§è¡ãä¸é¢ç示ä¾ä¸ count çå¼å°±ä¸º 5ã private int count; /** * Main barrier code, covering the various policies. */ private int dowait(boolean timed, long nanos) throws InterruptedException, BrokenBarrierException, TimeoutException { final ReentrantLock lock = this.lock; // éä½ lock.lock(); try { final Generation g = generation; if (g.broken) throw new BrokenBarrierException(); // å¦æçº¿ç¨ä¸æäºï¼æåºå¼å¸¸ if (Thread.interrupted()) { breakBarrier(); throw new InterruptedException(); } // coutå1 int index = --count; // å½ count æ°éå为 0 ä¹å说ææåä¸ä¸ªçº¿ç¨å·²ç»å°è¾¾æ æ äºï¼ä¹å°±æ¯è¾¾å°äºå¯ä»¥æ§è¡await æ¹æ³ä¹åçæ¡ä»¶ if (index == 0) { // tripped boolean ranAction = false; try { final Runnable command = barrierCommand; if (command != null) command.run(); ranAction = true; // å° count é置为 parties 屿§çåå§åå¼ // å¤éä¹åçå¾ ççº¿ç¨ // ä¸ä¸æ³¢æ§è¡å¼å§ nextGeneration(); return 0; } finally { if (!ranAction) breakBarrier(); } } // loop until tripped, broken, interrupted, or timed out for (;;) { try { if (!timed) trip.await(); else if (nanos > 0L) nanos = trip.awaitNanos(nanos); } catch (InterruptedException ie) { if (g == generation && ! g.broken) { breakBarrier(); throw ie; } else { // We're about to finish waiting even if we had not // been interrupted, so this interrupt is deemed to // "belong" to subsequent execution. Thread.currentThread().interrupt(); } } if (g.broken) throw new BrokenBarrierException(); if (g != generation) return index; if (timed && nanos <= 0L) { breakBarrier(); throw new TimeoutException(); } } } finally { lock.unlock(); } } ``` #### 宿 ç¤ºä¾ 1ï¼ ```java /** * * @author Snailclimb * @date 2018å¹´10æ1æ¥ * @Description: æµè¯ CyclicBarrier ç±»ä¸å¸¦åæ°ç await() æ¹æ³ */ public class CyclicBarrierExample1 { // 请æ±çæ°é private static final int threadCount = 550; // éè¦åæ¥ççº¿ç¨æ°é private static final CyclicBarrier cyclicBarrier = new CyclicBarrier(5); public static void main(String[] args) throws InterruptedException { // åå»ºçº¿ç¨æ± ExecutorService threadPool = Executors.newFixedThreadPool(10); for (int i = 0; i < threadCount; i++) { final int threadNum = i; Thread.sleep(1000); threadPool.execute(() -> { try { test(threadNum); } catch (InterruptedException e) { // TODO Auto-generated catch block e.printStackTrace(); } catch (BrokenBarrierException e) { // TODO Auto-generated catch block e.printStackTrace(); } }); } threadPool.shutdown(); } public static void test(int threadnum) throws InterruptedException, BrokenBarrierException { System.out.println("threadnum:" + threadnum + "is ready"); try { /**çå¾ 60ç§ï¼ä¿è¯å线ç¨å®å ¨æ§è¡ç»æ*/ cyclicBarrier.await(60, TimeUnit.SECONDS); } catch (Exception e) { System.out.println("-----CyclicBarrierException------"); } System.out.println("threadnum:" + threadnum + "is finish"); } } ``` è¿è¡ç»æï¼å¦ä¸ï¼ ``` threadnum:0is ready threadnum:1is ready threadnum:2is ready threadnum:3is ready threadnum:4is ready threadnum:4is finish threadnum:0is finish threadnum:1is finish threadnum:2is finish threadnum:3is finish threadnum:5is ready threadnum:6is ready threadnum:7is ready threadnum:8is ready threadnum:9is ready threadnum:9is finish threadnum:5is finish threadnum:8is finish threadnum:7is finish threadnum:6is finish ...... ``` å¯ä»¥çå°å½çº¿ç¨æ°éä¹å°±æ¯è¯·æ±æ°éè¾¾å°æä»¬å®ä¹ç 5 ä¸ªçæ¶åï¼ `await()` æ¹æ³ä¹åçæ¹æ³æè¢«æ§è¡ã å¦å¤ï¼`CyclicBarrier` è¿æä¾ä¸ä¸ªæ´é«çº§çæé 彿° `CyclicBarrier(int parties, Runnable barrierAction)`ï¼ç¨äºå¨çº¿ç¨å°è¾¾å±éæ¶ï¼ä¼å æ§è¡ `barrierAction`ï¼æ¹ä¾¿å¤çæ´å¤æçä¸å¡åºæ¯ã ç¤ºä¾ 2ï¼ ```java /** * * @author SnailClimb * @date 2018å¹´10æ1æ¥ * @Description: æ°å»º CyclicBarrier çæ¶åæå®ä¸ä¸ª Runnable */ public class CyclicBarrierExample2 { // 请æ±çæ°é private static final int threadCount = 550; // éè¦åæ¥ççº¿ç¨æ°é private static final CyclicBarrier cyclicBarrier = new CyclicBarrier(5, () -> { System.out.println("------å½çº¿ç¨æ°è¾¾å°ä¹åï¼ä¼å æ§è¡------"); }); public static void main(String[] args) throws InterruptedException { // åå»ºçº¿ç¨æ± ExecutorService threadPool = Executors.newFixedThreadPool(10); for (int i = 0; i < threadCount; i++) { final int threadNum = i; Thread.sleep(1000); threadPool.execute(() -> { try { test(threadNum); } catch (InterruptedException e) { // TODO Auto-generated catch block e.printStackTrace(); } catch (BrokenBarrierException e) { // TODO Auto-generated catch block e.printStackTrace(); } }); } threadPool.shutdown(); } public static void test(int threadnum) throws InterruptedException, BrokenBarrierException { System.out.println("threadnum:" + threadnum + "is ready"); cyclicBarrier.await(); System.out.println("threadnum:" + threadnum + "is finish"); } } ``` è¿è¡ç»æï¼å¦ä¸ï¼ ``` threadnum:0is ready threadnum:1is ready threadnum:2is ready threadnum:3is ready threadnum:4is ready ------å½çº¿ç¨æ°è¾¾å°ä¹åï¼ä¼å æ§è¡------ threadnum:4is finish threadnum:0is finish threadnum:2is finish threadnum:1is finish threadnum:3is finish threadnum:5is ready threadnum:6is ready threadnum:7is ready threadnum:8is ready threadnum:9is ready ------å½çº¿ç¨æ°è¾¾å°ä¹åï¼ä¼å æ§è¡------ threadnum:9is finish threadnum:5is finish threadnum:6is finish threadnum:8is finish threadnum:7is finish ...... ```