Semaphore(信号量)

📅 发布时间:2026/8/22 12:15:09
Semaphore(信号量) Semaphore 介绍synchronized和ReentrantLock都是一次只允许一个线程访问某个资源而Semaphore(信号量)可以用来控制同时访问特定资源的线程数量。Semaphore的使用简单我们这里假设有N(N5)个线程来获取Semaphore中的共享资源下面的代码表示同一时刻 N 个线程中只有 5 个线程能获取到共享资源其他线程都会阻塞只有获取到共享资源的线程才能执行。等到有线程释放了共享资源其他阻塞的线程才能获取到。// 初始共享资源数量 final Semaphore semaphore new Semaphore(5); // 获取1个许可 semaphore.acquire(); // 释放1个许可 semaphore.release();当初始的资源个数为 1 的时候Semaphore 退化为排他锁。Semaphore 有两种模式公平模式 调用 acquire() 方法的顺序就是获取许可证的顺序遵循 FIFO非公平模式 抢占式的。Semaphore 对应的两个构造方法如下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 原理Semaphore是共享锁的一种实现它默认构造 AQS 的state值为permits你可以将permits的值理解为许可证的数量只有拿到许可证的线程才能执行。以无参acquire方法为例调用semaphore.acquire()线程尝试获取许可证如果state 0的话则表示可以获取成功如果state 0的话则表示许可证数量不足获取失败。如果可以获取成功的话(state 0)会尝试使用 CAS 操作去修改state的值statestate-1。如果获取失败则会创建一个 Node 节点加入等待队列挂起当前线程。// 获取1个许可证 public void acquire() throws InterruptedException { sync.acquireSharedInterruptibly(1); } // 获取一个或者多个许可证 public void acquire(int permits) throws InterruptedException { if (permits 0) throw new IllegalArgumentException(); sync.acquireSharedInterruptibly(permits); }acquireSharedInterruptibly方法是AbstractQueuedSynchronizer中的默认实现。// 共享模式下获取许可证获取成功则返回失败则加入等待队列挂起线程 public final void acquireSharedInterruptibly(int arg) throws InterruptedException { if (Thread.interrupted()) throw new InterruptedException(); // 尝试获取许可证arg为获取许可证个数当获取失败时,则创建一个节点加入等待队列挂起当前线程。 if (tryAcquireShared(arg) 0) doAcquireSharedInterruptibly(arg); }这里再以非公平模式NonfairSync的为例看看tryAcquireShared方法的实现。// 共享模式下尝试获取资源(在Semaphore中的资源即许可证): protected int tryAcquireShared(int acquires) { return nonfairTryAcquireShared(acquires); } // 非公平的共享模式获取许可证 final int nonfairTryAcquireShared(int acquires) { for (;;) { // 当前可用许可证数量 int available getState(); /* * 尝试获取许可证当前可用许可证数量小于等于0时返回负值表示获取失败 * 当前可用许可证大于0时才可能获取成功CAS失败了会循环重新获取最新的值尝试获取 */ int remaining available - acquires; if (remaining 0 || compareAndSetState(available, remaining)) return remaining; } }以无参release方法为例调用semaphore.release();线程尝试释放许可证并使用 CAS 操作去修改state的值statestate1。释放许可证成功之后同时会唤醒等待队列中的一个线程。被唤醒的线程会重新尝试去修改state的值statestate-1如果state 0则获取令牌成功否则重新进入等待队列挂起线程。// 释放一个许可证 public void release() { sync.releaseShared(1); } // 释放一个或者多个许可证 public void release(int permits) { if (permits 0) throw new IllegalArgumentException(); sync.releaseShared(permits); }releaseShared方法是AbstractQueuedSynchronizer中的默认实现。// 释放共享锁 // 如果 tryReleaseShared 返回 true就唤醒等待队列中的一个或多个线程。 public final boolean releaseShared(int arg) { //释放共享锁 if (tryReleaseShared(arg)) { //释放当前节点的后置等待节点 doReleaseShared(); return true; } return false; }tryReleaseShared方法是Semaphore的内部类Sync重写的一个方法AbstractQueuedSynchronizer中的默认实现仅仅抛出UnsupportedOperationException异常。// 内部类 Sync 中重写的一个方法 // 尝试释放资源 protected final boolean tryReleaseShared(int releases) { for (;;) { int current getState(); // 可用许可证1 int next current releases; if (next current) // overflow throw new Error(Maximum permit count exceeded); // CAS修改state的值 if (compareAndSetState(current, next)) return true; } }可以看到上面提到的几个方法底层基本都是通过同步器sync实现的。Sync是CountDownLatch的内部类 , 继承了AbstractQueuedSynchronizer重写了其中的某些方法。并且Sync 对应的还有两个子类NonfairSync对应非公平模式 和FairSync对应公平模式。private static final class Sync extends AbstractQueuedSynchronizer { // ... } static final class NonfairSync extends Sync { // ... } static final class FairSync extends Sync { // ... }Semaphore 实战public class SemaphoreExample { // 请求的数量 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/120 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经常用于限制获取某种资源的线程数量。当然一次也可以一次拿取和释放多个许可不过一般没有必要这样做semaphore.acquire(5);// 获取5个许可所以可运行线程数量为20/54 test(threadnum); semaphore.release(5);// 释放5个许可除了acquire()方法之外另一个比较常用的与之对应的方法是tryAcquire()方法该方法如果获取不到许可就立即返回 false。Semaphore 补充Semaphore与CountDownLatch一样也是共享锁的一种实现。它默认构造 AQS 的state为permits。当执行任务的线程数量超出permits那么多余的线程将会被放入等待队列Park,并自旋判断state是否大于 0。只有当state大于 0 的时候阻塞的线程才能继续执行,此时先前执行任务的线程继续执行release()方法release()方法使得 state 的变量会加 1那么自旋的线程便会判断成功。如此每次只有最多不超过permits数量的线程能自旋成功便限制了执行任务线程的数量。