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.
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