Operation coordination system
- Java synchronize entity access across threads
- Java semaphore vs reentrantlock for multi-thread coordination
- Prevent memory leaks with entity locks in java
- Concurrent access control for dynamically created entities in java
- Entity-wide and global locking strategies in multithreaded java apps
- Weakhashmap pitfalls for java entity lock management
- Reentrantlock ownership change
- Two-phase concurrency control with handoff
- IllegalMonitorStateException
Let’s consider the task of designing a system for coordinating operations on entities in a multithreaded Java application.
Requirements and limitations
- Each operation must obtain exclusive access to an entity.
- An operation can start in one thread and end in another.
- Entities can be added and removed in large quantities; this should not lead to leaks of locks / mutexes / other objects.
- There are operations that require exclusive access to all entities, even those not yet created (essentially β a lock on entity creation).
- There are bulk operations on multiple entities.
Step 1. Synchronization by interned string
If we hear the phrase “synchronization of access to entities by ID” β we might want to use synchronization by the interned ID of the entity.
Let’s recall that string.intern() is a string interning operation. Such strings will be deduplicated and placed in the String Pool.
As a result, two different String objects will be turned into one, provided that the strings have the same content:
String s1 = new String("abc").intern();
String s2 = new String("abc").intern();
System.out.println(s1 == s2); // true!
In theory, we can synchronize on such interned strings:
synchronized (entityName.intern()) {
doWork();
}
But we, of course, will not do this because:
- In Java < 8, interned strings are not cleaned up at all.
- In 8 <= Java <= 17, there are peculiarities with String Pool cleanup, often depending on the JVM implementation.
Moreover, if there is another lock on interned strings in the JVM, such locks will intersect (addresses of strings with the same content are equivalent).
Step 2. Attempt to use ReentrantLock
The synchronization primitive β ReentrantLock β immediately comes to mind.
- We will assign its own lock to each entity.
- Before executing an operation, each of them must acquire the lock.
- After successful execution or in case of an error, each operation must release the lock.
But the problem with ReentrantLock is that the lock owner cannot be changed “on the fly”.
One of the requirements is that an operation can start executing in one thread and end in another.
This means that lock() will be called in one thread, and unlock() in another, to which the JVM will say no and upon attempting to call unlock() in another thread, will throw an IllegalMonitorStateException at us with the comment “thread doesn’t own the lock”.
Step 3. Replacing Lock with Semaphore
The problem of changing the lock owner has a solution β we can use Semaphore(1) instead of ReentrantLock.
ThreadA can obtain a permit by calling acquire(), and ThreadB can return it using release().
An attentive reader will say: wait, Semaphore(1) != Lock and will be right.
It is important to note that a semaphore does not ensure permit consistency by itself. Let’s consider such a scenario:
- We create a semaphore with one permit for a certain entity.
- ThreadA wants to work with the entity and acquires one permit.
- Semaphore state β 0 permits.
- ThreadC mistakenly releases 2 permits.
- ThreadA releases the previously acquired permit.
That is, in the fourth step, nothing will prevent a “rogue” ThreadC from releasing 2 permits, even though it never acquired them. As a result, the state of the semaphore after executing the scenario above is 3 available permits.
This means that the task of ensuring permit consistency falls on us:
- We must be clearly sure β how many permits (in our case β no more than one) and who holds them.
- Who and when acquires and releases permits.
This does not mean that when using ReentrantLock we don’t need to think. The peculiarity is that if errors are made in the lock lifecycle:
- With ReentrantLock, we get an Exception, immediately see the problem, and fix it.
- When using Semaphore β the system will continue to work with a masked problem that will blow up in a completely different place.
So, we will use Semaphore(1) for each entity and be extremely careful.
Step 4. Lock storage β new Semaphore[size]
There can be many entities, so we need a lock storage.
One possible solution is to store semaphores in a fixed-size array. Idempotency is achieved through “sharding”:
private final int STORAGE_SIZE = 5000;
Semaphore[] locks = new Semaphore[STORAGE_SIZE];
public synchronized Semaphore getLock(String entityName) {
return locks[entityName.hashCode() % STORAGE_SIZE];
}
The solution is simple and generally viable, but there is a big problem β collisions. If there are many operations, even a small percentage of collisions can significantly slow down the application.
Step 5. ConcurrentHashMap<String, Semaphore>
What about ConcurrentHashMap<String, Semaphore>?
When we create a new entity β we can create a new Semaphore and put it in the hash map.
To protect against a race condition where two operations simultaneously try to create the same entity, we can use concurrentMap.computeIfAbsent():
- The method will provide a critical section for creating a new lock, which will protect against race conditions.
- The method will return an already existing lock, if one exists.
Let’s remember that computeIfAbsent() synchronizes access only within ConcurrentHashMap. When calling computeIfAbsent() on a regular HashMap β there is no synchronization.
Let’s not forget one of the requirements: entities can be created and deleted in large quantities.
This means we have a leak β if we create 10 million entities, delete 9 of them, then create another 10 million β then the lock storage will contain 20 million locks, not 10 million.
Step 6. What about WeakHashMap?
WeakHashMap comes to mind. But what should be the key, and what β the value?
WeakHashMap<String, Semaphore>WeakHashMap<SemaphoreWrapper, Object>WeakHashMap<SemaphoreWrapper, SemaphoreWrapper>
where SemaphoreWrapper is a pair of Semaphore and entity ID, with overridden hashCode and equals methods based on ID.
The essence of WeakHashMap is that if GC is triggered β the JVM will clean up those entries whose keys have no direct references within the rest of the program.
WeakHashMap<String, Semaphore>
Suppose ThreadA sees that the map is empty and creates a new lock (pseudocode):
locks.put(entityName, new Semaphore(1));
Then another ThreadB starts and takes the same lock by entityName:
locks.get(entityName);
But who will hold a strong reference to the lock?
The trap is that a user can start any operation at any time, so no one holds a strong reference to entityName.
entityName is constructed on the fly. This means that even though ThreadB will get the same lock as ThreadA β the JVM can still remove it from the map at any time.
This means that for a suddenly appearing ThreadC, the locks map might already be empty β then it will create a new lock instance at the same time ThreadB holds the “old” lock.
Ok, what if we intern the identifiers? For the first access:
String key = entityName.intern();
locks.put(key, new Semaphore(1));
For subsequent ones:
String key = entityName.intern();
Semaphore lock = locks.get(key);
But as we already remember β interning millions of strings is not a very good idea.
WeakHashMap<SemaphoreWrapper, Object>
In such a WeakHashMap configuration β the lock will remain in storage as long as any thread has a reference to the EntityLock object.
public synchronized SemaphoreWrapper getLock(String entityName) {
SemaphoreWrapper key = new SemaphoreWrapper(entityName);
return locks.putIfAbsent(probe, key -> new Object());
}
But we are in a trap β if we look closely, there are several problems, but the main one is that putIfAbsent() returns the value, not the key, i.e., an empty Object, not our lock.
- If we always return key β then the storage loses its meaning β the lock will be new every time.
- It is not possible to “extract” the original key object from
WeakHashMapβ we can only “compare” with it via equals() and hashCode().
WeakHashMap<SemaphoreWrapper, SemaphoreWrapper>
What if we put the lock not only in the key, but also in the entry value? There are two options.
First option:
public synchronized SemaphoreWrapper getLock(String entityName) {
SemaphoreWrapper probe = new SemaphoreWrapper(entityName);
return locks.computeIfAbsent(probe, en -> new SemaphoreWrapper(entityName));
}
An attentive reader will again notice the catch β no one holds a strong reference to the key (probe). And therefore, it can be removed at any moment, even if one of the threads is actively using SemaphoreWrapper (the one in the value).
Second option:
public synchronized SemaphoreWrapper getLock(String entityName) {
SemaphoreWrapper probe = new SemaphoreWrapper(entityName);
return locks.computeIfAbsent(probe, en -> probe);
}
But this is also a trap. The keys of WeakHashMap are wrapped in WeakReference, but the values remain strong references.
This means that if both the key and the value are the same object, then there will always be one strong reference to such a key in the JVM β the map itself. As a result, such an entry will never be deleted.
Step 7. ConcurrentHashMap<String, WeakReference>
What if we return to ConcurrentHashMap, but use WeakReference as values?
Such a lock will be removed by the JVM if no one is using it.
This is true, but the key in such a map β the entity ID β will remain, which is still a leak.
Step 8. ConcurrentHashMap with “on delete” cleanup
Okay, what if we manually delete locks when an entity is deleted?
If someone called removeOperation(entityName), then we can do the following:
ConcurrentHashMap<String, ConcurrentHashMap> locks;
public void removeOperation(String entityName) {
// Simplified
Semaphore lock = locks.get(entityName);
lock.acquire();
try {
removeEntityInternalOperation(entityName);
} finally {
lock.release();
locks.remove(entityName);
}
}
Besides the obvious problems, there is also the following: right at the moment of executing the removeEntityInternalOperation operation, another operation on the same entity can arrive!
This is why we are developing the access coordination system β to prevent such cases.
This means there can be such a scenario:
[delete start] --------------> [delete finish]
[create start] --------------> [create finish]
This means that the delete operation will remove the lock from storage right while it is held by the create operation. That is, for another operation (the third one in sequence) β the lock storage will be empty, even though the create operation holds the lock.
Step 9. ConcurrentHashMap with “by reference counter” cleanup
We are approaching the key concept β reference counting.
The idea is that when obtaining a lock, an operation must “register in a log”, and upon completion β cancel the registration.
Schematically, it will look like this.
Lock with support for ownership change and reference counting:
class RefCountingSemaphore {
AtomicInteger refCount = new AtomicInteger(0);
Semaphore semaphore = new Semaphore(1);
String entityName;
}
Lock coordinator:
class LockManager {
ConcurrentHashMap<String, RefCountingSemaphore> locks;
private synchronized RefCountingSemaphore getRef(String entityName) {
RefCountingSemaphore lock = locks.get(entityName);
if (lock == null) {
lock = new RefCountingSemaphore(entityName);
}
// Registration in the entity log
lock.refCount.incrementAndGet();
return lock;
}
private synchronized void releaseRef(RefCountingSemaphore lock) {
// Cancellation of registration in the entity log
int refs = lock.refCount.decrementAndGet();
if (refs == 0) {
locks.remove(lock.entityName);
}
}
}
As can be seen, each operation must:
- “Register” before starting work β i.e., obtain the lock using
getRef(). - “Cancel registration” after completion β i.e., call
releaseRef()with the previously obtained lock.
Schematically, it will look like this:
public void createEntity(String entityName) {
RefCountingSemaphore lock = lockManager.getRef(entityName);
lock.acquire();
createEntityStageOne(entityName);
// As per condition β operation can end in another thread
new Thread(() -> {
try {
createEntityStageTwo(entityName);
} finally {
lock.release();
lockManager.releaseRef(lock);
}
}).start();
}
Step 10. Bulk-locking mechanism
So, at this point, we have solved the problem of synchronizing operations on a dynamic set of entities in the configuration “no more than 1 entity per operation”.
But among the requirements, there is also a point that there are bulk operations requiring locks on several entities simultaneously.
For example β bulkDelete(entityName1, entityName2).
Can we resolve the situation as follows?
public void bulkDelete(List<String> entityNames) {
List<RefCountingSemaphore> locks = new ArrayList<>();
for (String entityName : entityNames) {
RefCountingSemaphore lock = lockManager.getRef(entityName);
locks.add(lock);
lock.lock();
}
// ...
}
Will it work?
An attentive reader will notice β another trap. If we run several bulkDelete operations in parallel threads in a configuration like this:
new Thread(() -> {
bulkDelete("entity1", "entity2");
}).start();
new Thread(() -> {
bulkDelete("entity2", "entity1");
}).start();
Then we can very quickly run into a deadlock. ThreadA will get the lock on entity1 and wait for the lock on entity2. And ThreadB will get the lock on entity2 and wait for the lock on entity1.
The solution is as follows: bulk operations requiring exclusive access to multiple entities must be synchronized among themselves. Such a compromise.
We introduce an additional mutex for such operations. However, it does not negate the need to obtain individual locks:
final Object bulkMutex = new Object();
public void bulkDelete(List<String> entityNames) {
synchronized (bulkMutex) {
List<RefCountingSemaphore> locks = new ArrayList<>();
for (String entityName : entityNames) {
RefCountingSemaphore lock = lockManager.getRef(entityName);
locks.add(lock);
lock.lock();
}
// ...
}
}
Coordination of bulk operations is ready.
Step 11. Global lock
It remains to consider the last requirement β we must support global operations that require exclusive access to all entities, even those not yet created.
The bulk-locking mechanism we have already developed is not suitable here because it can only work with a fixed set of entities. And a global lock implies a ban even on the creation of new entities.
Introduce another mutex like globalMutex and put synchronized(globalMutex) in every regular operation and every global operation?
Obviously β this is not an option, as the level of parallelism will ultimately be one.
Here we need a ReadWriteLock. Regular operations will have to obtain globalLock.readLock(), thus not blocking each other, while global operations must obtain globalLock.writeLock().
Let’s not forget to set the fair parameter to true:
final ReadWriteLock globalLock = new ReentrantReadWriteLock(true);
Example of a global operation:
public void resetSystem() {
globalLock.writeLock().lock();
try {
doGlobalWork();
} finally {
globalLock.writeLock().unlock();
}
}
Example of a regular operation:
public void createEntity(String entityName) {
RefCountingSemaphore lock = lockManager.getRef(entityName);
globalLock.readLock().lock(); // First, we obtain the global lock
lock.acquire(); // Then we obtain the regular lock
// ...
}
Solved? No! We did not consider that operations can end in another thread:
public void createEntity(String entityName) {
RefCountingSemaphore lock = lockManager.getRef(entityName);
globalLock.readLock().lock();
lock.acquire();
// As per condition β operation can end in another thread
new Thread(() -> {
try {
createEntityStageTwo(entityName);
} finally {
lock.release();
lockManager.releaseRef(lock);
globalLock.readLock().unlock();
}
}).start();
}
This means that when attempting to execute unlock() for the global lock, the JVM will start throwing the already familiar IllegalMonitorStateException at us.
Ok, we’ll say, been there β done that, we can replace ReadWriteLock with Semaphore.
Or not? In this case β no; and all because Semaphore does not support the behavior we need: multiple readers, but a single writer.
How to “transfer” a lock from one thread to another? Locks in Java are not designed to do this.
One of the original solutions can be described as “two-phase concurrency control with handoff”.
The point is that we introduce a global operation counter:
final AtomicInteger runningOperations = new AtomicInteger(0);
Regular operations still need to obtain readLock(), but now only for registration in the log:
public void createEntity(String entityName) {
RefCountingSemaphore lock = lockManager.getRef(entityName);
globalLock.readLock().lock(); // First, we obtain the global lock
runningOperations.incrementAndGet(); // We register in the log
globalLock.readLock().unlock(); // We release the global lock
lock.acquire(); // We obtain the regular lock
// As per condition β operation can end in another thread
new Thread(() -> {
try {
createEntityStageTwo(entityName);
} finally {
lock.release();
lockManager.releaseRef(lock); // Cancellation of registration within the entity
runningOperations.decrementAndGet(); // Cancellation of global registration
}
}).start();
}
It is precisely at the moment when we register in the log and release the global lock that, in essence, a “transfer” of the lock occurs from the hands of the lock (globalLock) to the management of the counter (runningOperations).
Global operations, however, even after obtaining writeLock(), must wait until the log becomes empty:
public void resetSystem() {
globalLock.writeLock().lock();
try {
while (runningOperations.get() != 0) {
Thread.sleep(100);
}
doGlobalWork();
} finally {
globalLock.writeLock().unlock();
}
}
With this, all requirements are met.
Final operation coordination system
Let’s put all the developments together.
Necessary objects:
final Object bulkMutex = new Object();
final ReadWriteLock globalLock = new ReentrantReadWriteLock(true);
final AtomicInteger runningOperations = new AtomicInteger(0);
final LockManager lockManager = new LockManager();
Lock with support for ownership change and reference counting:
class RefCountingSemaphore {
AtomicInteger refCount = new AtomicInteger(0);
Semaphore semaphore = new Semaphore(1);
String entityName;
}
Lock coordinator:
class LockManager {
ConcurrentHashMap<String, RefCountingSemaphore> locks;
private synchronized RefCountingSemaphore getRef(String entityName) {
RefCountingSemaphore lock = locks.get(entityName);
if (lock == null) {
lock = new RefCountingSemaphore(entityName);
}
lock.refCount.incrementAndGet();
return lock;
}
private synchronized void releaseRef(RefCountingSemaphore lock) {
int refs = lock.refCount.decrementAndGet();
if (refs == 0) {
locks.remove(lock.entityName);
}
}
}
Example of a regular operation:
public void createEntity(String entityName) {
RefCountingSemaphore lock = lockManager.getRef(entityName);
globalLock.readLock().lock(); // First, we obtain the global lock
runningOperations.incrementAndGet(); // We register in the log
globalLock.readLock().unlock(); // We release the global lock
lock.acquire(); // We obtain the regular lock
// As per condition β operation can end in another thread
new Thread(() -> {
try {
createEntityStageTwo(entityName);
} finally {
lock.release();
lockManager.releaseRef(lock); // Cancellation of registration within the entity
runningOperations.decrementAndGet(); // Cancellation of global registration
}
}).start();
}
Example of a bulk operation:
public void bulkDelete(List<String> entityNames) {
synchronized (bulkMutex) {
List<RefCountingSemaphore> locks = new ArrayList<>();
for (String entityName : entityNames) {
RefCountingSemaphore lock = lockManager.getRef(entityName);
locks.add(lock);
lock.lock();
}
// ...
}
}
Example of a global operation:
public void resetSystem() {
globalLock.writeLock().lock();
try {
while (runningOperations.get() != 0) {
Thread.sleep(100);
}
doGlobalWork();
} finally {
globalLock.writeLock().unlock();
}
}
That’s all.