NOTE

6.38 Fair ReadWriteLock

1. What it is. A fair read-write lock follows queue order: if threads are already waiting ahead, the current thread should not barge in. 2. How to use it. 3. Source-code analysis: constructor, read locking/unlocking, write locking/unlocking, and AQS integration.

JavaCreated Updated 1 min readhistorical

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

1. What It Is

A fair read-write lock follows queue order. If threads are already waiting ahead, the current thread should not barge in and compete for the lock directly.

2. How to Use It

public class ReadWriteLockTest
{
    private static ReadWriteLock lock = new ReentrantReadWriteLock(true);// true means fair
    private static Lock readLock = lock.readLock();
    private static Lock writeLock = lock.writeLock();

    private static List<Integer> data = new ArrayList<>();

    public static void main(String[] args) throws InterruptedException
    {
        Thread readThread = new Thread(() -> {
            while (true)
            {
                try
                {
                    TimeUnit.MILLISECONDS.sleep(500);
                    readLock.lock();
                    System.out.println(Thread.currentThread().getName() + " read: " + data);
                }
                catch (InterruptedException e)
                {
                    e.printStackTrace();
                }
                finally
                {
                    readLock.unlock();
                }
            }
        });
        Thread readThread2 = new Thread(() -> {
            while (true)
            {
                try
                {
                    TimeUnit.MILLISECONDS.sleep(300);

                    readLock.lock();
                    System.out.println(Thread.currentThread().getName() + " read: " + data);
                }
                catch (InterruptedException e)
                {
                    e.printStackTrace();
                }
                finally
                {
                    readLock.unlock();
                }
            }
        });
        Thread writeThread = new Thread(() -> {

            int i = 0;
            while (true)
            {
                try
                {
                    TimeUnit.MILLISECONDS.sleep(200);

                    writeLock.lock();
                    if (i % 2 == 0)
                    {
                        data.add(i);
                    }else
                    {
                        data.remove(0);
                    }
                    i++;
                }
                catch (InterruptedException e)
                {
                    e.printStackTrace();
                }
                finally
                {
                    writeLock.unlock();
                }
            }
        });

        readThread.start();
        readThread2.start();
        writeThread.start();

        readThread.join();
        readThread2.join();
        writeThread.join();

    }

}

3. Source-Code Analysis

3.1. UML

PlantUML 图表

3.2. Constructor

  • ReentrantReadWriteLock
public ReentrantReadWriteLock(boolean fair) {
    // Initialize Sync
    // If true, use FairSync
    sync = fair ? new FairSync() : new NonfairSync();
    // Initialize the read and write locks
    readerLock = new ReadLock(this);
    writerLock = new WriteLock(this);
}
  • ReentrantReadWriteLock.ReadLock
protected ReadLock(ReentrantReadWriteLock lock) {
    // Essentially just save ReentrantReadWriteLock's Sync
    sync = lock.sync;
}
  • ReentrantReadWriteLock.WriteLock
protected WriteLock(ReentrantReadWriteLock lock) {
    // Essentially just save ReentrantReadWriteLock's Sync
    sync = lock.sync;
}

3.3. Acquiring the Read Lock

  • ReentrantReadWriteLock.readLock
// Returns the read lock
public ReentrantReadWriteLock.ReadLock  readLock()  { return readerLock; }
  • ReadLock.lock
public void lock() {
    // Use AQS to acquire a shared lock
    sync.acquireShared(1);
}

3.3.1. Use AQS to Acquire a Shared Lock

  • AQS.acquireShared
public final void acquireShared(int arg) {
    // ReentrantReadWriteLock.Sync overrides tryAcquireShared,
    // so Sync.tryAcquireShared is called here.

    // A return value < 0 means lock acquisition failed:
    // execute doAcquireShared, join the AQS queue, block, and wait to be woken.
    // A return value >= 0 means lock acquisition succeeded, so continue with the business logic.
    if (tryAcquireShared(arg) < 0)
        doAcquireShared(arg);
}
3.3.1.1. Use Sync to Try to Acquire the Shared Lock
  • ReentrantReadWriteLock.Sync.tryAcquireShared
