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.
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
-
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.
-
A separate field is needed to represent the current lock state: free or occupied.
-
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.
-
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
- Consider using a field
Thread lockHolderto represent the thread currently holding the lock. In a multithreaded environment, usevolatileso 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
- Use
int stateto 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, usevolatileso 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
- Use
LockSupport.parkto block a thread andLockSupport.unparkto wake a thread.
2.4. Use a Queue to Store Threads That Failed to Acquire the Lock
- 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