NOTE

6.25 Implementing a Simple AQS by Hand

Build a simple AQS-like lock by hand to better understand the real AQS source code: requirements, fields, blocking/waking, waiter queue, lock/unlock flow, fairness, final implementation, test, and flow.

JavaCreated Updated 2 min readhistorical

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

We can write a simple AQS ourselves to better understand the actual AQS source code.

1. Requirements

  1. The lock is exclusive. Once a thread holds this lock, no other thread can hold it until it is released. Therefore, the thread currently holding the lock needs to be stored.

  2. A separate field is needed to represent the current lock state: free or occupied.

  3. Many threads may compete for the lock at the same time, but only one thread can succeed. What should happen to the other threads?

    First, the other threads should temporarily stop competing for the lock, meaning they should be blocked. Since blocking exists, there must also be a wake-up operation. When the thread holding the lock releases it, it needs to wake other threads that are blocked waiting for the lock.

  4. A queue is also needed to store threads that failed to acquire the lock so that they can later be woken and continue competing for it. Note that operations on this queue must be thread-safe.

2. Define the Fields

2.1. Lock Exclusivity

  1. Consider using a field Thread lockHolder to represent the thread currently holding the lock. In a multithreaded environment, use volatile so that changes to this variable can be observed by other threads in time.
public class MyLock
{
    // Represents the thread currently holding the lock
    private volatile Thread lockHolder;
}

2.2. Lock State

  1. Use int state to represent the lock state: 0 means free/not held, and 1 means held. Operations on it must be atomic, so use CAS. Likewise, in a multithreaded environment, use volatile so that changes can be observed by other threads in time.
public class MyLock
{
    // Represents lock state. Records the number of times the lock is acquired
    private volatile int state = 0;
    //...
}

In Java, CAS operations require an instance of Unsafe:

public class MyLock
{
    //...

    private static final Unsafe unsafe = UnsafeInstance.getInstance();
    // Offset address of the state variable. Needed when using Unsafe to perform CAS
    private static final long stateOffset;

    static
    {
        try
        {
            stateOffset = unsafe.objectFieldOffset(MyLock.class.getDeclaredField("state"));
        }catch (Exception e)
        {
            throw new Error();
        }
    }
    // CAS-set state
    public final boolean compareAndSwapState(int except, int update)
    {
        return unsafe.compareAndSwapInt(this, stateOffset, except, update);
    }


    // Obtain the Unsafe instance through reflection
    private static class UnsafeInstance
    {
        public static Unsafe getInstance()
        {
            try
            {
                Field field = Unsafe.class.getDeclaredField("theUnsafe");
                field.setAccessible(true);
                return (Unsafe) field.get(null);
            }
            catch (Exception e)
            {
                e.printStackTrace();
            }
            return null;
        }
    }
}

2.3. Block and Wake Threads

  1. Use LockSupport.park to block a thread and LockSupport.unpark to wake a thread.

2.4. Use a Queue to Store Threads That Failed to Acquire the Lock

  1. For a thread-safe queue, consider using ConcurrentLinkedQueue:
public class MyLock
{
    // Store threads that have not acquired the lock
    private ConcurrentLinkedQueue<Thread> waiters = new ConcurrentLinkedQueue<>();
    //...
}
  • Putting the code above together:
public class MyLock
{
    // Represents lock state. Records the number of times the lock is acquired
    private volatile int state = 0;
    // Represents the thread currently holding the lock
    private volatile Thread lockHolder;
    // Store threads that have not acquired the lock
    private ConcurrentLinkedQueue<Thread> waiters = new ConcurrentLinkedQueue<>();

    //==================UNSAFE======================//
    private static final Unsafe unsafe = UnsafeInstance.getInstance();
    // Offset address of the state variable. Needed when using Unsafe to perform CAS
    private static final long stateOffset;

    static
    {
        try
        {
            stateOffset = unsafe.objectFieldOffset(MyLock.class.getDeclaredField("state"));
        }catch (Exception e)
        {
            throw new Error();
        }
    }
    // CAS-set state
    public final boolean compareAndSwapState(int except, int update)
    {
        return unsafe.compareAndSwapInt(this, stateOffset, except, update);
    }


    // Obtain the Unsafe instance through reflection
    private static class UnsafeInstance
    {
        public static Unsafe getInstance()
        {
            try
            {
                Field field = Unsafe.class.getDeclaredField("theUnsafe");
                field.setAccessible(true);
                return (Unsafe) field.get(null);
            }
            catch (Exception e)
            {
                e.printStackTrace();
            }
            return null;
        }
    }

3. Add Lock and Unlock Operations

3.1. Basic Flow

The pseudocode is as follows:

