NOTE

6.31 Nonfair ReadWriteLock

1. What it is. A nonfair read-write lock may let a newly arriving thread compete without strictly following queue order. 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

Regardless of whether threads are already queued ahead waiting for the lock, a newly arriving thread may still try to compete for it directly.

2. How to Use It

public class ReadWriteLockTest
{
    private static ReadWriteLock lock = new ReentrantReadWriteLock();// Nonfair by default
    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() {
    // false by default
    this(false);
}

public ReentrantReadWriteLock(boolean fair) {
    // Initialize Sync
    // If false, use NonfairSync
    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);
    // Nonfair lock: call NonfairSync.readerShouldBlock to determine 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 [Nonfair]
  • ReentrantReadWriteLock.NonfairSync#readerShouldBlock
final boolean readerShouldBlock() {

    // A read blocks only when the first queued node is an exclusive node
    // (a node waiting for the write lock).
    // This is the characteristic of the nonfair lock:
    // even if other read-lock threads are already waiting ahead in the queue,
    // the current reader does not care.
    return apparentlyFirstQueuedIsExclusive();
}
  • apparentlyFirstQueuedIsExclusive
final boolean apparentlyFirstQueuedIsExclusive() {
    Node h, s;

    // The queue is not empty
    return (h = head) != null &&
        // And the actual queue head is not null
        // (the AQS head itself is a placeholder)
        (s = h.next)  != null &&
        // And the actual queue head is not a shared node
        // (meaning it is not waiting for a read lock)
        !s.isShared()         &&
        // And the actual queue-head thread is not null
        s.thread != null;
    // Only when all conditions above hold does the current read-lock thread need to block
}
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 [Nonfair]
  • ReentrantReadWriteLock.NonfairSync#writerShouldBlock
final boolean writerShouldBlock() {
    // A writer does not block by default
    return false; // writers can always barge
}
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