protected final int tryAcquireShared(int unused) {
    // Current thread
    Thread current = Thread.currentThread();
    // Current state value (the number of locks, including both read and write locks)
    int c = getState();
    // Get the number of exclusive locks from state (the number of write locks).
    // A nonzero value (> 0) means a write lock has already been acquired.
    if (exclusiveCount(c) != 0 &&
        // Check whether the thread holding the write lock is the current thread
        getExclusiveOwnerThread() != current)
        // If not, return -1: another thread holds the write lock.
        // Later, the current thread must join the AQS queue and block until it is woken.
        return -1;

    // Reaching here means no thread holds the write lock,
    // or the thread holding it is the current thread.

    // Get the number of shared locks from state (the number of read locks)
    int r = sharedCount(c);
    // Fair lock: readerShouldBlock determines whether the read should block.
    // If false, it does not need to block and the rest of the && expression is evaluated.
    if (!readerShouldBlock() &&
        // Check whether the number of read locks is below the maximum
        r < MAX_COUNT &&
        // CAS-acquire a read lock
        compareAndSetState(c, c + SHARED_UNIT)) {

        // After the if condition succeeds, the read-lock portion of state has already been updated.
        // The following logic updates the other related fields.

        // First read-lock acquisition
        if (r == 0) {
            // Record the first reader and initialize its read-lock count
            firstReader = current;
            firstReaderHoldCount = 1;
        // The nth acquisition is still by the first reader
        } else if (firstReader == current) {
            // Only update the read-lock count
            firstReaderHoldCount++;
        // The nth acquisition is by a thread other than the first reader
        } else {
            // cachedHoldCounter is a cache
            HoldCounter rh = cachedHoldCounter;
            if (rh == null || rh.tid != getThreadId(current))
                cachedHoldCounter = rh = readHolds.get();
            else if (rh.count == 0)
                readHolds.set(rh);
            rh.count++;
        }
        // Return 1 (positive), indicating that the read lock was acquired successfully
        return 1;
    }

    // Reaching here means one of the following happened:
    // 1. the read needs to block
    // 2. the read-lock count has exceeded the maximum
    // 3. CAS acquisition of the read lock failed
    return fullTryAcquireShared(current);
}
3.3.1.1.1. Determine Whether the Read Should Block [Fair]
  • ReentrantReadWriteLock.FairSync#readerShouldBlock
final boolean readerShouldBlock() {
    // A read needs to block when there is a queued predecessor
    return hasQueuedPredecessors();
}
  • AQS.hasQueuedPredecessors
public final boolean hasQueuedPredecessors() {
    Node t = tail; // Read fields in reverse initialization order
    Node h = head;
    Node s;
    return h != t &&
        // If the actual queue head is null, or the thread at the actual head is not the current thread,
        // return true
        ((s = h.next) == null || s.thread != Thread.currentThread());
}
3.3.1.1.2. If the Fast Acquisition Attempt Fails, Use an Infinite Loop to Acquire
  • ReentrantReadWriteLock.Sync#fullTryAcquireShared
final int fullTryAcquireShared(Thread current) {

    HoldCounter rh = null;
    // Infinite loop
    for (;;) {
        int c = getState();
        // There is a write lock
        if (exclusiveCount(c) != 0) {
            // And it is not held by the current thread
            if (getExclusiveOwnerThread() != current)
                // Return -1
                return -1;
            // else we hold the exclusive lock; blocking here
            // would cause deadlock.
        // Nobody holds a write lock; if reads need to block
        } else if (readerShouldBlock()) {
            // Make sure we're not acquiring read lock reentrantly
            if (firstReader == current) {
                // assert firstReaderHoldCount > 0;
            } else {
                if (rh == null) {
                    rh = cachedHoldCounter;
                    if (rh == null || rh.tid != getThreadId(current)) {
                        rh = readHolds.get();
                        if (rh.count == 0)
                            readHolds.remove();
                    }
                }
                if (rh.count == 0)
                    return -1;
            }
        }
        // Reaching here means there is no write lock and the read does not need to block
        if (sharedCount(c) == MAX_COUNT)
            // Lock count has exceeded the maximum
            throw new Error("Maximum lock count exceeded");
        // Acquire the read lock, same as tryAcquireShared above
        if (compareAndSetState(c, c + SHARED_UNIT)) {
            if (sharedCount(c) == 0) {
                firstReader = current;
                firstReaderHoldCount = 1;
            } else if (firstReader == current) {
                firstReaderHoldCount++;
            } else {
                if (rh == null)
                    rh = cachedHoldCounter;
                if (rh == null || rh.tid != getThreadId(current))
                    rh = readHolds.get();
                else if (rh.count == 0)
                    readHolds.set(rh);
                rh.count++;
                cachedHoldCounter = rh; // cache for release
            }
            return 1;
        }
    }
}
3.3.1.2. If Lock Acquisition Fails, Join the AQS Queue, Block, and Wait to Be Woken to Continue Competing for the Lock
private void doAcquireShared(int arg) {
    // Construct a SHARED node, join the AQS queue, block, and wait to be woken
    final Node node = addWaiter(Node.SHARED);
    boolean failed = true;
    try {
        boolean interrupted = false;
        // Infinite loop competing for the lock
        for (;;) {
            // The predecessor of the current node is the head
            final Node p = node.predecessor();
            if (p == head) {
                // Continue trying to acquire the shared lock
                int r = tryAcquireShared(arg);
                if (r >= 0) {
                    setHeadAndPropagate(node, r);
                    p.next = null; // help GC
                    if (interrupted)
                        selfInterrupt();
                    failed = false;
                    return;
                }
            }
            // Determine whether blocking is needed
            if (shouldParkAfterFailedAcquire(p, node) &&
                // Block if necessary
                parkAndCheckInterrupt())
                interrupted = true;
        }
    } finally {
        if (failed)
            cancelAcquire(node);
    }
}
Join the AQS Queue and Block
  • addWaiter
