Semaphore

Semaphore 源码分析

Semaphore 解决的问题:限制同时进入某段代码/某个资源的线程数量(比如”这个接口最多 10 个并发”)。
核心思路是维护一个许可计数,acquire 扣一个(不够就排队阻塞),release 还一个——许可不是对象,只是 AQS 里的一个 int:state = 剩余许可数。
整个类没有一行阻塞代码,排队、挂起、唤醒全部委托给 AQS 的共享模式。

// 基于 JDK 17, java.util.concurrent.Semaphore —— 同步器部分
public class Semaphore implements java.io.Serializable {
    // ...
    private final Sync sync;

    abstract static class Sync extends AbstractQueuedSynchronizer {
        // ...

        Sync(int permits) {
            setState(permits);      // 许可数直接放进 AQS 的 state;允许传负数,表示"先欠着",release 够了才能 acquire
        }

        final int getPermits() {
            return getState();
        }

        final int nonfairTryAcquireShared(int acquires) {
            for (;;) {
                int available = getState();
                int remaining = available - acquires;
                if (remaining < 0 ||
                    compareAndSetState(available, remaining))
                    return remaining;   // 这个 || 是关键:不够时不做 CAS 直接返回负数(失败零副作用,state 不会被 acquire 减成负);够但 CAS 输了就重试
            }
        }

        protected final boolean tryReleaseShared(int releases) {
            for (;;) {
                int current = getState();
                int next = current + releases;
                if (next < current) // overflow   // int 溢出防护:release 没有上限校验,能加出比初始更多的许可,只拦回绕
                    throw new Error("Maximum permit count exceeded");
                if (compareAndSetState(current, next))
                    return true;    // 永远返回 true → AQS 每次 release 都会去唤醒队头,不像锁那样"减到 0 才唤醒"
            }
        }

        final void reducePermits(int reductions) {
            for (;;) {
                int current = getState();
                int next = current - reductions;
                if (next > current) // underflow
                    throw new Error("Permit count underflow");
                if (compareAndSetState(current, next))
                    return;         // 和 acquire 的本质区别:不排队不阻塞,允许把 state 扣成负数。用于"资源永久坏了一个"的场景
            }
        }

        final int drainPermits() {
            for (;;) {
                int current = getState();
                if (current == 0 || compareAndSetState(current, 0))
                    return current; // 一口气清零并返回抢到的数量;state 为负时相当于把欠账归零
            }
        }
    }

    static final class NonfairSync extends Sync {
        // ...

        NonfairSync(int permits) {
            super(permits);
        }

        protected int tryAcquireShared(int acquires) {
            return nonfairTryAcquireShared(acquires);   // 非公平(默认):上来就 CAS 抢,可能插队,吞吐高
        }
    }

    static final class FairSync extends Sync {
        // ...

        FairSync(int permits) {
            super(permits);
        }

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

    // ...
}
// 基于 JDK 17, java.util.concurrent.Semaphore —— 对外 API,全是对 sync 的一行转发
public class Semaphore implements java.io.Serializable {
    // ...

    public Semaphore(int permits) {
        sync = new NonfairSync(permits);    // 默认非公平;Doug Lea 的注释建议:用来控制资源访问时应显式传 true 用公平版,防饿死
    }

    public Semaphore(int permits, boolean fair) {
        sync = fair ? new FairSync(permits) : new NonfairSync(permits);
    }

    public void acquire() throws InterruptedException {
        sync.acquireSharedInterruptibly(1);
    }

    public void acquireUninterruptibly() {
        sync.acquireShared(1);      // 不可中断版:中断只记标记不抛异常,返回时补上中断状态
    }

    public boolean tryAcquire() {
        return sync.nonfairTryAcquireShared(1) >= 0;    // 坑:写死 nonfair,公平模式下它照样插队。想守公平用 tryAcquire(0, TimeUnit.SECONDS)
    }

    public boolean tryAcquire(long timeout, TimeUnit unit)
        throws InterruptedException {
        return sync.tryAcquireSharedNanos(1, unit.toNanos(timeout));
    }

    public void release() {
        sync.releaseShared(1);      // 不校验"你还的是不是你拿的"——信号量没有 owner 概念,谁都能 release
    }

