ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Java Semaphore 信号量源码解析:基于 AQS 的并发许可控制与限流实现

Java Semaphore 信号量源码解析:基于 AQS 的并发许可控制与限流实现 Java Semaphore 信号量源码解析基于 AQS 的并发许可控制与限流实现【免费下载链接】source-code-hunter 从源码层面剖析挖掘互联网行业主流技术的底层实现原理为广大开发者 “提升技术深度” 提供便利。目前开放 Spring 全家桶Mybatis、Netty、Dubbo 框架及 Redis、Tomcat 中间件等项目地址: https://gitcode.com/GitHub_Trending/so/source-code-hunterSemaphore信号量是 JUC 中用于控制一定时间内并发执行线程数的同步工具其内部完全基于 AbstractQueuedSynchronizerAQS实现在 source-code-hunter 仓库的 Semaphore.md 中给出了核心内部类 Sync、公平/非公平模式以及 acquire/release 全套 API 的源码级剖析。本文以该文档为主体骨架结合仓库中 详解AbstractQueuedSynchronizer.md 对 AQS 共享锁机制的讲解深入拆解 Semaphore 的许可获取、释放、唤醒传播等底层原理并给出网关限流、连接数限制等可直接落地的实战示例。读完本文你将掌握 Semaphore 的完整源码脉络理解公平与非公平模式的本质差异并能基于它写出正确的限流与资源管控代码。Semaphore 是什么信号量的核心语义Semaphore 翻译为“信号量”用于控制一定时间内并发执行的线程数。它内部维护一个许可permit计数器线程执行任务前必须先获取一个或多个许可许可耗尽后后续线程将被阻塞排队直到其他线程释放许可。其典型应用场景包括网关限流限制某一接口同一时刻最多有多少个请求进入处理逻辑超出部分排队或拒绝资源限制限制可同时发起的数据库连接数、HTTP 连接数、线程池外的额外资源占用等互斥退化为二进制信号量当初始许可数为 1 时Semaphore 退化为一把互斥锁但需要注意的是它不具备可重入性同一线程重复 acquire 会阻塞自己。从 JUC 全量 UML 类图见上可以看到Semaphore 与 ReentrantLock、CountDownLatch 等工具一样都挂靠在 AbstractQueuedSynchronizer 这颗“大树”之下许可证的计数与线程排队复用 AQS 的同步状态state与同步队列。一个容易忽略的关键特性原文档已明确指出release() 释放许可时并未对释放许可数做限制因此可以通过该方法动态增加总的许可数量。这意味着 Semaphore 的许可总数不是固定不变的这一点在阅读tryReleaseShared的源码时会得到印证。Semaphore 与 AQS基于共享锁模式的实现Semaphore 的全部机制都委托给内部类Sync继承自AbstractQueuedSynchronizer这正是 AQS 的经典用法——通过子类覆写若干模板方法让 AQS 框架完成排队、阻塞、唤醒等通用逻辑。从行为模式看AQS 分为独占锁和共享锁两种模式独占锁同一时刻只允许一个线程持有如 ReentrantLock共享锁允许多个线程同时持有如 Semaphore、CountDownLatch、ReentrantReadWriteLock 的读锁。Semaphore 属于共享锁模式多个线程可以同时获取许可只要剩余许可数不为负。对应地它覆写的是 AQS 的tryAcquireShared/tryReleaseShared两个共享模式方法而 AQS 提供的acquireSharedInterruptibly、releaseShared、doAcquireSharedInterruptibly、doReleaseShared、setHeadAndPropagate等模板与工具方法则在 详解AbstractQueuedSynchronizer.md 中有完整剖析获取共享锁 / 释放共享锁的实现见该文档“获取共享锁的实现”与“释放共享锁的实现”两节。本文后续会反复回到这些方法说明它们如何与 Semaphore 的许可逻辑咬合。核心内部类 Sync许可状态的管理者先看 Semaphore 的骨架——所有许可逻辑都收敛在Sync中源码对应仓库文档 Semaphore.mdabstract static class Sync extends AbstractQueuedSynchronizer { private static final long serialVersionUID 1192457210091910933L; /* 赋值state为总许可数 */ Sync(int permits) { setState(permits); } /* 剩余许可数 */ final int getPermits() { return getState(); } /* 自旋 CAS非公平获取 */ final int nonfairTryAcquireShared(int acquires) { for (;;) { // 剩余可用许可数 int available getState(); // 本次获取许可后剩余许可 int remaining available - acquires; // 如果获取后剩余许可大于0则CAS更新剩余许可否则获取失败失败 if (remaining 0 || compareAndSetState(available, remaining)) return remaining; } } /** * 自旋 CAS 释放许可 * 由于未对释放许可数做限制所以可以通过release动态增加许可数量 */ 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); // CAS更新许可数量 if (compareAndSetState(current, next)) return true; } } /* 自旋 CAS 减少许可数量 */ final void reducePermits(int reductions) { for (;;) { // 当前剩余许可 int current getState(); // 更新值 int next current - reductions; // 较少许可数错误抛出异常 if (next current) // underflow throw new Error(Permit count underflow); // CAS更新许可数 if (compareAndSetState(current, next)) return; } } /* 丢弃所有许可 */ final int drainPermits() { for (;;) { int current getState(); if (current 0 || compareAndSetState(current, 0)) return current; } } }这里有几个值得深入的关键点许可数即 AQS 的同步状态 state。构造 Semaphore 时调用setState(permits)把初始许可数写入 AQS 的state之后getPermits()/getState()返回的永远是当前剩余许可数。AQS 对 state 提供了 volatile 可见性保证与 CAS 更新原语这是整个并发安全性的基石。获取许可 自旋 CAS 扣减。nonfairTryAcquireShared是一个无锁算法循环读取available计算出扣减后的remaining只有remaining 0且 CAS 成功才会返回否则继续自旋重试。返回值语义与 AQS 共享锁的约定完全一致详见 AQS 文档返回非负数表示获取成功0 表示成功但后继争用线程不会成功正数表示成功且后继也可能成功返回负数表示获取失败由上层据此决定是否入队阻塞。释放许可 自旋 CAS 累加且不设上限。tryReleaseShared直接把current releases写回 state。注意它只检查了整数溢出next current时抛出Error(Maximum permit count exceeded)却没有限制累加后的值不能超过初始许可数——这就是原文档强调的“可以通过 release 动态增加许可数量”的源码出处。这在业务上是一种灵活性例如允许运维在运行时临时扩容“许可池”让更多线程并发执行。许可的运维操作reducePermits(int reductions)自旋 CAS 扣减许可检查下溢即next current时抛Error(Permit count underflow)drainPermits()一次性把许可清零并返回清零前的数量可用于“暂停放行”等场景。这两个方法在Sync中以final方法存在对外通过Semaphore的公有 API如drainPermits()、availablePermits()暴露属子类可控的辅助能力。公平与非公平两种许可获取策略Semaphore通过构造器参数fair在FairSync与NonfairSync之间选择策略二者都继承自Sync源码对应 Semaphore.md/** * 非公平模式 */ static final class NonfairSync extends Sync { private static final long serialVersionUID -2694183684443567898L; NonfairSync(int permits) { super(permits); } protected int tryAcquireShared(int acquires) { return nonfairTryAcquireShared(acquires); } } /** * 公平模式 */ static final class FairSync extends Sync { private static final long serialVersionUID 2014338818796000944L; FairSync(int permits) { super(permits); } /** * 公平模式获取许可 * 公平模式不论许可是否充足都会判断同步队列中是否有线程在等地如果有获取失败排队阻塞 */ protected int tryAcquireShared(int acquires) { for (;;) { // 如果有线程在排队立即返回 if (hasQueuedPredecessors()) return -1; // 自旋 cas获取许可 int available getState(); int remaining available - acquires; if (remaining 0 || compareAndSetState(available, remaining)) return remaining; } } }两种策略的差异体现在tryAcquireShared的第一步模式获取许可策略适用场景非公平默认无论当前是否有线程在同步队列中排队都直接自旋 CAS 抢许可抢不到再入队追求吞吐量允许插队许可总量大、争用不激烈时效率高公平先调用hasQueuedPredecessors()判断同步队列中是否有线程在排队只要有排队线程无论许可是否充足都直接返回 -1老老实实入队追求公平性避免线程饥饿如必须按请求到达顺序放行一句话总结原文档的表述公平模式无论是否有许可都会先判断是否有线程在排队如果有线程排队则进入排队否则尝试获取许可非公平模式无论许可是否充足直接尝试获取许可。hasQueuedPredecessors()是 AQS 提供的队列探测方法返回 true 表示同步队列中存在排队线程当前线程排在队首前驱为 head 的情况除外。获取许可acquire 的完整调用链Semaphore对外提供两组获取许可的 API全部委托给内部的sync源码对应 Semaphore.md// --------------------- 获取许可 -------------------- /* 获取指定数量的许可 */ public void acquire(int permits) throws InterruptedException { if (permits 0) throw new IllegalArgumentException(); sync.acquireSharedInterruptibly(permits); } /* 获取一个许可 */ public void acquire() throws InterruptedException { sync.acquireSharedInterruptibly(1); } public final void acquireSharedInterruptibly(int arg) throws InterruptedException { if (Thread.interrupted()) throw new InterruptedException(); if (tryAcquireShared(arg) 0) // 获取许可剩余许可0则获取许可成功0获取许可失败进入排队 doAcquireSharedInterruptibly(arg); } protected int tryAcquireShared(int acquires) { return nonfairTryAcquireShared(acquires); } /** * return 剩余许可数量。非负数获取许可成功负数获取许可失败 */ final int nonfairTryAcquireShared(int acquires) { for (;;) { int available getState(); int remaining available - acquires; if (remaining 0 || compareAndSetState(available, remaining)) return remaining; } } /** * 获取许可失败当前线程进入同步队列排队阻塞 */ private void doAcquireSharedInterruptibly(int arg) throws InterruptedException { // 创建同步队列节点并入队 final Node node addWaiter(Node.SHARED); boolean failed true; try { for (;;) { // 如果当前节点是第二个节点尝试获取锁 final Node p node.predecessor(); if (p head) { int r tryAcquireShared(arg); if (r 0) { setHeadAndPropagate(node, r); p.next null; // help GC failed false; return; } } // 阻塞当前线程 if (shouldParkAfterFailedAcquire(p, node) parkAndCheckInterrupt()) throw new InterruptedException(); } } finally { if (failed) cancelAcquire(node); } }调用链拆解入口校验acquire(int permits)对负数参数抛IllegalArgumentException随后进入 AQS 的acquireSharedInterruptibly。中断响应acquireSharedInterruptibly首先检查Thread.interrupted()若线程已被中断则立即抛InterruptedException保证acquire的可中断语义。快速路径调用tryAcquireShared(arg)非公平模式即Sync.nonfairTryAcquireShared。返回值 0说明剩余许可充足且 CAS 成功线程直接持有许可继续执行无需排队。慢速路径入队阻塞返回值为负则进入doAcquireSharedInterruptibly——这是 AQS 共享模式的标准排队逻辑与 详解AbstractQueuedSynchronizer.md 中doAcquireShared的实现同源addWaiter(Node.SHARED)把当前线程包装成共享模式节点插入同步队列尾部循环中只允许前驱为 head 的节点尝试再次获取许可保证队列纪律防止队列内部插队获取成功则调用setHeadAndPropagate(node, r)——这是共享锁与独占锁的关键区别它不仅把自己设为新 head还会根据返回值 r 与节点状态决定是否继续唤醒后继的共享节点实现“许可有余量则依次放行多个线程”的传播效应获取失败则通过shouldParkAfterFailedAcquire把前驱节点状态置为 SIGNAL再由parkAndCheckInterrupt阻塞当前线程等待被unpark唤醒或响应中断若因中断退出循环finally中的cancelAcquire(node)会将该节点标记为 CANCELLED 并从队列中摘除。许可放行与唤醒传播的底层细节setHeadAndPropagate与doReleaseShared是共享锁“一放多醒”的核心其源码细节在 详解AbstractQueuedSynchronizer.md“获取共享锁的实现”一节中有完整注释。要点如下setHeadAndPropagate(node, r)设置新 head 后只要满足propagate 0还有剩余许可或 head 节点等待状态为 SIGNAL / PROPAGATE就检查后继节点是否为共享节点s.isShared()是则调用doReleaseShared唤醒之doReleaseShared循环读取 head若状态为 SIGNAL 则 CAS 清 0 后unparkSuccessor唤醒第一个等待线程若读到状态 0恰逢并发释放的中间态则 CAS 置为PROPAGATE由后续获取到许可的线程代为继续传播唤醒唤醒是链式传播的被唤醒的线程拿到许可成为新 head 后又会执行setHeadAndPropagate把唤醒继续向后传递直到队列中所有能拿到许可的线程都被放行。这套机制回答了“为什么 Semaphore 一次 release 可以让多个排队线程依次拿到许可”的问题唤醒沿着同步队列逐级传播而不是像独占锁那样一次只唤醒一个。释放许可release 与许可的动态扩容释放侧 API 与实现源码对应 Semaphore.md// --------------------- 释放归还许可 ------------------------- /* 释放指定数量的许可 */ public void release(int permits) { if (permits 0) throw new IllegalArgumentException(); sync.releaseShared(permits); } /* 释放一个许可 */ public void release() { sync.releaseShared(1); } public final boolean releaseShared(int arg) { // 归还许可成功 if (tryReleaseShared(arg)) { doReleaseShared(); return true; } return false; } /** * 释放许可 * 由于未对释放许可数做限制所以可以通过release动态增加许可数量 */ 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; } } private void doReleaseShared() { // 自旋唤醒等待的第一个线程(其他线程将由第一个线程向后传递唤醒) for (;;) { Node h head; if (h ! null h ! tail) { int ws h.waitStatus; if (ws Node.SIGNAL) { if (!compareAndSetWaitStatus(h, Node.SIGNAL, 0)) continue; // loop to recheck cases // 唤醒第一个等待线程 unparkSuccessor(h); } else if (ws 0 !compareAndSetWaitStatus(h, 0, Node.PROPAGATE)) continue; // loop on failed CAS } if (h head) // loop if head changed break; } }释放链路的要点release(int permits)同样对负数参数做防御性校验tryReleaseShared用自旋 CAS 把许可累加回去成功返回 true 后进入doReleaseShareddoReleaseShared唤醒同步队列中第一个等待线程head 的后继其余等待线程由被唤醒者沿队列向后传递唤醒注释“唤醒等待的第一个线程(其他线程将由第一个线程向后传递唤醒)”正是对这一传播机制的精确描述动态扩容tryReleaseShared只做溢出检查、不限制累加上限因此release()可以释放比初始许可更多的许可。例如初始许可为 3两个线程各 release 一次后许可数会变成 5后续即可容纳最多 5 个并发线程。这是原文档反复强调的核心行为也是 Semaphore 区别于“固定容量资源池”的一个重要特性——若业务需要严格固定上限必须自行约束 release 的调用次数或配合reducePermits收紧。其他常用 API 与注意事项结合Sync提供的底层能力Semaphore还向外暴露了一系列实用 API从源码结构可以确认其存在与语义API作用备注availablePermits()返回当前剩余许可数委托sync.getPermits()仅作监控参考非线程安全快照drainPermits()一次性清空所有许可并返回清空前的数量可用于“熔断放行”拒绝一切新任务直到许可被恢复reducePermits(int)减少指定数量许可底层为Sync.reducePermits下溢时抛ErrorisFair()判断是否为公平模式返回sync instanceof FairSynchasQueuedThreads()/getQueueLength()是否有排队线程 / 排队线程数委托 AQS 队列查询能力用于监控使用注意事项许可数必须非负构造器Semaphore(int permits)及Semaphore(int permits, boolean fair)在 permits 为负时会抛IllegalArgumentExceptionacquire 与 release 要成对与锁不同Semaphore 不强制要求由同一线程释放许可这也正是它能跨线程传递许可的原因但业务上必须保证 try/finally 或 try-with-resources 模式成对调用否则许可泄漏会导致可用并发数持续下降注意中断与不可中断变体acquire是可中断的若不想被中断打断排队可使用acquireUninterruptibly()tryAcquire()则是非阻塞尝试拿不到许可立即返回 false。实战网关限流与连接数控制下面给出两个贴合原文档应用场景网关限流、资源限制如最大可发起连接数的完整示例。示例一网关限流限制某个接口同一时刻最多放行 3 个请求超出部分阻塞等待import java.util.concurrent.Semaphore; public class GatewayLimiter { // 初始许可 3公平模式保证请求按到达顺序放行 private final Semaphore semaphore new Semaphore(3, true); public void handleRequest(String requestId) throws InterruptedException { // 排队等待许可若线程被中断则放弃本次请求 semaphore.acquire(); try { System.out.println([ requestId ] 进入处理剩余许可: semaphore.availablePermits()); // 模拟业务处理 Thread.sleep(500); } finally { // 无论业务是否异常必须归还许可 semaphore.release(); } } public static void main(String[] args) { GatewayLimiter limiter new GatewayLimiter(); for (int i 1; i 10; i) { final String id req- i; new Thread(() - { try { limiter.handleRequest(id); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }).start(); } } }运行后观察输出可以看到任意时刻最多只有 3 个线程处于“进入处理”状态其余线程排队等待这就是 Semaphore 限流的最直观效果。示例二限制最大并发连接数并动态扩容模拟连接池初始最多 5 个连接使用过程中通过额外release实现“扩容”import java.util.concurrent.Semaphore; public class ConnectionLimiter { private final Semaphore semaphore new Semaphore(5); public void acquireConnection() throws InterruptedException { semaphore.acquire(); System.out.println(获取连接成功当前可用许可: semaphore.availablePermits()); } public void releaseConnection() { semaphore.release(); } public void expandCapacity(int extra) { // 动态增加许可总数扩容连接池 semaphore.release(extra); System.out.println(扩容 extra 个许可当前可用许可: semaphore.availablePermits()); } public void shrinkCapacity(int reduction) { // 通过 reducePermits 收紧容量 semaphore.reducePermits(reduction); System.out.println(缩减 reduction 个许可当前可用许可: semaphore.availablePermits()); } }这个例子演示了原文档强调的“release 可动态增加许可数量”特性以及配套的reducePermits收紧手段——在真实系统中这对应着运行时调整资源池容量的弹性扩缩容需求。总结Semaphore 设计要点一览设计点说明源码位置许可计数直接复用 AQS 的同步状态 statesetState(permits)初始化Semaphore.md并发安全全程自旋 CAS无锁实现许可的扣减与累加Semaphore.md公平性FairSync通过hasQueuedPredecessors()保证先来先得默认NonfairSync直接抢Semaphore.md排队阻塞失败线程以 SHARED 节点入 AQS 同步队列park阻塞、unpark唤醒详解AbstractQueuedSynchronizer.md唤醒传播setHeadAndPropagatedoReleaseShared PROPAGATE 状态实现“一放多醒”详解AbstractQueuedSynchronizer.md动态扩容tryReleaseShared仅做溢出检查release 可无上限累加许可Semaphore.mdSemaphore 是理解 AQS 共享锁模式的绝佳样本它用最简的许可状态 共享节点队列支撑起限流、资源管控等高频业务需求。其“许可可动态增长”的语义在并发工具中独树一帜使用时务必与业务容量模型对齐。若想进一步吃透其底层排队与唤醒机制建议配合仓库中的 详解AbstractQueuedSynchronizer.md重点关注共享锁的获取与释放以及 Lock锁组件.md独占锁对照一起阅读JUC并发包UML全量类图.md 则可帮助你把 Semaphore 放进整个 JUC 体系中建立全局认知。【免费下载链接】source-code-hunter 从源码层面剖析挖掘互联网行业主流技术的底层实现原理为广大开发者 “提升技术深度” 提供便利。目前开放 Spring 全家桶Mybatis、Netty、Dubbo 框架及 Redis、Tomcat 中间件等项目地址: https://gitcode.com/GitHub_Trending/so/source-code-hunter创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表