NOTE

6.39 Fair Semaphore

1. What it is. Rate limiting using a fair strategy. 2. Usage. 3. Principle analysis: constructor, FairSync, acquire, shared acquisition through AQS, fair acquisition, release, and shared release.

JavaCreated Updated 1 min readhistorical

This is a historical learning note and may contain outdated or incomplete understanding.

1. What It Is

Rate limiting using a fair strategy.

2. Usage

public class SemaphoreTest
{
    private final static int THREAD_COUNT = 100;
    private final static CountDownLatch countDownLatch = new CountDownLatch(THREAD_COUNT);

    public static void main(String[] args) throws InterruptedException
    {
        Semaphore semaphore = new Semaphore(10, true);// true means fair

        for (int i = 0; i < THREAD_COUNT; i++)
        {
            new Thread(()->{
                try
                {
                    semaphore.acquire();
                    System.out.println(Thread.currentThread().getName() + " is accessing the resource...");
                    TimeUnit.SECONDS.sleep(3);

                }
                catch (Exception e)
                {
                    e.printStackTrace();
                }
                finally
                {
                    semaphore.release();
                    countDownLatch.countDown();
                }
            }).start();
        }

        countDownLatch.await();
    }
}

3. Principle Analysis

3.1. Constructor

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

3.1.1. Fair Sync

  • FairSync
static final class FairSync extends Sync {
    private static final long serialVersionUID = 2014338818796000944L;

    FairSync(int permits) {
        // Semaphore.Sync#Sync
        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;
        }
    }
}
  • Sync
abstract static class Sync extends AbstractQueuedSynchronizer {
    private static final long serialVersionUID = 1192457210091910933L;

    Sync(int permits) {
        // Ultimately, set permits permits
        setState(permits);
    }
}

3.2. acquire

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

3.2.1. Call AQS to Acquire a Shared Lock

  • AQS.acquireSharedInterruptibly
public final void acquireSharedInterruptibly(int arg)
        throws InterruptedException {
    if (Thread.interrupted())
        throw new InterruptedException();
    // Semaphore.FairSync overrides tryAcquireShared.
    // If there are not enough permits, it returns a negative value.
    // doAcquireSharedInterruptibly then joins the AQS queue and blocks while waiting to be woken.
    // If there are enough permits, it returns a value >= 0, and the caller of acquire can continue executing business logic.
    if (tryAcquireShared(arg) < 0)
        doAcquireSharedInterruptibly(arg);
}
3.2.1.1. Try to Acquire [Fair: If Someone Is Queued Ahead, Fail Directly]
  • Semaphore.FairSync#tryAcquireShared
protected int tryAcquireShared(int acquires) {
    for (;;) {
        // If someone is queued ahead of me, return -1
        if (hasQueuedPredecessors())
            return -1;
        // Current number of permits
        int available = getState();
        // Are there enough permits for me to acquire?
        int remaining = available - acquires;
        // If < 0, there are not enough permits, so return this value
        if (remaining < 0 ||
             // >= 0 means there are enough permits, so CAS-update the remaining permits
            compareAndSetState(available, remaining))
            return remaining;
    }
}

3.3. release

  • Semaphore#release()
public void release() {
    // AQS.releaseShared
    sync.releaseShared(1);
}

3.3.1. Call AQS to Release a Shared Lock

  • AQS#releaseShared
public final boolean releaseShared(int arg) {
    // Semaphore.Sync overrides tryReleaseShared
    if (tryReleaseShared(arg)) {
        doReleaseShared();
        return true;
    }
    return false;
}
3.3.1.1. Try to Release the Shared Lock
  • Semaphore.Sync#tryReleaseShared
protected final boolean tryReleaseShared(int releases) {
    for (;;) {
        // Get the current number of permits
        int current = getState();
        // Add the permits back
        int next = current + releases;
        // Throw an exception on overflow
        if (next < current) // overflow
            throw new Error("Maximum permit count exceeded");
        // CAS-update the number of permits
        if (compareAndSetState(current, next))
            return true;
    }
}

Discussion

Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub