Skip to content
JackSparrow414
Go back

Locks in Java (Part 2): Implementing a Custom Lock

Table of contents

Open Table of contents

Locks in Java (Part 2): Implementing a Custom Lock

Lock Semantics

  1. When a thread releases a lock, the JMM flushes shared variables from that thread’s local memory to main memory.
  2. When a thread acquires a lock, the JMM invalidates its local memory, forcing code in the monitor-protected critical section to read shared variables from main memory.

When thread A releases a lock, it effectively sends a message about its changes to shared variables to whichever thread acquires the lock next. When thread B acquires that lock, it receives the message from a thread that previously released it. Thus, A releasing a lock and B subsequently acquiring it amounts to A sending a message to B through main memory.

Custom Synchronization Components

Approach

Custom synchronization component = custom synchronizer + implementation of the Lock interface

The core acquire and release operations are implemented through AbstractQueuedSynchronizer, AQS, the queued synchronizer. Extend AQS and override the methods required by your scenario to create a custom synchronizer. As its documentation describes:

AbstractQueuedSynchronizer documentation describing custom synchronizers implemented through subclassing

Then implement Lock and call the AQS template methods from its methods to complete the custom synchronization component.

AQS Uses a FIFO Doubly Linked Queue. Why Doubly Linked?

AQS FIFO doubly linked queue with head, tail, prev, and next links

When multiple threads compete for an exclusive lock, those that fail to acquire it wait in the queue above and spin. When a thread releases the lock, other threads begin competing for it.

  1. Suppose it were a singly linked queue: head -> next -> tail. How could a thread check whether its predecessor is the head? The first-in, first-out principle would be difficult to enforce. Could we simply hand the lock to the next thread in order? Unfortunately, we cannot determine which thread will ultimately win it.
  2. With a doubly linked queue, head <-> next/prev <-> tail, threads check whether their predecessor is the head before competing for a released lock. If not, they continue spinning instead of competing; if so, they compete. FIFO is maintained because we can control enqueue order by appending nodes before tail. This must be atomic, using compareAndSetTail. The resulting queue has our intended order, and prev maintains dequeue order.

AQS Methods That Can Be Overridden

  1. protected boolean tryAcquire(int arg): exclusively acquire synchronization state, allowing only one thread at a time. Check whether the current state satisfies the requirements, then update it with CAS.
  2. protected boolean tryRelease(int arg): release synchronization state exclusively.
  3. protected int tryAcquireShared(int arg): acquire synchronization state in shared mode; a return value >= 0 indicates success.
  4. protected boolean tryReleaseShared(int arg): release synchronization state in shared mode.
  5. protected boolean isHeldExclusively(): generally checks whether the current thread is the one holding the lock.

AQS documentation listing five methods to override for a custom synchronizer The screenshot above comes from the AbstractQueuedSynchronizer documentation. Subclasses only need to override the five highlighted methods, using setState(), getState(), and compareAndSetState() to acquire and release locks.

AQS represents synchronization state with an int, state. Setting, reading, and modifying it in those methods requires atomic operations. Use the provided setState(), getState(), and compareAndSetState() methods.

Implement the Lock Interface

Implement Lock methods by calling AQS template methods. Internally, these call the five overridden methods above, implementing lock acquisition and release for the custom component.

Available AQS Template Methods

  1. void acquire(int arg) calls tryAcquire.
  2. boolean tryAcquireNanos(int arg, long nanos) adds a timeout.
  3. void acquireShared(int arg) calls tryAcquireShared.
  4. boolean tryAcquireSharedNanos(int arg, long nanos) adds a timeout.
  5. boolean release(int arg) calls tryRelease.
  6. boolean releaseShared(int arg) calls tryReleaseShared.
  7. Collection getQueuedThreads() returns the threads in the waiting queue.

Implementation

import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.AbstractQueuedSynchronizer;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.Lock;

/**
 * Allows at most two threads to access concurrently; additional threads are blocked.
 *
 * Implement the Lock interface to create a custom lock.
 *
 * Extend the synchronizer and override the required
 * methods, then compose Sync into the custom TwinsLock component
 * and call the synchronizer’s template methods.
 * These template methods call the overridden methods.
 * @author jacksparrow414
 * @date 2020/10/19
 */
public class TwinsLock implements Lock {

    /**
     * All custom synchronizers extend AQS.
     *
     * Override the following five methods selectively according to the required functionality.
     *
     * tryAcquire: exclusive; only one thread can acquire the lock.
     *
     * tryAcquireShared: shared; multiple threads can acquire the lock.
     *
     * For the two methods above, values >= 0 indicate success; otherwise acquisition fails.
     *
     * tryRelease: exclusive lock release.
     *
     * tryReleaseShared: release shared synchronization state.
     *
     * For both release methods, true indicates success; otherwise release fails.
     *
     * isHeldExclusively checks whether the current thread matches the thread in the synchronizer.
     */
    private static final class Sync extends AbstractQueuedSynchronizer {
        Sync(int count) {
            if (count <= 0) {
                throw new IllegalArgumentException("count must large then zero");
            }
            // Use the AQS methods to read, set, and modify synchronization state
            setState(count);
        }

        /**
         * Acquire the lock.
         * @param reduceCount
         * @return
         */
        @Override
        protected int tryAcquireShared(int reduceCount) {
            for (;;) {
	              int current = getState();
	              int newCount = current + returnCount;
                  // Use AQS CAS to set the value atomically
	              if (compareAndSetState(current, newCount)) {
	                  return true;
	              }
            }
        }

        /**
         * Release the lock.
         * @param returnCount
         * @return
         */
        @Override
        protected boolean tryReleaseShared(int returnCount) {
            for (;;) {
	              int current = getState();
	              int newCount = current - reduceCount;
	              if (newCount < 0 || compareAndSetState(current, newCount)) {
	                  return newCount;
	              }
            }
        }
    }

    private final Sync sync = new Sync(2);

    /**
     * Implement Lock’s lock and unlock methods.
     *
     * Call the AQS template methods acquire, acquireShared, release, and releaseShared from these methods.
     *
     * tryAcquireNanos and tryAcquireSharedNanos support timeouts.
     *
     * These methods ultimately call the overridden methods above.
     */
    @Override
    public void lock() {
        sync.acquireShared(1);
    }

    @Override
    public void unlock() {
        sync.releaseShared(1);
    }

    @Override
    public void lockInterruptibly() throws InterruptedException {

    }

    @Override
    public boolean tryLock() {
        return false;
    }

    @Override
    public boolean tryLock(final long time, final TimeUnit unit) throws InterruptedException {
        return false;
    }

    @Override
    public Condition newCondition() {
        return null;
    }


}

Reference Implementations

  1. Consult the AbstractQueuedSynchronizer documentation, which provides two worthwhile examples. They are much simpler and easier to understand than the implementation above.
  2. Study ReentrantLock and ReentrantReadWriteLock source code to learn how custom locks are implemented.

Notification Between Threads

Scenario

Two threads may need to communicate, meaning coordinate to complete a task. Suppose object A is accessed by multiple threads under certain conditions. When thread M changes state, thread N needs to print something. How can we implement this?

Use the thread wait/notification mechanism. Every Java object has wait(), wait(long timeout), notify(), and notifyAll(). Why? Every class inherits Object by default, and these methods belong to Object.

After N acquires A’s lock, it finds that flag does not satisfy the condition. It should release the lock so other threads can run. Otherwise N keeps holding it, flag may never satisfy the condition, and other threads cannot acquire A’s lock. Once N releases it, M can acquire it, modify flag, and release it. If N then acquires the lock again, flag now satisfies the condition, so N continues with its logic.

Possible approaches:

  1. N continually checks state and competes for the lock only when flag satisfies the condition. Otherwise it can sleep briefly.
  2. N waits for notification from M and competes for the lock after being notified.

Approach One

import java.util.concurrent.TimeUnit;
import lombok.SneakyThrows;
import org.junit.Test;

/**
 * Does not use waiting/notification; simply checks in a loop and sleeps briefly when the condition is not satisfied.
 *
 * @author jacksparrow414
 * @date 2020/10/20
 */
public final class ThreadWaitTest {
    static boolean flag = true;
    static Object lock = new Object();

    @Test
    @SneakyThrows
    public void assertThreadWait() {
        Thread threadN = new Thread(new ThreadN(), "threadN");
        Thread threadM = new Thread(new ThreadM(), "threadM");
        threadN.start();
        threadM.start();
        // Sleep for 10 seconds so the test output can finish
        TimeUnit.SECONDS.sleep(10);
    }

    static class ThreadN implements Runnable {

        @Override
        public void run() {
            while (flag) {
                try {
                    System.out.println(Thread.currentThread().getName() + "Condition not satisfied; sleeping before checking again");
                    Thread.sleep(1000);
                } catch (InterruptedException exception) {
                    exception.printStackTrace();
                }
            }
            synchronized (lock) {
                System.out.println(Thread.currentThread().getName() + "acquire lock");
            }
        }
    }

    static class ThreadM implements Runnable {
        @Override
        public void run() {
            while (flag) {
                synchronized (lock) {
                    try {
                        System.out.println(Thread.currentThread().getName() + "acquire lcok");
                        // Alternative sleeping method
                        TimeUnit.SECONDS.sleep(5);
                    } catch (InterruptedException exception) {
                        exception.printStackTrace();
                    }
                    flag = false;
                }
            }

            synchronized (lock) {
                System.out.println("acquire again");
            }
        }
    }
}