  • Check the lock state.
  • If no thread holds the lock, try to acquire it (CAS atomic operation).
  • If acquisition succeeds, set the thread currently holding the lock.
  • If acquisition fails, enqueue and block.
// Lock operation
public void lock()
{

    Thread currentThread = Thread.currentThread();
    int state = getState();
    // state == 0 means not locked
    if (state == 0)
    {
        // If CAS succeeds, lock acquisition succeeds
        if (compareAndSwapState(0, 1))
        {
            setLockHolder(currentThread);// Set the thread currently holding the lock
            return;
        }
    }


    // Failed to acquire the lock; enqueue
    waiters.add(currentThread);

    // Failed to acquire the lock; block the current thread
    LockSupport.park(currentThread);// Woken by unpark

}
  • Unlock

The pseudocode is as follows:

  • The thread unlocking (the current thread) must be the thread holding the lock.
  • If so, use CAS to unlock.
  • If unlocking succeeds, clear the current lock-holder thread and wake blocked threads.
public void unlock()
{
    // The locking and unlocking thread must be the same
    if (Thread.currentThread() != lockHolder)
    {
        throw new RuntimeException("current thread is not lockHolder");
    }

    // CAS-set state successfully: unlock
    int state = getState();
    if (compareAndSwapState(state, 0))
    {
        setLockHolder(null);
        // Wake all threads waiting for the lock
        for (Thread waiter : waiters)
        {
           LockSupport.unpark(head);// Wake a parked thread
        }

    }
}

3.2. After Being Woken, Continue Competing for the Lock

After the locking code above is woken, it needs to try to acquire the lock again. If it fails, it blocks; after being woken, it tries again…

This continues until it successfully sets the current thread as the lock holder and removes the current thread from the waiting queue.

This is clearly an infinite loop. Modify the locking code as follows:

public void lock()
{
    // Lock acquisition succeeded
    if (acquire())
    {
        return;
    }

    // Failed to acquire the lock; enqueue
    Thread currentThread = Thread.currentThread();
    waiters.add(currentThread);

    for (;;)
    {
        // Keep trying to acquire the lock
        if (acquire())
        {
            return;
        }
        // Failed to acquire the lock; block the current thread
        LockSupport.park(currentThread);// Woken by unpark
    }
}

private boolean acquire()
{
    Thread currentThread = Thread.currentThread();
    int state = getState();
    // state == 0 means not locked
    if (state == 0)
    {
        // If CAS succeeds, lock acquisition succeeds
        if (compareAndSwapState(0, 1))
        {
            waiters.remove(Thread.currentThread());// Remove from the wait queue after acquiring the lock
            setLockHolder(currentThread);// Set the thread currently holding the lock
            return true;
        }
    }
    return false;
}

3.3. Add the Fair-Lock Property

A fair lock follows a first-come, first-served principle. Threads that fail to acquire the lock are placed in a queue and then blocked. When the thread holding the lock releases it successfully, the thread that has waited the longest (the queue head) should be woken so it can try to acquire the lock.

Modify the code as follows:

public void lock()
{
    // Lock acquisition succeeded
    if (acquire())
    {
        return;
    }

    // Failed to acquire the lock; enqueue
    Thread currentThread = Thread.currentThread();
    waiters.add(currentThread);

    for (;;)
    {
        // Only when the queue is empty (nobody is waiting) or the current thread is the queue head
        // (meaning I am the first waiter) can the thread try to acquire the lock; that is what makes it fair
        // Keep trying to acquire the lock
        if ((waiters.isEmpty() || currentThread == waiters.peek()) && acquire())
        {
            return;
        }
        // Failed to acquire the lock; block the current thread
        LockSupport.park(currentThread);// Woken by unpark
    }
}

private boolean acquire()
{
    Thread currentThread = Thread.currentThread();
    int state = getState();
    // state == 0 means not locked
    if (state == 0)
    {
        // Only try to acquire the lock when the queue is empty (nobody is waiting ahead),
        // or when the current thread is the queue head; this makes it fair
        // If CAS succeeds, lock acquisition succeeds
        if ((waiters.isEmpty() || currentThread == waiters.peek()) && compareAndSwapState(0, 1))
        {
            waiters.poll();// Remove from the wait queue after acquiring the lock
            setLockHolder(currentThread);// Set the thread currently holding the lock
            return true;
        }
    }
    return false;
}

public void unlock()
{
    // The locking and unlocking thread must be the same
    if (Thread.currentThread() != lockHolder)
    {
        throw new RuntimeException("current thread is not lockHolder");
    }

    // CAS-set state successfully: unlock
    int state = getState();
    if (compareAndSwapState(state, 0))
    {
        setLockHolder(null);
        // Wake the queue-head thread waiting for the lock
        Thread head = waiters.peek();
        if (head != null)
        {
            LockSupport.unpark(head);// Wake a parked thread
        }
    }
}

4. Final Version

public class MyLock
{
    // Represents lock state. Records the number of times the lock is acquired
    private volatile int state = 0;
    // Represents the thread currently holding the lock
    private volatile Thread lockHolder;
    // Store threads that have not acquired the lock
    private ConcurrentLinkedQueue<Thread> waiters = new ConcurrentLinkedQueue<>();