    public void acquire(int permits) throws InterruptedException {
        if (permits < 0) throw new IllegalArgumentException();
        sync.acquireSharedInterruptibly(permits);   // 批量获取是原子的:要么一次拿够,要么一个不扣地等
    }

    public void release(int permits) {
        if (permits < 0) throw new IllegalArgumentException();
        sync.releaseShared(permits);
    }

    public int availablePermits() {
        return sync.getPermits();
    }

    public int drainPermits() {
        return sync.drainPermits();
    }

    protected void reducePermits(int reduction) {
        if (reduction < 0) throw new IllegalArgumentException();
        sync.reducePermits(reduction);  // protected:只给子类用的"资源永久缩容"钩子
    }

    public boolean isFair() {
        return sync instanceof FairSync;
    }
    // ...
}

原理串讲

以默认的非公平信号量为例,走一遍”许可不够时 acquire 再被 release 救活”的完整链路。
线程 A 调 acquire(),转发给 AQS 的 acquireSharedInterruptibly(1):先检查中断,然后调子类的 tryAcquireShared(1),也就是 nonfairTryAcquireShared——读 state、算 remaining,不够则直接返回负数。
返回负数后 AQS 进入 acquire(null, arg, shared=true, ...):把 A 包成 SharedNode 挂到队尾,LockSupport.park 挂起。
注意失败路径上 state 一个字节都没动,这是第一处设计:为什么许可不够时不做 CAS 就返回? 因为 acquire 失败必须零副作用——state 若被减成负数,后续 release 的加法语义、availablePermits 的读数全都会乱掉;同时也省掉一次注定无意义的 CAS。
想要”负许可”的语义,走的是另一个专门口子 reducePermits

线程 B 调 release() → AQS 的 releaseShared(1) → 子类 tryReleaseShared(1) CAS 把 state 加回去,成功后 AQS 调 signalNext(head) unpark 队头的 A。
A 醒来后并不是”直接拿到许可”,而是回到循环里重新调 tryAcquireShared 抢——非公平模式下这一刻若有新线程插队,A 可能再次失败再次 park。
抢成功后 A 把自己设为新 head,并调 signalNextIfShared(node) 看下一个节点是不是也是共享节点、要不要接着唤醒。
这依赖第二处设计:为什么 tryAcquireShared 返回 int 而不像独占模式返回 boolean? 因为共享模式需要三态:负数=失败去排队,0=成功但许可耗尽,正数=成功且还有剩余。
返回正数时 AQS 才知道”后面的兄弟可能也能过”,于是唤醒链式传播下去——比如 release(3) 一次补三个许可,队里三个等 1 个许可的线程会被接力唤醒,而不需要 release 方唤醒三次。

还有一处容易被问到:为什么 release 不校验归属? 信号量的语义是”计数”不是”锁”,没有 getExclusiveOwnerThread 这种 owner 状态(对比 ReentrantLock 的 tryRelease 会先验持锁线程)。
这让 Semaphore 能做跨线程信号传递——A 挂起等许可、B 来发放,甚至能故意多 release 实现动态扩容;代价是”忘了 release”或”多 release”编译器和运行时都不会拦,只能靠 try/finally 的编程约定兜底。
排队、park/unpark、唤醒传播这些骨架逻辑全在 AbstractQueuedSynchronizer 里,Semaphore 本体只回答一个问题:“state 该怎么增减”。

设计取舍

  • state 三个入口三种规则:acquire 只减不许变负、release 只加防溢出、reducePermits 可以扣成负——同一个 int 靠调用方语义区分,没有额外字段。
  • CountDownLatch 同是”借 state 当计数的 AQS 共享模式”,区别是 Semaphore 的计数可增可减可复用,latch 只减不增、一次性。
  • 公平/非公平的实现差异只有一行 hasQueuedPredecessors(),和 ReentrantLock 完全同一套路;但无参 tryAcquire() 例外,公平模式下也插队。
  • 初始 permits=1 可当互斥锁用(binary semaphore),但它不可重入、也没有 owner:同一线程 acquire 两次会自己锁死自己,这是和 ReentrantLock 的高频对比题。
  • 批量 acquire(n) 原子生效,但公平模式下队头一个要 3 个许可的线程会挡住后面只要 1 个的——大额请求可能拖慢整体吞吐。

延伸阅读