private Node addWaiter(Node mode) {
    Node node = new Node(Thread.currentThread(), mode);
    // Try the fast path of enq; backup to full enq on failure
    // Quickly try to append to the tail
    Node pred = tail;
    if (pred != null) {
        node.prev = pred;
        if (compareAndSetTail(pred, node)) {
            pred.next = node;
            return node;
        }
    }
    // Fast-path append failed, so use enq to append to the tail
    enq(node);
    return node;
}

See 5.AQS.md.

3.4. Releasing the Read Lock

  • ReentrantReadWriteLock.ReadLock#unlock
public void unlock() {
    // Use AQS to release the shared lock
    sync.releaseShared(1);
}

3.4.1. Use AQS to Release the Shared Lock

  • AbstractQueuedSynchronizer#releaseShared
public final boolean releaseShared(int arg) {
    // Call Sync to try to release the shared lock.
    // If all shared locks [read locks] have been released, execute doReleaseShared.
    if (tryReleaseShared(arg)) {
        doReleaseShared();
        return true;
    }
    return false;
}
3.4.1.1. Use Sync to Try to Release the Lock
  • ReentrantReadWriteLock.Sync#tryReleaseShared
protected final boolean tryReleaseShared(int unused) {
    Thread current = Thread.currentThread();
    // The current thread is the first thread that acquired the read lock
    if (firstReader == current) {
        // All read locks have been released
        if (firstReaderHoldCount == 1)
            // Clear the first-reader thread
            firstReader = null;
        // Not all released, so decrement the read-lock count
        else
            firstReaderHoldCount--;
    // The current thread is not the first thread that acquired the read lock
    } else {
        HoldCounter rh = cachedHoldCounter;
        if (rh == null || rh.tid != getThreadId(current))
            rh = readHolds.get();
        int count = rh.count;
        if (count <= 1) {
            readHolds.remove();
            if (count <= 0)
                throw unmatchedUnlockException();
        }
        --rh.count;
    }
    // Infinite loop modifying the state count
    for (;;) {
        int c = getState();
        int nextc = c - SHARED_UNIT;
        // CAS-update the state count
        if (compareAndSetState(c, nextc))
            return nextc == 0;// Return true when it reaches 0
    }
}
3.4.1.2. After All Shared Locks Are Released, Wake the Node After the Head in the AQS Queue
  • AbstractQueuedSynchronizer#doReleaseShared
private void doReleaseShared() {

    for (;;) {
        Node h = head;
        // The queue is not empty
        if (h != null && h != tail) {
            int ws = h.waitStatus;
            // The head's state is SIGNAL [it promises to wake the next node]
            if (ws == Node.SIGNAL) {
                // CAS-set the state to 0
                if (!compareAndSetWaitStatus(h, Node.SIGNAL, 0))
                    continue;            // loop to recheck cases
                // Wake the node after the head
                unparkSuccessor(h);
            }
            else if (ws == 0 &&
                     !compareAndSetWaitStatus(h, 0, Node.PROPAGATE))
                continue;                // loop on failed CAS
        }
        if (h == head)                   // loop if head changed
            break;
    }
}
3.6.1.2. Wake the Node After the Head in the AQS Queue
  • unparkSuccessor
private void unparkSuccessor(Node node) {

    int ws = node.waitStatus;
    // If the current node's state is negative
    if (ws < 0)
        // Change it to 0.
        // 0 is an empty state because after this node's thread releases the lock,
        // there is nothing else it needs to do
        compareAndSetWaitStatus(node, ws, 0);


     // Get the current node's next node
    Node s = node.next;
    // If the next node is null (the current node is the tail) or its state > 0 (cancelled)
    if (s == null || s.waitStatus > 0) {
        s = null;
        // Traverse backward from the tail to the current node
        for (Node t = tail; t != null && t != node; t = t.prev)
            // Find the nearest node to the current node whose state <= 0 (not cancelled)
            if (t.waitStatus <= 0)
                s = t;
    }
    // Wake the current node's successor
    if (s != null)
        LockSupport.unpark(s.thread);
}

3.5. Acquiring the Write Lock

  • writeLock
// Returns the write lock
public ReentrantReadWriteLock.WriteLock writeLock() { return writerLock; }
  • ReentrantReadWriteLock.WriteLock#lock