Test output:

threadNCondition not satisfied; sleeping before checking again threadM acquire lcok threadNCondition not satisfied; sleeping before checking again threadNCondition not satisfied; sleeping before checking again threadNCondition not satisfied; sleeping before checking again threadNCondition not satisfied; sleeping before checking again threadM acquire again threadN acquire lock

Initially, N sees that flag does not satisfy its condition, so it sleeps instead of competing for the lock. After sleeping, it checks again. M’s condition is initially satisfied, so M runs, modifies flag, and releases the object lock. N is still sleeping, so M acquires the lock again, finishes, and releases it. Once no other thread is competing and N wakes up, N successfully acquires the lock and executes.

Problems

  1. Threads competing for a lock should do so at the same time. Since N sleeps, it is still asleep when M first releases the lock and does not compete. N should immediately compete with other threads after release, but our code prevents that. With many competing threads, N sleeps for one second each time. If every release occurs during that sleep, N may have to wait until all others finish before it gets a chance. This clearly conflicts with the intention of multithreaded programming: submit tasks and let whichever thread acquires the lock execute, rather than forcing one to run last.
  2. Sleeping threads do not consume CPU resources. Since sleeping causes the problem, we could shorten it to 50 ms or 1 ms. But if tasks take time, frequent sleeping and checking will consume substantial CPU. We can never determine a precise ideal sleep duration, because each thread’s execution time is uncertain.

Approach Two

Use the thread wait/notification mechanism to solve these problems. First acquire the object lock. If the current thread’s condition is not satisfied, it waits and releases the lock. Another thread holding the lock then calls notifyAll or notify near the end of its synchronized block to notify all waiting threads or one of them. Notified threads compete for the lock; the winner leaves the wait state.

import java.util.concurrent.TimeUnit;
import lombok.SneakyThrows;
import org.junit.Test;

/**
 * @author jacksparrow414
 * @date 2020/10/20
 */
public final class ThreadWaitNotifyTest {
    static boolean flag = true;
    static Object lock = new Object();

    @Test
    @SneakyThrows
    public void assertThreadWaitNotify() {
        Thread threadN = new Thread(new threadN(), "threadN");
        Thread threadM = new Thread(new threadM(), "threadM");
        threadN.start();
        TimeUnit.SECONDS.sleep(1);
        threadM.start();
        TimeUnit.SECONDS.sleep(10);
    }

    static class threadN implements Runnable {

        @Override
        public void run() {
            synchronized (lock) {
                while (flag) {
                    System.out.println(Thread.currentThread().getName() + " flag is true wait");
                    try {
                        lock.wait();
                    } catch (InterruptedException exception) {
                        exception.printStackTrace();
                    }
                }

                System.out.println(Thread.currentThread().getName() + "falg is false running");
            }
        }
    }

    static class threadM implements  Runnable {

        @Override
        public void run() {
            synchronized (lock) {
                System.out.println(Thread.currentThread().getName() + " acquire lock");
                flag = false;

                try {
                    TimeUnit.SECONDS.sleep(3);
                    // Wake all waiting threads
                    lock.notifyAll();
                } catch (InterruptedException exception) {
                    exception.printStackTrace();
                }
            }
            // Notify waiting threads inside the synchronized block, or an error occurs
            //lock.notifyAll();
            synchronized (lock) {
                System.out.println(Thread.currentThread().getName() + " acquire again");
            }
        }

    }
}

Important Points

  1. Call wait, notifyAll, and notify inside a synchronized block, with the current thread already holding the lock. Calling them without the object lock causes an error.
Exception in thread "threadM" java.lang.IllegalMonitorStateException
 	at java.lang.Object.notifyAll(Native Method)
 	at com.example.mybatis.demomybatis.thread.ThreadWaitNotifyTest$threadM.run(ThreadWaitNotifyTest.java:60)
	at java.lang.Thread.run(Thread.java:748)
  1. After a lock-holding thread calls notify or notifyAll, waiting threads do not immediately return from wait. They only leave the wait state after acquiring the lock.
  2. After N is notified and reacquires the lock, it must still check the while condition and execute only when the condition holds. This is understandable: after obtaining the lock, it needs to recheck the thread’s condition rather than simply proceeding from the last point of execution. Flowchart of thread waiting, notification, lock reacquisition, and condition checks Sample code location

Share this post:

Continue this series

Locks in Java

  1. Locks in Java (Part 1): Lock Categories
  2. Locks in Java (Part 2): Implementing a Custom LockYou are here
  3. Locks in Java (Part 3): Implementing Wait and Notify with Condition
  4. Locks in Java (Part 4): High-Performance Web Page Caching with a Read-Write Lock

Comments

Questions, corrections, and experiences are welcome. Sign in with GitHub to comment; both language versions share this discussion.

Comments are available on the live site only.