ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Java并发同步器源码解析:AQS与Semaphore/CountDownLatch/CyclicBarrier原理

Java并发同步器源码解析:AQS与Semaphore/CountDownLatch/CyclicBarrier原理 先说个题外话。最近技术社区里“源码”相关的内容热度一直不减但真正值得反复翻的源码其实来来去去就那么几个。Java并发包里的Semaphore、CountDownLatch、CyclicBarrier这三个同步器属于典型的“面试必背、背完就忘”的组件很多人能用会写但问到实现原理就卡壳。我自己读源码的习惯是先把AQSAbstractQueuedSynchronizer这个底层框架吃透再看这三个同步器只是“几层皮”而已。这篇文章就把三者的源码链路梳理一遍讲清楚state、CLH队列、Condition这些核心机制是怎么配合的同时附上我在实际项目中用它们做限流、压测、分批计算的踩坑经验。无论你是刚开始学并发的小白还是准备面试跳槽的开发者希望这篇能帮你把浮在表面的“会用”变成“心里有底”的“懂原理”。1. 先啃地基AQS为什么是这三个同步器的共同底牌1.1 一小时搞懂AQS的两大核心state与CLH变种队列AQS这个东西说白了就是一把“并发锁的模板”。它内部维护了一个volatile int state和一个双向等待队列。state的含义由子类自己定Semaphore拿它当许可证数量CountDownLatch拿它当计数器ReentrantLock拿它当重入次数。队列则是CLH锁的一个变种每个等待的线程被封装成一个Node节点挂在队尾通过前驱节点的状态判断自己是否可以尝试获取资源。这里的Node节点有两个关键字段waitStatus和thread。waitStatus是一个int取值有CANCELLED1、SIGNAL-1、CONDITION-2、PROPAGATE-3。简单理解SIGNAL表示“当前节点释放时需要唤醒后继节点”这是整个队列得以“接力跑”的核心。CANCELLED表示节点因为超时或中断被取消了释放资源时会从队列里踢出去。队列的进出规则也很固定获取资源失败的线程通过addWaiter(Node.SHARED)或addWaiter(Node.EXCLUSIVE)入队然后在一个for循环里反复尝试如果前驱是head说明轮到它了再试一次获取资源成功就把自己变成新的head。整个过程没有使用synchronized全靠CAS保证并发安全所以性能才能撑得住高并发场景。1.2 模板方法模式tryAcquireShared与tryReleaseSharedAQS本身的acquireShared和releaseShared不会直接操作业务状态而是调用子类实现的两个方法。这就是模板方法模式也是读源码最舒服的地方——你只需要盯住子类重写的那几个方法就能理解整个同步逻辑。以共享模式为例AQS的acquireShared(int arg)大概长这样public final void acquireShared(int arg) { if (tryAcquireShared(arg) 0) { doAcquireSharedInterruptibly(arg); } }tryAcquireShared返回负数表示资源不足需要入队等待返回非负数表示获取成功直接往下走。releaseShared(int arg)则负责释放public final boolean releaseShared(int arg) { if (tryReleaseShared(arg)) { doReleaseShared(); return true; } return false; }doReleaseShared()的职责是唤醒队列中阻塞的线程它内部用CAS把head的waitStatus从SIGNAL改成0再调用LockSupport.unpark唤醒后继节点。这里有个细节唤醒操作会从head向后传播因为共享模式下可能有多个线程同时被唤醒这也是Semaphore能同时放行多个线程的原因。一直有人问“这三个同步器哪个最难”我的看法是理解顺序应该是先AQS再Semaphore再CountDownLatch最后CyclicBarrier。因为前两个是标准的AQS共享模式实现第三四个虽然名字像但CyclicBarrier压根没用AQS这点后面细说。2. Semaphore扒皮信号量如何用state完成限流与公平调度2.1 核心模型一个state当N张许可证Semaphore翻译成“信号量”确实贴切。它的构造器里直接setState(permits)state的大小就是许可证的总数。线程调用acquire()时其实是在“借”一张许可证调用release()时是把许可证“还”回去。如果state已经降到0再有人来借就只能去队列里排队等着。Semaphore有一个很有意思的设计它区分了公平和非公平两种模式。默认是非公平的对应的内部类是NonfairSync它的tryAcquireShared实现如下final int nonfairTryAcquireShared(int acquires) { for (;;) { int available getState(); int remaining available - acquires; if (remaining 0 || compareAndSetState(available, remaining)) { return remaining; } } }这段代码的逻辑是先读出当前state算出剩余许可数。如果剩余为负说明许可证不够直接返回负数如果够就用CAS把state更新成剩余值然后返回剩余值。这里return的语义要特别注意——返回负数代表失败返回非负数代表成功非负数的大小还暗示了“现在还剩多少许可证”。公平版的FairSync则多了一步检查protected int tryAcquireShared(int acquires) { for (;;) { if (hasQueuedPredecessors()) { return -1; } int available getState(); int remaining available - acquires; if (remaining 0 || compareAndSetState(available, remaining)) { return remaining; } } }hasQueuedPredecessors()就是检查等待队列里有没有比自己排得更早的线程有就直接返回-1不去抢。所以公平模式下新建线程必须老实排队避免“插队”现象。我在实际使用中默认都选非公平因为公平模式在大量线程竞争时会多出一次队列查询吞吐量会差一截。2.2 许可证的借与还acquire和release链路解析acquire()的入口是sync.acquireSharedInterruptibly(1)这个方法会响应中断。获取失败后进入doAcquireSharedInterruptibly线程被LockSupport.park挂起等着被唤醒。所以Semaphore的等待是“挂起-唤醒”式的阻塞不是自旋空转。release()的入口是sync.releaseShared(1)它的tryReleaseShared实现如下protected final boolean tryReleaseShared(int releases) { for (;;) { int current getState(); int next current releases; if (next current) { // overflow throw new Error(Maximum permit count exceeded); } if (compareAndSetState(current, next)) { return true; } } }这里用到了自旋CAS因为可能有多个线程同时释放许可证CAS失败就重试直到成功。next current是防溢出保护正常情况下永远不会触发。有个坑值得提一下release()不需要调用方持有过许可证也就是说一个没调用过acquire()的线程也能调用release()结果就是state被不断加高。这在代码里容易埋雷比如finally块里忘了判断是否实际获取过许可就会导致许可证数量虚增。2.3 实操用Semaphore给数据库连接池做并发限流我举个我之前在项目里实际写过的例子。一个订单系统需要访问数据库但数据库连接数有限超过一定并发就会报连接超时。最简单的方案就是用Semaphore限制同时执行的SQL数量public class SqlGuard { private final Semaphore semaphore new Semaphore(5); public void execute(Connection conn, String sql) throws InterruptedException { semaphore.acquire(); try { try (Statement stmt conn.createStatement()) { stmt.execute(sql); } } finally { semaphore.release(); } } }这里的核心是acquire和release必须成对出现release必须放在finally里。否则SQL执行过程中抛异常许可证就泄漏了跑一会儿会发现所有线程都堵在acquire上。这是个非常经典的线上事故我见过不止一次。用tryAcquire可以做成“优雅降级”的限流拿不到许可就走兜底逻辑比如直接返回缓存数据而不是硬等。if (semaphore.tryAcquire(200, TimeUnit.MILLISECONDS)) { try { // 执行数据库操作 } finally { semaphore.release(); } } else { // 降级读缓存 return cache.get(key); }这种写法在网关层做接口限流很常见既保证核心链路不被打垮又不至于因为等待许可把线程全部拖死。Semaphore还有一个看着反直觉的用法初始化为0配合release()做“开关”。所有线程都先执行acquire()挂起另一个线程执行N次release()逐个放行。不过这个场景我用CountDownLatch更多毕竟那个语义更清晰。2.4 避坑清单许可证泄漏、不可重入、公平性误判写Semaphore踩过的坑我整理成几条许可证泄漏acquire之后没在finally释放或者提前return了导致可用许可越来越少最后系统“假死”。排查办法是监控availablePermits()正常情况下应该回到初始值。Semaphore不可重入同一线程里再次acquire也会继续扣减许可不会像ReentrantLock那样放行。所以递归或循环里用Semaphore要特别小心容易把自己堵死。不要把限流数设得和连接数一样大如果底层连接池只有5个连接而Semaphore设成10等待许可的线程过了这一关还是会堵在连接池上等于白限流。我通常是设为连接数的70%左右留出余量。公平模式不等于绝对公平公平模式只是避免了新来的线程插队但如果是“先唤醒的线程被中断了下一个线程还没被唤醒”这种窗口期还是会有些许不公。并发场景没有银弹要看业务容忍度。3. CountDownLatch拆解一次性闸门背后的CLH队列等待机制3.1 state从count到0countDown与await源码精读CountDownLatch的构造器同样是把state设置成初始count。它有两个核心方法countDown()和await()。await()的实现是sync.acquireSharedInterruptibly(1)其内部的tryAcquireShared很简单protected int tryAcquireShared(int acquires) { return (getState() 0) ? 1 : -1; }只要state不是0就返回-1调用的线程进入等待队列被挂起。state数到0之后所有等待的线程都会被一次性放行。这里有个关键点countDown()用的也是releaseShared它的tryReleaseShared里有一个非常重要的判断返回值逻辑protected boolean tryReleaseShared(int releases) { for (;;) { int c getState(); if (c 0) { return false; } int nextc c - 1; if (compareAndSetState(c, nextc)) { return nextc 0; } } }注意这两处细节第一state已经到0时countDown()不会再做什么返回值是false第二只有当CAS成功且nextc 0时tryReleaseShared才返回true触发后续的doReleaseShared()去唤醒等待队列里的线程。也就是说最后一次countDown的线程承担了“唤醒所有等待者”的职责。3.2 等待队列里发生了什么doAcquireSharedInterruptibly逐行拆解await()真正复杂的部分在于等待逻辑。看AQS的doAcquireSharedInterruptiblyprivate void doAcquireSharedInterruptibly(int arg) throws InterruptedException { final Node node addWaiter(Node.SHARED); try { for (;;) { final Node p node.predecessor(); if (p head) { int r tryAcquireShared(arg); if (r 0) { setHeadAndPropagate(node, r); p.next null; return; } } if (shouldParkAfterFailedAcquire(p, node) parkAndCheckInterrupt()) { throw new InterruptedException(); } } } catch (Throwable t) { cancelAcquire(node); throw t; } }这段代码是整个等待机制的核心。我逐行解释一下addWaiter(Node.SHARED)创建当前线程的节点以SHARED模式挂到队列尾部。进入for循环后先看自己的前驱是不是head。如果是说明自己排到了队首再尝试一次获取资源。tryAcquireShared此时会再次检查stateCountDownLatch就是看state是否为0。如果获取成功执行setHeadAndPropagate(node, r)把当前节点设为新head。PROPAGATE状态就是在这个方法里传播的目的是让“闸门打开”的信号能连续传递下去唤醒后续所有等待的共享节点。如果获取失败执行shouldParkAfterFailedAcquire把前驱节点的waitStatus设为SIGNAL然后parkAndCheckInterrupt真正挂起线程。我一开始读这段代码的时候一直想不通为什么await()明明是个“一次性”的操作还要用循环反复尝试后来才明白挂起中的线程被唤醒后可能并不满足获取条件比如被虚假唤醒必须回到循环开头重新判断。这种循环加CAS的写法正是无锁编程的标准范式。3.3 实操CountDownLatch模拟并发压测以及多模块初始化协调CountDownLatch最常见的场景是“主线程等所有子任务完成”但其实用它做“同时起跑”也特别顺。我做过一个接口压测小工具需要让500个线程在同一时刻发请求。实现很简单两个CountDownLatch互相配合。int threadCount 500; CountDownLatch ready new CountDownLatch(threadCount); CountDownLatch start new CountDownLatch(1); CountDownLatch end new CountDownLatch(threadCount); for (int i 0; i threadCount; i) { new Thread(() - { ready.countDown(); // 通知主线程我准备好了 try { start.await(); // 阻塞等待开跑信号 // 在这里发送HTTP请求记录耗时 } catch (InterruptedException e) { Thread.currentThread().interrupt(); } finally { end.countDown(); // 通知主线程我跑完了 } }).start(); } ready.await(); // 等待所有线程就位 long begin System.currentTimeMillis(); start.countDown(); // 放行所有线程 end.await(); // 等待所有线程完成 long cost System.currentTimeMillis() - begin;这里start这个CountDownLatch初始值是1主线程执行一次countDown()就能让500个线程从await()里同时醒来。实际测试下来线程唤醒的延迟一般在几毫秒内对于绝大多数压测场景足够精确了。如果你有更严格的时间对齐要求那得上专门的压测框架而不是手动写Latch。另一个常见场景是服务启动时的多模块预热。比如一个服务要同时预热本地缓存、连接池、线程池三件事可以并行做等全部完成后才对外提供服务。用CountDownLatch写出来的代码非常直观而且好维护。3.4 拉闸不能复位为什么CountDownLatch是一次性的这个特性经常被拿来和CyclicBarrier做对比。CountDownLatch的state一旦减到0它的tryAcquireShared就永远返回非负数了后面的countDown()不会再产生任何影响。从源码里也能看到state 0时tryReleaseShared直接返回false后续代码直接跳过。这意味着如果你需要“多轮等待”能力不能用CountDownLatch硬撑要么每次创建新实例要么换成CyclicBarrier。我在项目里见到有人用for循环给CountDownLatch每个批次都new一个对象这没问题就是对象生命周期没管理好容易出小毛病。反观CyclicBarrier把这个需求内置了后面会讲到它如何通过generation实现复用。3.5 避坑清单countDown遗漏、超时设置、中断恢复countDown必须放在finally里如果任务执行抛异常当前线程直接退出count就永远少了一个主线程会一直挂在await()上。线上最常见的是线程池里任务被拒绝或者业务异常被吞掉导致计数一直凑不齐。await必须加超时尤其是跨系统调用的场景谁也不能保证子任务一定成功完成。我会习惯写成await(30, TimeUnit.SECONDS)超时后走降级逻辑而不是无限等。子线程捕获InterruptedException后要恢复中断标志很多代码在catch (InterruptedException e)里只打印日志中断信号就丢了。正确做法是Thread.currentThread().interrupt()把中断状态还回去。小心“信号提前丢失”如果某个子线程在countDown()之前就因为异常退出了主线程这边无法感知到具体是哪个任务失败只能靠超时兜底。所以任务内最好有明确的成功/失败标记配合结果汇总而不是只依赖计数。4. CyclicBarrier解剖它为什么不用AQS而是LockCondition4.1 换个思路ReentrantLock Condition如何实现“人齐发车”CyclicBarrier是我认为三个同步器里设计最有意思的一个。它没有用AQS而是组合了ReentrantLock Condition。为什么不用AQS你可以从语义上琢磨一下CountDownLatch等的是“数归零”Semaphore等的是“许可证”这些都可以抽象成“某个state满足条件就放行”。但CyclicBarrier等的是“所有参与线程到达同一个屏障点”而且到达之后会自动重置进入下一轮。这种“多轮协作”的场景借助Condition的等待/通知机制实现起来更直白。CyclicBarrier内部有几个关键字段lock一个ReentrantLock所有对屏障状态的操作都要拿锁。triplock.newCondition()等待线程的休息室。generation内部类Generation记录当前是第几轮。count当前轮还差多少个线程到达初始是parties。barrierCommand每次屏障打破时执行的Runnable任务由最后一个到达的线程执行。4.2 await源码精读最后一个线程做了什么每个参与者都会调用await()真正干活的都在dowait(boolean timed, long nanos)里。我把关键路径拆出来了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(); } int index --count; if (index 0) { // 我就是最后一个到达的线程 boolean ranAction false; try { final Runnable command barrierCommand; if (command ! null) { command.run(); } ranAction true; nextGeneration(); // 唤醒所有等待线程并开启新一轮 return 0; } finally { if (!ranAction) { breakBarrier(); // 如果barrierCommand执行失败也要打破屏障 } } } // 不是最后一个就进入等待循环 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 { 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(); } }这段代码信息量很大我挑三个重点第一最后一个线程执行了额外的两件事执行barrierCommand如果有的话然后调用nextGeneration()它内部先trip.signalAll()唤醒所有等待线程再把count重置回parties最后generation new Generation()。整个过程都在持锁状态下完成保证原子性。第二**等待线程被唤醒后怎么知道自己是“正常完成”还是“屏障坏了”**它拿到Generation g后先检查g.broken如果被打破就抛BrokenBarrierException再检查g ! generation不等说明已经换代说明上一轮成功结束安全返回。这两层检查就是中断、超时与正常唤醒之间的分水岭。第三中断与超时的处理规则如果某个等待线程中断了它会调用breakBarrier()把屏障标记为broken并唤醒所有伙伴其他线程随后要么抛BrokenBarrierException要么继续被唤醒后也抛异常。所以栅栏只要有一个参与者“出事”整批人都会知道。Personal note我在读这段代码前一直以为CyclicBarrier内部也维护了一个计数器且复用的是state。看完才发现它用generation解决了“多轮复用”这个难题这个抽象非常干净。4.3 实操CyclicBarrier做分批并行计算多次复用同一个屏障讲一个我实际接触过的数据处理需求一个定时任务需要把一批数据分批次写入外部系统每批要等四个线程都准备好了再合并批次之间有依赖必须按顺序处理。用CyclicBarrier实现非常顺int batchThreads 4; CyclicBarrier barrier new CyclicBarrier(batchThreads, () - { // 这个回调会由每批最后一个到达的线程执行 // 在这里做汇总、落库、发消息 System.out.println(第 batchNo.incrementAndGet() 批完成); }); // 每个工作线程的run方法大概长这样 for (int batch 0; batch batchCount; batch) { // 1. 处理属于自己的那一段数据 processSlice(batch); // 2. 等本批其他线程都处理完 barrier.await(); // 3. 屏障被打破后所有线程自动进入下一轮 }第一轮跑完后所有线程在barrier.await()返回这里的关键是返回后线程并没有被销毁而是继续执行for循环进入下一批。所以CyclicBarrier特别适合“多个线程在同一份数据上分片处理每处理完一片同步一次然后继续”的多轮任务。另外要注意回调barrierCommand是同步执行的由最后一个到达的线程执行。如果这个回调很耗时它会拖累整批线程的进度因为只有它返回后nextGeneration()才会执行其他线程才能被唤醒。所以回调里尽量别做重IO要么异步化要么把汇总逻辑拆到后续步骤。4.4 避坑清单broken状态、超时、reset时机CyclicBarrier的坑比前两个更隐蔽我列一下我踩过的BrokenBarrierException不是偶发异常是“系统性故障”的预警只要有一个线程中断或超时整个屏障就broken了其他线程全部抛异常。排查这类问题不要只看单个线程的日志要看全量线程的异常分布。reset() 不要在等待过程中调用源码里reset()先执行breakBarrier()再nextGeneration()如果此时有其他线程正在await()它们会因为broken抛异常。正确姿势是先确认没有线程在等待或者通过try-catch兜底。await别裸奔超时一定要加我给合作方讲一次事故就是因为有个线程长时间卡在IO上屏障迟迟凑不齐其他三个线程干等。后来全换成await(10, TimeUnit.SECONDS)超时后做降级或重试才稳住。parties的数量必须和实际参与线程数一致多喊一个少喊一个都容易出问题。少一个还行最后一个线程会等很久多一个就直接错过“人齐”的判断永远等不到。这个从源码逻辑上想就明白了--count到0才会换代总数不对永远到不了0。子线程数比parties少时的处理如果线程池只有3个线程而parties是4那屏障等不到第4个线程会一直挂起。遇到这种情况可以用超时重试的方式兜底或者干脆用CountDownLatch更合适。5. 三兄弟怎么选一张对比表搞定并发场景的选型判断5.1 底层实现与核心行为对比把三个组件放在一张表里差异一目了然维度SemaphoreCountDownLatchCyclicBarrier底层实现AQS共享模式AQS共享模式ReentrantLock Condition核心状态许可证数量state计数器state归零放行参与线程数count generation可重用性可反复acquire/release一次性不能重置可重复使用自动换代等待线程数多个线程同时等待许可多个线程同时等待归零多个线程互相等待人齐放行中断处理抛出InterruptedException抛出InterruptedException中断会让屏障broken伙伴线程抛BrokenBarrierException超时处理tryAcquire(timeout)await(timeout)await(timeout)超时打破屏障是否区分公平性支持公平/非公平不区分不区分典型场景限流、资源池、信号量锁多任务汇总、压测同时起跑、服务预热多轮分批计算、并行任务对齐后继续从表上能看出CyclicBarrier和前两者完全不是一个路由前两个是“容器里的状态变化放行外部线程”CyclicBarrier是“参与线程互相帮扶到位后一起走”。5.2 选型逻辑先想清楚你的业务到底需要“等待什么”我在实际写代码的时候会先问自己三个问题我到底在等一个什么条件如果条件是“某个共享资源数量够不够”选Semaphore如果条件是“一坨子任务有没有全部完事”选CountDownLatch如果条件是“每个参与线程有没有各自到达本批次的终点”选CyclicBarrier。这个流程是一次性的还是多轮的一轮定生死选CountDownLatch需要反复同步选CyclicBarrier。等待的线程是“被动等被通知”还是“主动互相等”前者是CountDownLatch/Semaphore有一个释放方后者是CyclicBarrier没有专门的释放方每个线程都是参与者。有朋友问过“CyclicBarrier能替代CountDownLatch吗”严格讲不行。CountDownLatch等的是别人被等待对象不感知等待者CyclicBarrier等的是自己人每个线程都知道别人在等自己。语义不一样硬套很容易出bug。5.3 源码阅读方法论我是怎么顺着一条主线把并发包读薄的最后聊点学习方法。我自己读并发包源码的路径是这样的先找最典型的AQS实现ReentrantLock或Semaphore读它的tryAcquire、tryRelease、tryAcquireShared、tryReleaseShared把AQS的模板方法跑通。再看AQS的doAcquireSharedInterruptibly和doReleaseShared理解排队、挂起、唤醒、传播这几步。最后看CyclicBarrier用“Lock Condition 换代”的思路去理解它和AQS的差异。读源码时我习惯先用一个demo把流程跑起来然后打断点看线程状态。比如CountDownLatch的demo主线程和子线程都打上await()和countDown()的断点观察主线程在LockSupport.park前后的状态一下子就能把“挂起-唤醒”这个过程看成动态的。纸面上读十遍不如断点盯一遍。如果你对这个主题有兴趣下一步可以顺手把ReentrantLock源码也读了因为你读CyclicBarrier时已经接触了Condition而ConditionObject本身也是AQS内部类两者的关联会让你对AQS的理解再上一层楼。我个人在实际操作中还有一个体会这三种同步器虽然常被当成面试题但它们背后的设计思路——模板方法、状态机、队列化等待、Condition与中断协作——才是真正值钱的东西。把这些吃透了以后碰到再冷门的并发组件也基本能蒙着源码猜出八九分。
返回列表