public void lock() {
    // Use AQS to acquire the lock
    sync.acquire(1);
}

3.5.1. Call AQS to Acquire the Exclusive Lock

  • AQS.acquire
public final void acquire(int arg) {
    // If tryAcquire succeeds, it returns true and business logic continues.
    // If it fails, it returns false; then join the blocking queue, block, and wait to be woken.
    if (!tryAcquire(arg) &&
        acquireQueued(addWaiter(Node.EXCLUSIVE), arg))
        selfInterrupt();
}
3.5.1.1. Use Sync to Try to Acquire the Exclusive Lock
  • ReentrantReadWriteLock.Sync#tryAcquire
protected final boolean tryAcquire(int acquires) {

    // Get the current thread
    Thread current = Thread.currentThread();
    // Get the state count [the number of read/write locks]
    int c = getState();
    // Get the number of write locks from state
    int w = exclusiveCount(c);
    // A write lock has already been acquired
    if (c != 0) {
        // If the write lock is not held by me, return false to indicate acquisition failure.
        // (Note: if c != 0 and w == 0 then shared count != 0)
        if (w == 0 || current != getExclusiveOwnerThread())
            return false;
        // I hold the lock: reenter
        if (w + exclusiveCount(acquires) > MAX_COUNT)
            throw new Error("Maximum lock count exceeded");
        // Update state
        setState(c + acquires);
        return true;
    }
    // Reaching here means nobody currently holds a write lock.
    // The write may need to block
    if (writerShouldBlock() ||
        // Or CAS acquisition failed
        !compareAndSetState(c, c + acquires))
        // Return false
        return false;

    // Reaching here means the write lock was acquired successfully.
    // Set the current thread as the lock owner.
    setExclusiveOwnerThread(current);
    return true;
}
3.5.1.1.1. Determine Whether the Write Should Block [Fair]
  • ReentrantReadWriteLock.FairSync#writerShouldBlock
final boolean writerShouldBlock() {
    return hasQueuedPredecessors();
}
  • AQS.hasQueuedPredecessors
public final boolean hasQueuedPredecessors() {
    Node t = tail; // Read fields in reverse initialization order
    Node h = head;
    Node s;
    return h != t &&
        // If the actual queue head is null, or its thread is not the current thread,
        // return true
        ((s = h.next) == null || s.thread != Thread.currentThread());
}
3.5.1.2. If Lock Acquisition Fails, Join the AQS Queue, Block, and Wait to Be Woken to Continue Competing
  • addWaiter
//...
  • acquireQueued
//...

See 5.AQS.md.

3.6. Releasing the Write Lock

  • ReentrantReadWriteLock.WriteLock#unlock
public void unlock() {
    // Use AQS to release the lock
    sync.release(1);
}

3.6.1. Call AQS to Release the Exclusive Lock

  • AbstractQueuedSynchronizer#release
public final boolean release(int arg) {
    // Call Sync to try to release the exclusive lock.
    // If all exclusive state has been released, return true and execute the logic below.
    if (tryRelease(arg)) {
        Node h = head;
        if (h != null && h.waitStatus != 0)
            // Wake the node after the head in the AQS queue
            unparkSuccessor(h);
        return true;
    }
    return false;
}
3.6.1.1. Call Sync to Try to Release the Exclusive Lock
  • ReentrantReadWriteLock.Sync#tryRelease
protected final boolean tryRelease(int releases) {
    if (!isHeldExclusively())
        throw new IllegalMonitorStateException();
    // Decrement state
    int nextc = getState() - releases;
    // Has the write-lock exclusive count reached 0?
    boolean free = exclusiveCount(nextc) == 0;
    if (free)
        // If it reached 0, clear the thread that owns the lock
        setExclusiveOwnerThread(null);
    // Update state
    setState(nextc);
    return free;
}
3.6.1.2. After Unlocking Succeeds, Wake the Node After the Head in the AQS Queue
  • unparkSuccessor
private void unparkSuccessor(Node node) {

    int ws = node.waitStatus;
    // If the current node's state is negative
    if (ws < 0)
        // Change it to 0.
        // 0 is an empty state because after this node's thread releases the lock,
        // there is nothing else it needs to do
        compareAndSetWaitStatus(node, ws, 0);


     // Get the current node's next node
    Node s = node.next;
    // If the next node is null (the current node is the tail) or its state > 0 (cancelled)
    if (s == null || s.waitStatus > 0) {
        s = null;
        // Traverse backward from the tail to the current node
        for (Node t = tail; t != null && t != node; t = t.prev)
            // Find the nearest node to the current node whose state <= 0 (not cancelled)
            if (t.waitStatus <= 0)
                s = t;
    }
    // Wake the current node's successor
    if (s != null)
        LockSupport.unpark(s.thread);
}

Discussion

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