    public void lock()
    {
        // Lock acquisition succeeded
        if (acquire())
        {
            return;
        }

        // Failed to acquire the lock; enqueue
        Thread currentThread = Thread.currentThread();
        waiters.add(currentThread);

        for (;;)
        {
            // Only when the queue is empty (nobody is waiting) or the current thread is the queue head
            // (meaning I am the first waiter) can the thread try to acquire the lock; that is what makes it fair
            // Keep trying to acquire the lock
            if ((waiters.isEmpty() || currentThread == waiters.peek()) && acquire())
            {
                return;
            }
            // Failed to acquire the lock; block the current thread
            LockSupport.park(currentThread);// Woken by unpark
        }
    }

    private boolean acquire()
    {
        Thread currentThread = Thread.currentThread();
        int state = getState();
        // state == 0 means not locked
        if (state == 0)
        {
            // Only try to acquire the lock when the queue is empty (nobody is waiting ahead),
            // or when the current thread is the queue head; this makes it fair
            // If CAS succeeds, lock acquisition succeeds
            if ((waiters.isEmpty() || currentThread == waiters.peek()) && compareAndSwapState(0, 1))
            {
                waiters.poll();// Remove from the wait queue after acquiring the lock
                setLockHolder(currentThread);// Set the thread currently holding the lock
                return true;
            }
        }
        return false;
    }

    public void unlock()
    {
        // The locking and unlocking thread must be the same
        if (Thread.currentThread() != lockHolder)
        {
            throw new RuntimeException("current thread is not lockHolder");
        }

        // CAS-set state successfully: unlock
        int state = getState();
        if (compareAndSwapState(state, 0))
        {
            setLockHolder(null);
            // Wake the queue-head thread waiting for the lock
            Thread head = waiters.peek();
            if (head != null)
            {
                LockSupport.unpark(head);// Wake a parked thread
            }
        }
    }

    public int getState()
    {
        return state;
    }

    public void setState(int state)
    {
        this.state = state;
    }

    public Thread getLockHolder()
    {
        return lockHolder;
    }

    public void setLockHolder(Thread lockHolder)
    {
        this.lockHolder = lockHolder;
    }


    //==================UNSAFE======================//

    private static final Unsafe unsafe = UnsafeInstance.getInstance();
    // Offset address of the state variable. Needed when using Unsafe to perform CAS
    private static final long stateOffset;

    static
    {
        try
        {
            stateOffset = unsafe.objectFieldOffset(MyLock.class.getDeclaredField("state"));
        }catch (Exception e)
        {
            throw new Error();
        }
    }
    // CAS-set state
    public final boolean compareAndSwapState(int except, int update)
    {
        return unsafe.compareAndSwapInt(this, stateOffset, except, update);
    }


    // Obtain the Unsafe instance through reflection
    private static class UnsafeInstance
    {
        public static Unsafe getInstance()
        {
            try
            {
                Field field = Unsafe.class.getDeclaredField("theUnsafe");
                field.setAccessible(true);
                return (Unsafe) field.get(null);
            }
            catch (Exception e)
            {
                e.printStackTrace();
            }
            return null;
        }
    }


}

5. Test

private static int result = 0;
public static void main(String[] args) throws InterruptedException
{
    final int threadCound = 10000;
    final CyclicBarrier barrier = new CyclicBarrier(threadCound);
    final CountDownLatch countDownLatch = new CountDownLatch(threadCound);
    final MyLock lock = new MyLock();

    for (int i = 0; i < threadCound; i++)
    {
        String name = "thread-" + i;
        new Thread(()->{
            try
            {
                barrier.await();
                lock.lock();
                result++;
                System.out.println(Thread.currentThread() + " result: " + result);
            }
            catch (Exception e)
            {
                e.printStackTrace();
            }
            finally
            {
                lock.unlock();
                countDownLatch.countDown();
            }
        }, name).start();
    }

    countDownLatch.await();

    System.out.println(result);

}

6. Flow

Discussion

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