Java Semaphore 信号量源码解析:基于 AQS 的并发许可控制与限流实现
2026/9/12 20:13:44 网站建设 项目流程

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

【免费下载链接】source-code-hunter😱 从源码层面,剖析挖掘互联网行业主流技术的底层实现原理,为广大开发者 “提升技术深度” 提供便利。目前开放 Spring 全家桶,Mybatis、Netty、Dubbo 框架,及 Redis、Tomcat 中间件等项目地址: https://gitcode.com/GitHub_Trending/so/source-code-hunter

Semaphore(信号量)是 JUC 中用于控制一定时间内并发执行线程数的同步工具,其内部完全基于 AbstractQueuedSynchronizer(AQS)实现,在 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 提供的acquireSharedInterruptiblyreleaseShareddoAcquireSharedInterruptiblydoReleaseSharedsetHeadAndPropagate等模板与工具方法则在 详解AbstractQueuedSynchronizer.md 中有完整剖析(获取共享锁 / 释放共享锁的实现见该文档“获取共享锁的实现”与“释放共享锁的实现”两节)。本文后续会反复回到这些方法,说明它们如何与 Semaphore 的许可逻辑咬合。

核心内部类 Sync:许可状态的管理者

先看 Semaphore 的骨架——所有许可逻辑都收敛在Sync中(源码对应仓库文档 Semaphore.md):

abstract 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; } } }

这里有几个值得深入的关键点:

  1. 许可数即 AQS 的同步状态 state。构造 Semaphore 时调用setState(permits),把初始许可数写入 AQS 的state;之后getPermits()/getState()返回的永远是当前剩余许可数。AQS 对 state 提供了 volatile 可见性保证与 CAS 更新原语,这是整个并发安全性的基石。

  2. 获取许可 = 自旋 + CAS 扣减nonfairTryAcquireShared是一个无锁算法:循环读取available,计算出扣减后的remaining,只有remaining >= 0且 CAS 成功才会返回;否则继续自旋重试。返回值语义与 AQS 共享锁的约定完全一致(详见 AQS 文档):返回非负数表示获取成功(0 表示成功但后继争用线程不会成功,正数表示成功且后继也可能成功),返回负数表示获取失败,由上层据此决定是否入队阻塞。

  3. 释放许可 = 自旋 + CAS 累加,且不设上限tryReleaseShared直接把current + releases写回 state。注意它只检查了整数溢出next < current时抛出Error("Maximum permit count exceeded")),却没有限制累加后的值不能超过初始许可数——这就是原文档强调的“可以通过 release 动态增加许可数量”的源码出处。这在业务上是一种灵活性:例如允许运维在运行时临时扩容“许可池”,让更多线程并发执行。

  4. 许可的运维操作reducePermits(int reductions)自旋 CAS 扣减许可(检查下溢,即next > current时抛Error("Permit count underflow"));drainPermits()一次性把许可清零并返回清零前的数量,可用于“暂停放行”等场景。这两个方法在Sync中以final方法存在,对外通过Semaphore的公有 API(如drainPermits()availablePermits())暴露,属子类可控的辅助能力。

公平与非公平:两种许可获取策略

Semaphore通过构造器参数fairFairSyncNonfairSync之间选择策略,二者都继承自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); } }

调用链拆解

  1. 入口校验acquire(int permits)对负数参数抛IllegalArgumentException,随后进入 AQS 的acquireSharedInterruptibly

  2. 中断响应acquireSharedInterruptibly首先检查Thread.interrupted(),若线程已被中断则立即抛InterruptedException,保证acquire的可中断语义。

  3. 快速路径:调用tryAcquireShared(arg)(非公平模式即Sync.nonfairTryAcquireShared)。返回值>= 0说明剩余许可充足且 CAS 成功,线程直接持有许可继续执行,无需排队。

  4. 慢速路径(入队阻塞):返回值为负则进入doAcquireSharedInterruptibly——这是 AQS 共享模式的标准排队逻辑,与 详解AbstractQueuedSynchronizer.md 中doAcquireShared的实现同源:

    • addWaiter(Node.SHARED)把当前线程包装成共享模式节点插入同步队列尾部;
    • 循环中只允许前驱为 head 的节点尝试再次获取许可(保证队列纪律,防止队列内部插队);
    • 获取成功则调用setHeadAndPropagate(node, r)——这是共享锁与独占锁的关键区别:它不仅把自己设为新 head,还会根据返回值 r 与节点状态决定是否继续唤醒后继的共享节点,实现“许可有余量则依次放行多个线程”的传播效应;
    • 获取失败则通过shouldParkAfterFailedAcquire把前驱节点状态置为 SIGNAL,再由parkAndCheckInterrupt阻塞当前线程,等待被unpark唤醒或响应中断;
    • 若因中断退出循环,finally中的cancelAcquire(node)会将该节点标记为 CANCELLED 并从队列中摘除。

许可放行与唤醒传播的底层细节

setHeadAndPropagatedoReleaseShared是共享锁“一放多醒”的核心,其源码细节在 详解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; } }

释放链路的要点:

  1. release(int permits)同样对负数参数做防御性校验;
  2. tryReleaseShared用自旋 CAS 把许可累加回去,成功返回 true 后进入doReleaseShared
  3. doReleaseShared唤醒同步队列中第一个等待线程(head 的后继),其余等待线程由被唤醒者沿队列向后传递唤醒,注释“唤醒等待的第一个线程(其他线程将由第一个线程向后传递唤醒)”正是对这一传播机制的精确描述;
  4. 动态扩容tryReleaseShared只做溢出检查、不限制累加上限,因此release()可以释放比初始许可更多的许可。例如初始许可为 3,两个线程各 release 一次后,许可数会变成 5,后续即可容纳最多 5 个并发线程。这是原文档反复强调的核心行为,也是 Semaphore 区别于“固定容量资源池”的一个重要特性——若业务需要严格固定上限,必须自行约束 release 的调用次数,或配合reducePermits收紧。

其他常用 API 与注意事项

结合Sync提供的底层能力,Semaphore还向外暴露了一系列实用 API(从源码结构可以确认其存在与语义):

API作用备注
availablePermits()返回当前剩余许可数委托sync.getPermits(),仅作监控参考,非线程安全快照
drainPermits()一次性清空所有许可并返回清空前的数量可用于“熔断放行”,拒绝一切新任务直到许可被恢复
reducePermits(int)减少指定数量许可底层为Sync.reducePermits,下溢时抛Error
isFair()判断是否为公平模式返回sync instanceof FairSync
hasQueuedThreads()/getQueueLength()是否有排队线程 / 排队线程数委托 AQS 队列查询能力,用于监控

使用注意事项:

  • 许可数必须非负:构造器Semaphore(int permits)Semaphore(int permits, boolean fair)在 permits 为负时会抛IllegalArgumentException
  • acquire 与 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 的同步状态 state,setState(permits)初始化Semaphore.md
并发安全全程自旋 + CAS,无锁实现许可的扣减与累加Semaphore.md
公平性FairSync通过hasQueuedPredecessors()保证先来先得,默认NonfairSync直接抢Semaphore.md
排队阻塞失败线程以 SHARED 节点入 AQS 同步队列,park阻塞、unpark唤醒详解AbstractQueuedSynchronizer.md
唤醒传播setHeadAndPropagate+doReleaseShared+ PROPAGATE 状态实现“一放多醒”详解AbstractQueuedSynchronizer.md
动态扩容tryReleaseShared仅做溢出检查,release 可无上限累加许可Semaphore.md

Semaphore 是理解 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询