ARTICLE DETAIL

资讯详情

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

Semaphore源码剖析:从AQS到公平非公平,彻底拆解并发控制

Semaphore源码剖析:从AQS到公平非公平,彻底拆解并发控制 第一次用 Semaphore 的时候我把它当成一个带计数器的锁acquire 就是拿锁release 就是放锁。直到线上一个服务出现莫名的线程阻塞排查日志时发现一堆线程卡在 WAITING(parking)我才真正意识到Semaphore 的价值根本不在于“锁”而在于它背后那套基于 AQSAbstractQueuedSynchronizer的并发控制框架。这篇文章我就从源码层面把 Semaphore 彻底拆开看看 state 字段怎样从“许可总数”变成“剩余许可”、公平和非公平到底差在哪一行代码、共享模式下唤醒为什么需要 PROPAGATE 这个状态以及真正生产环境里哪些用法会踩坑。我读的是 JDK 8 的 AQS 实现JDK 17 里 Semaphore 和 AQS 的核心逻辑基本一致只是部分方法做了微调。所以这篇文章的讲解对你本地的源码同样适用。想深入理解 JUC 的并发控制、或者被 Semaphore 的阻塞问题折磨过这篇源码剖析能帮你把这些“黑盒”变成“白盒”。1. 先回答Semaphore在这套并发工具箱里到底解决什么问题1.1 从一次“假死”事故说起我曾经见过一个数据导出服务需求是“最多同时跑8个导出任务”。最初实现用的是 synchronized结果一个任务在等待远程接口响应的过程中其他任务全部被挡在锁外面整个导出队的列吞吐惨不忍睹。后来把 synchronized 换成new Semaphore(8)问题立刻消失——因为Semaphore控制的不是“同一时刻只有一个线程”而是“同一时刻最多只有 N 个线程”。信号量Semaphore这个概念最早是 Dijkstra 在 1965 年提出的Java 的java.util.concurrent.Semaphore把它移植到了 JVM 世界。核心模型简单得不能再简单一个计数器表示当前可用的许可permit数量线程拿许可前必须先acquire()用完必须release()还回去。不过模型简单实现却一点也不简单——并发场景下要处理 CAS 竞争、线程排队、挂起、唤醒以及各种竞态条件。这些都压在 AQS 的肩膀上。1.2 三行代码说清楚并发度先看最基本的用法Semaphore semaphore new Semaphore(3); // 同时最多 3 个线程 for (int i 0; i 10; i) { new Thread(() - { try { semaphore.acquire(); // 获取一个许可不够就阻塞等待 System.out.println(Thread.currentThread().getName() 开始执行); Thread.sleep(200); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } finally { semaphore.release(); // 释放许可让给后面的线程 System.out.println(Thread.currentThread().getName() 释放许可); } }, worker- i).start(); }上面这段代码只有 3 个线程能同时进入临界区其余 7 个线程会阻塞在acquire()上。注意 3 是并发度上限不是总执行次数——第 9 个线程只要有人释放许可它就能拿到。这些语义都是源码里明确写死的理解这一点后面看源码才有方向。1.3 类结构四个类的分工其实很清晰Semaphore内部有一个抽象类Sync它继承AbstractQueuedSynchronizer。Sync又分裂出两个子类FairSync公平信号量和NonfairSync非公平信号量。构造器里的fair参数决定实例化哪一个。类职责Semaphore对外暴露acquire、release、tryAcquire等 API内部全部委托给syncSync继承 AQS实现共享模式的获取/释放逻辑管理剩余许可数FairSync公平模式队列中有排队的线程时新来的不去抢许可NonfairSync非公平模式只要许可够新来的线程可以直接 CAS 抢AbstractQueuedSynchronizer提供 state、CLH 等待队列、park/unpark 等基础设施换句话说Semaphore只负责定义“信号量”的业务语义真正排版、阻塞、唤醒这些脏活累活全部下沉到 AQS。Semaphore的源码不到 400 行核心逻辑更少可它的行为却完全取决于 AQS 那几百行模板方法。下一章我们就看 AQS 到底给Semaphore递了什么底牌。2. AQS给Semaphore的底牌state、模板方法和CLH队列2.1 state许可证数量以一种“反直觉”的方式存放AQS 里有一个volatile int state所有同步器都围绕它工作。ReentrantLock里state1表示锁被持有state0表示锁空闲而Semaphore的state表示的是“当前还剩多少个许可”。比如new Semaphore(3)时AQS 构造器会执行setState(3)。之后每次acquire()让 state 减 1每次release()让 state 加 1。用表格看更直观操作state 变化含义new Semaphore(3)3还有 3 个许可线程A acquire()2还剩 2 个许可线程B acquire()1还剩 1 个许可线程C acquire()0许可耗尽下一个 acquire 要排队线程A release()1释放一个许可所有对state的修改都必须通过compareAndSetState(expect, update)完成这是一套基于 CAS 的乐观锁逻辑保证并发环境下修改安全。Semaphore里“获取许可”这个动作本质上就是对 state 做一次受保护的自减。2.2 模板方法为什么子类只写 try 方法就够了AQS 的设计哲学是模板方法模式。它把同步器拆成两部分固定流程和可变逻辑。固定流程包括“尝试获取失败后就入队、挂起、被唤醒后重新尝试”这一整套可变逻辑只有两个小方法——tryAcquireShared和tryReleaseShared。AQS 的模板方法长这样public final void acquireSharedInterruptibly(int arg) throws InterruptedException { if (Thread.interrupted()) throw new InterruptedException(); if (tryAcquireShared(arg) 0) doAcquireSharedInterruptibly(arg); }tryAcquireShared的返回值不是布尔值而是一个 int大于等于 0 表示获取成功且返回剩余许可数小于 0 表示获取失败需要进入等待队列。Semaphore要做的事就是告诉 AQS“我的获取逻辑是什么”至于失败后怎么排队、怎么挂起、怎么唤醒都是 AQS 的模板方法在管。用点外卖来类比你只需要告诉商家要什么菜try方法用户下单、商家接单、骑手配送这些流程由平台统一处理模板方法。如果不采用模板方法设计每个同步器都得自己实现一套排队和唤醒逻辑代码冗余且容易出错JUC 的架构也不会这么优雅。2.3 CLH等待队列排队不是“往里一坐”那么简单AQS 内部维护着一个双向链表队列每个等待线程被包装成一个Node节点。节点上主要存这些信息thread当前线程prev、next前后节点引用waitStatus节点状态后面详谈nextWaiter标记当前节点是共享模式还是独占模式Semaphore里所有节点都是SHARED队列的规矩是head节点代表正在占用许可的线程在Semaphore里也可以理解成一个空的队头占位符真正的等待者全部排在head后面。新来的线程在队尾插入成为新的tail。只有prev是head的节点才有资格去尝试获取——也就是说队列内部还是严格的 FIFO非公平模式所谓的“不公平”只发生在入队之前。这个结构很像银行柜台排队柜台 A 正在办业务BC 在后面排队A 办完叫 BB 办完叫 C。需要注意的是“叫号”不是当场发生的而是通过LockSupport.unpark把 B 唤醒B 醒来后再走一遍tryAcquireShared去抢许可。2.4 waitStatus一个 int 藏着四种含义每个Node都有一个waitStatus字段它是后续唤醒机制的核心取值如下值常量含义1CANCELLED节点已取消比如线程被中断-1SIGNAL当前节点的后继需要被唤醒-2CONDITION节点在条件队列中等待Semaphore不用-3PROPAGATE共享模式下唤醒需要继续向后传播0默认节点处于正常初始状态后面读doAcquireSharedInterruptibly和doReleaseShared源码时你会反复看到这些状态怎么被修改。特别是PROPAGATE它是共享模式里最容易被忽略、但又最关键的一个状态。3. 公平和非公平的源码岔路tryAcquireShared的两种写法3.1 非公平锁抢到就是赚到打开NonfairSync里面只有这么一段代码protected int tryAcquireShared(int acquires) { for (;;) { int available getState(); // 当前剩余许可 int remaining available - acquires; // 拿掉 acquires 个后还剩多少 if (remaining 0 || compareAndSetState(available, remaining)) return remaining; } }逐行读先读当前state也就是剩余许可数。减掉要获取的许可数得到remaining。如果remaining 0说明许可不够了直接返回负数不修改任何状态。如果许可够就 CAS 尝试把state从available改成remaining。CAS 失败说明有其他线程抢先了继续循环重试。这就是“非公平”的本义新来的线程不需要看队列里有没有人只要许可数还有富余它就可以通过 CAS 抢走一个许可。哪怕队列里已经有一个线程等了几个小时新来的线程照样可以插队。只要remaining 0资格瞬间生效。这种设计追求的是吞吐牺牲的是公平性——极端竞争下队列里某些线程可能长时间拿不到许可也就是所谓的饥饿。3.2 公平锁先看队列里有没有人FairSync只比NonfairSync多了一行判断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()是 AQS 提供的方法它的语义是“是否有线程排在我前面”。注意一个容易忽略的细节如果队列里第一个排队节点正好是当前线程这个方法返回 false允许当前线程再尝试一次。这是为了避免“线程刚被唤醒抢锁资格又被剥夺”的尴尬情况。两种实现放在一起对比非常有意思非公平和公平只差一个hasQueuedPredecessors()。但正是这一个方法把“插队”行为彻底禁止了。所有新线程看到队列有人直接返回 -1 去队尾排队绝不尝试 CAS。3.3 一个容易忽略的返回值为什么是 int 而不是 boolean这是很多初读源码的人会懵的地方tryAcquireShared明明叫“尝试获取”为什么不返回 boolean原因在于共享模式需要传递“剩余量”。Semaphore一次性可能获取 1 个、2 个甚至 N 个许可返回 int 既能表达成功与否0 成功0 失败又能告诉 AQS“现在还剩多少个许可可以继续分配”。这个剩余量到后面setHeadAndPropagate会被用到如果剩得多唤醒行为就要像接力赛一样传下去让多个等待线程都在本次释放中获益。3.4 公平与非公平的性能差本质差在上下文切换为什么非公平模式往往吞吐更高因为它避免了一次“入队→挂起→唤醒”的完整链路。想象一个短任务的场景任务执行只需 2ms许可刚好够用新线程直接 CAS 成功连队列都不用进。而公平模式下每个新线程都要先检查队列一旦发现队列非空立刻入队线程被 park 之后又要等 unpark这中间涉及操作系统线程调度的开销成本远高于一次 CAS。所以我的选型建议是默认new Semaphore(n)的非公平模式完全没问题适合绝大多数短任务、高并发的场景一旦任务耗时长、对响应时间敏感或者你希望每个线程都有确定的执行机会就用new Semaphore(n, true)公平模式。但要注意公平模式在高竞争下吞吐下降可能非常明显建议压测对比后再定。4. acquire全链路追踪从入口到LockSupport.park4.1 acquire 到底掉进了哪个方法Semaphore.acquire()内部就是一行委托public void acquire() throws InterruptedException { sync.acquireSharedInterruptibly(1); }也就是说Semaphore走的是 AQS 的共享模式、响应中断版本。模板方法里如果Thread.interrupted()为 true会直接抛InterruptedException这个动作意味着“中断优先于获取”。之后如果tryAcquireShared(1) 0才进入doAcquireSharedInterruptibly(arg)开始排队等待。如果不想响应中断可以用acquireUninterruptibly()它走的是另一个不检查中断的模板方法。不过生产环境我强烈建议保留中断处理方便线程池shutdownNow时能及时退出阻塞。4.2 addWaiter尾插不是简单的 addLast进入排队流程后AQS 会通过addWaiter把当前线程包装成节点插到队列尾部private Node addWaiter(Node mode) { Node node new Node(Thread.currentThread(), mode); Node pred tail; if (pred ! null) { node.prev pred; if (compareAndSetTail(pred, node)) { // CAS 替换尾节点 pred.next node; return node; } } enq(node); // CAS 失败则走自旋入队 return node; }这里必须用 CAS因为 tail 是从队尾插入的当前线程的 prev 先指向旧 tail然后 CAS 把 tail 更新为自己。如果 CAS 失败说明有并发线程抢先入队了那就走enq里的for(;;)循环反复尝试直到成功。所有并发容器里的“尾插”本质都是这种 CAS 自旋的套路。4.3 只有 prev 是 head 的节点才有资格 trydoAcquireSharedInterruptibly的核心循环是这样private void doAcquireSharedInterruptibly(int arg) throws InterruptedException { final Node node addWaiter(Node.SHARED); try { for (;;) { final Node p node.predecessor(); if (p head) { // 只有前驱是 head 才能抢 int r tryAcquireShared(arg); if (r 0) { setHeadAndPropagate(node, r); p.next null; // 帮助 GC 清理 return; } } if (shouldParkAfterFailedAcquire(p, node) parkAndCheckInterrupt()) throw new InterruptedException(); } } catch (Throwable t) { cancelAcquire(node); throw t; } }为什么必须是p head因为 head 是“正在占用许可”的节点只有它的后继才处于队首才有资格去尝试获取。这个设计保证了队列内部仍然是 FIFO——你再着急也必须等到前一个人下队。获取成功后会调用setHeadAndPropagate把当前节点设为新的 head并且尝试唤醒后面的线程。这个方法的细节我会在 release 那章一起讲因为它是共享模式“唤醒接力”的核心。4.4 shouldParkAfterFailedAcquire两次自旋后才真正挂起获取失败后能不能立刻park不行。AQS 先要做一件准备动作确保前驱节点处于SIGNAL状态这样将来前驱被释放时才会来唤醒我。源码如下private static boolean shouldParkAfterFailedAcquire(Node pred, Node node) { int ws pred.waitStatus; if (ws Node.SIGNAL) return true; // 前驱保证会唤醒我可以放心 park if (ws 0) { // 前驱已取消 do { node.prev pred pred.prev; } while (pred.waitStatus 0); pred.next node; } else { compareAndSetWaitStatus(pred, ws, Node.SIGNAL); // 把前驱状态改成 SIGNAL } return false; // 还没准备好下一轮循环再决定 }三种情况前驱是SIGNAL说明前驱承诺“等我执行完会来叫你”当前线程可以安心睡觉。前驱是CANCELLED说明前驱已经取消等待当前线程必须跳过它重新建立链表连接顺便把取消节点“剪”出去。前驱是 0 或其他状态就 CAS 把它改成SIGNAL然后返回 false触发外层循环再自旋一轮。第二轮回来时前驱已经是SIGNAL这才真正park。这就是为什么你经常看到“连续两次循环才挂起”的说法第一次循环负责设置 SIGNAL第二次循环才能 park。4.5 park 之后的两种醒来方式parkAndCheckInterrupt的实现非常短private final boolean parkAndCheckInterrupt() { LockSupport.park(this); return Thread.interrupted(); }线程执行到LockSupport.park(this)后会进入WAITING(parking)状态。醒来有两种可能被前驱节点通过unpark唤醒——这是正常情况继续循环尝试获取。因为interrupt()被唤醒——Thread.interrupted()返回 true注意它会清除中断标志然后抛InterruptedException进入cancelAcquire取消当前节点的排队资格。这里有一个很微妙的点park不会自动响应中断它只是让线程醒过来是否抛异常由上层代码判断。所以你在catch (InterruptedException e)里往往要补一句Thread.currentThread().interrupt()——因为中断标志已经被Thread.interrupted()清掉了不补回去上层就可能感知不到中断发生过。5. release的另一半真相doReleaseShared与PROPAGATE状态5.1 tryReleaseShared一个简单的 CAS 加法释放许可的入口是Semaphore.release()最终调用sync.releaseShared(1)。AQS 的releaseShared模板方法会先调用子类的tryReleaseSharedprotected final boolean tryReleaseShared(int releases) { for (;;) { int current getState(); int next current releases; if (next current) // int 溢出保护 throw new Error(Maximum permit count exceeded); if (compareAndSetState(current, next)) return true; } }这个方法和tryAcquireShared正好对称循环读当前 state加上要释放的许可数CAS 更新。注意next current的判断——如果许可数加起来超过Integer.MAX_VALUE会抛出一个很罕见的Error。源代码作者显然考虑到极端情况了虽然实际很难触发。release()还有一个特点它不像acquire那样可能返回负数去排队释放操作只会成功不存在失败分支。5.2 doReleaseShared唤醒的不是“下一个”而是“值得唤醒的后继”真正干唤醒活的是 AQS 的doReleaseSharedprivate 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; unparkSuccessor(h); // 唤醒后继 } else if (ws 0 !compareAndSetWaitStatus(h, 0, Node.PROPAGATE)) continue; } if (h head) // 如果 head 没变结束循环 break; } }逻辑分三条路如果 head 的waitStatus是SIGNAL说明后面有人等着被唤醒。先把状态从SIGNAL改成 0然后调用unparkSuccessor(h)唤醒真正的后继。为什么要先改成 0因为要防止自己还没唤醒别人又被其他线程重复唤醒一次。如果 head 的状态是 0说明“没有明确需要唤醒的人”但为了不让唤醒信号随便丢失会把状态改成PROPAGATE这就是共享模式接力赛的“火种”。如果并发中 head 被换掉了循环会发现h ! head继续处理新的 head。unparkSuccessor的细节也值得说它找后继不是简单从node.next拿而是从 tail 往前遍历找到最靠近 head 的未取消节点。因为队列里可能有CANCELLED节点破坏next链从后往前遍历才能可靠找到真正需要唤醒的线程。5.3 PROPAGATE专治共享模式下唤醒信号丢失这应该是整篇文章里最硬核的部分。为什么需要PROPAGATE状态因为共享模式存在一个著名竞态多个线程同时release如果唤醒信号在队头切换的瞬间丢了队列里明明有等待者却没有人去叫醒它们。想象这种场景线程 A 调用release看到 head 是 SIGNALCAS 把 head 改成 0正准备unparkSuccessor。就在这个瞬间队列里的线程 B 被前一次唤醒获取许可成功把自己设置成了新 head。旧 head 的 waitStatus 还是 0。线程 C 也调用release看到新 head 的 waitStatus 是 0没有去唤醒任何线程因为“没有人需要被唤醒”但实际上队列里还有线程 D 在等待它的前驱也就是旧的 head已经不是 headD 没被信号覆盖到。如果没有PROPAGATED 就可能一直睡下去。有了它以后doReleaseShared在面对状态为 0 的 head 时会标记成PROPAGATE而获取侧的setHeadAndPropagate看到这个状态就会继续调用doReleaseShared把唤醒行为向后传播。对应地setHeadAndPropagate源码如下private void setHeadAndPropagate(Node node, int propagate) { Node h head; setHead(node); if (propagate 0 || h null || h.waitStatus 0 || (h head) null || h.waitStatus 0) { Node s node.next; if (s null || s.isShared()) doReleaseShared(); } }四个触发继续唤醒的条件剩余许可数大于 0、旧 head 状态小于 0说明有过传播标记、新 head 状态小于 0、后继节点是共享模式。只要满足任意一个就继续doReleaseShared把唤醒接力下去。这就是共享模式下多个等待线程可以“一波波醒来”的底层原理。5.4 获取侧与释放侧如何握手把整个闭环串起来acquire时tryAcquireShared返回负数线程入队前驱状态被设为SIGNAL线程park。release时tryReleaseShared成功加回 state进入doReleaseShared。doReleaseShared发现 head 是SIGNAL改为 0unparkSuccessor唤醒队首等待线程。被唤醒线程从park返回重新走tryAcquireShared拿到许可后setHeadAndPropagate。如果还有剩余许可或传播条件满足继续唤醒下一个队列逐步清空。到这里Semaphore的获取和释放源码就完全串起来了。这些机制看起来复杂但理解了“CLH 队列 waitStatus 共享传播”这三板斧JUC 里其他同步器比如CountDownLatch、ReentrantReadWriteLock的源码也就通了一半。6. 读完源码才会明白的实战细节开关、精度与限流边界6.1 acquire/release 必须成对跨线程释放才是信号量的灵魂读完源码你会发现Semaphore根本不关心“谁获得了许可”也不检查“谁释放许可”。state只是一堆数字线程 A 获取的许可完全可以由线程 B 释放。这是它和synchronized的本质区别——跨线程释放是合法的。但这也带来了风险任何线程获取到许可后如果因为异常、忘记写 finally、或者逻辑分支漏了release()许可就真的“消失”了一个。长期积累剩余许可会降到 0所有线程永久阻塞服务进入假死状态。所以我在使用Semaphore时有一条铁律acquire之后必须try-finally包裹且在finally中release。看起来是常识线上事故却往往就出在这一行遗漏上。6.2 tryAcquire 与超时生产代码的首选裸写acquire()的最大问题是如果许可一直拿不到线程就会一直等着。在生产环境我更推荐带超时的版本Semaphore semaphore new Semaphore(8); if (semaphore.tryAcquire(2, TimeUnit.SECONDS)) { try { // 执行业务 } finally { semaphore.release(); } } else { // 拿不到许可走降级、重试或告警 }源码路径上tryAcquire(long timeout, TimeUnit unit)会进入doAcquireSharedNanos和doAcquireSharedInterruptibly几乎一样只是多了超时时间LockSupport.parkNanos(this, nanosTimeout)到点醒来后如果还没获得许可就返回 false。它不会无限阻塞也不会让你连一个“拿不到许可”的处理机会都没有。6.3 Semaphore 限流为什么限不住 QPSSemaphore是并发度闸门不是速率闸门。举一个具体例子假设接口平均耗时 2ms用Semaphore(50)限流理论最大吞吐是 50 / 0.002 25000 QPS。虽然瞬时并发被限制到了 50但每秒能处理的请求总量依然非常可观根本限不住“每秒多少次”的维度。从源码也能印证这一点release之后 state 立刻加回完全不存在时间维度。如果你要限制的是 QPS、或者平滑流量突发应该用令牌桶算法比如 Guava RateLimiter或者 Redis 分布式限流而不是Semaphore。我见过不少团队用Semaphore做限流最后发现“限了个寂寞”就是因为没想清楚这两个概念的边界。6.4 drainPermits 与 reducePermits两个少有人知的内部接口drainPermits()的作用是直接清空剩余许可并返回清空的许可数。典型场景是应用准备优雅停机先调用drainPermits()把现有许可全部收走新请求拿不到许可自然走降级通路已经进入临界区的任务则继续执行完。reducePermits(int reduction)用于减少总许可数它是protected方法需要自定义子类暴露class MySemaphore extends Semaphore { MySemaphore(int permits) { super(permits); } void shrink(int reduction) { super.reducePermits(reduction); // 动态把并发度从 10 降到 5 } }注意这个方法仅缩减许可数量不触发任何阻塞或唤醒。如果你要动态调整并发度又不想让排队线程全部取消等待这是一个很趁手的“并发度开关”。6.5 一次 jstack 定位的真实复盘最后分享一次让我印象深刻的排查。线上某个服务在流量高峰时突然大面积超时线程池打满jstack一看大量线程堆栈停留在AbstractQueuedSynchronizer.park和Semaphore.acquire上。availablePermits()打印出来是 0但整条业务链路并没有显式的死锁。顺着代码一路找下去才发现某个分支在 catch 异常后直接 return 而忘了release()把许可“吞”掉了。这个事故给我的教训是Semaphore的坑往往不在并发代码本身而在业务异常路径。你如果决定在一个高并发项目里用Semaphore第一件事就是在类注释里写清楚“谁获取、谁释放、并发度是多少、是否允许动态收缩”同时加上availablePermits()的监控。源码都能看明白但线上事故往往就出在“默认它不会错”和“没看明白”之间。希望这篇源码拆解能让你少踩一个Semaphore的坑。
返回列表