Система координации операций
- Cинхронизация доступа к сущности между потоками Java
- Semaphore vs ReentrantLock для координации между потоками
- Предотвращение утечек памяти при блокировке сущностей в Java
- Контроль конкурентного доступа к динамически создаваемым сущностям в Java
- Стратегии блокировки на уровне сущности и глобального уровня в многопоточных Java-приложениях
- WeakHashMap для управления блокировками сущностей в Java
- Смена владельца у ReentrantLock
- Двухфазное управление конкурентным доступом с передачей управления (handoff)
- IllegalMonitorStateException
Рассмотрим задачу проектирования системы координации операций над сущностями в многопоточном Java приложении.
Требования и ограничения
- Каждая операция должна получать эксклюзивный доступ к сущности.
- Операция может начинаться в одном потоке, а заканчиваться в другом.
- Сущности могут добавляться и удаляться в большом количестве, это не должно приводить к утечкам локов / мьютексов / других объектов.
- Есть операции, которые требуют эксклюзивный доступ ко всем сущностям, даже еще не созданным (по сути – блокировка на создание сущностей).
- Есть bulk-операции над несколькими сущностями.
Шаг 1. Синхронизация по интернированной строке
Если услышать фразу “синхронизация доступа к сущностям по ID” – может захотеться использовать синхронизацию по интернированному ID сущности.
Вспомним, что string.intern() – операция интернирования строки. Такие строки будут дедуплицированны и помещены в String Pool.
В итоге два разных объекта String будут превращены в один, при условии, что у строк одинаковый контент:
String s1 = new String("abc").intern();
String s2 = new String("abc").intern();
System.out.println(s1 == s2); // true!
В теории можно синхронизироваться по таким интернированным строкам:
synchronized (entityName.intern()) {
doWork();
}
Но делать мы так конечно же не будем потому, что:
- В Java < 8 интернированные строки не чистятся вообще.
- В 8 <= Java <= 17 есть особенности с очисткой String Pool, часто это зависит от имплементации JVM.
Более того, если в JVM есть другая блокировка по интернированным строкам – такие блокировки будут пересекаться (адресы строк с одним и тем же контентом эквивалентны).
Шаг 2. Попытка использовать ReentrantLock
На ум сразу же приходит примитив синхронизации – ReentrantLock.
- Каждой сущности мы назначим свой lock.
- Перед выполнением операции – каждая из них должна захватить лок.
- После успешного выполнения или в случае ошибки – каждая из операций должна освободить лок.
Но проблема с ReentrantLock в том, что владелец лока (lock owner) не может быть изменен “на лету”.
Одно из требований заключается в том, что операция может начать выполняться в одном потоке, а закончится в другом.
Это значит, что lock() будет вызван в одном потоке, а unlock() в другом, на что JVM скажет “фи” и при попытке вызвать unlock() в другом потоке кинет в нас IllegalMonitorStateException с комментарием “thread doesn’t own the lock”.
Шаг 3. Заменяем Lock на Semaphore
У проблемы смены владельца лока есть решение – можно использовать Semaphore(1) вместо ReentrantLock.
ThreadA может получить пермит вызвав acquire(), а ThreadB может вернуть его обратно с помощью release().
Внимательный читатель скажет: wait, Semaphore(1) != Lock и будет прав.
Важно учесть, что семафор не обеспечивает согласованность пермитов сам. Рассмотрим такой сценарий:
- Мы создаем семафор с одним пермитом для некой сущности.
- ThreadA хочет поработать с сущностью, и захватывает один пермит.
- Состояние семафора – 0 пермитов.
- ThreadC по ошибке высвобождает 2 пермита.
- ThreadA высвобождает захваченный ранее пермит.
То есть в четвертом шаге никто не помешает “левому” ThreadC высвободить 2 пермита, хотя он их даже не получал. В итоге состояние семаформа после выполнения сценария выше – 3 доступных пермита.
Это значит, что задача обеспечения согласованности пермитов ложится на нас:
- Мы четко должны быть уверенны – сколько пермитов (в нашем случае – не больше одного) и у кого они находятся.
- Кто и когда и получает пермиты и высвобождает их.
Это не значит, что при использовании ReentrantLock можно не думать. Особенность в том, что если допущены ошибки в жизненном цикле блокировок:
- При ReentrantLock мы получаем Exception, сразу видим проблему и исправляем ее.
- При использовании Semaphore – система продолжит работать с замаскированной проблемой, которая стрельнет совсем в другом месте.
Итак, мы будем использовать Semaphore(1) для каждой сущности и будем предельно осторожным.
Шаг 4. Хранилище блокировок – new Semaphore[size]
Сущностей может быть много, поэтому нам необходимо хранилище блокировок.
Одно из возможных решений – хранить семафоры в массиве фиксированного размера. Идемпотентность реализуется через “шардирование”:
private final int STORAGE_SIZE = 5000;
Semaphore[] locks = new Semaphore[STORAGE_SIZE];
public synchronized Semaphore getLock(String entityName) {
return locks[entityName.hashCode() % STORAGE_SIZE];
}
Решение простое и в целом жизнеспособное, но есть большая проблема – коллизии. Если операций много – даже небольшой процент коллизий может сильно замедлить работу приложения.
Шаг 5. ConcurrentHashMap<String, Semaphore>
Что насчет ConcurrentHashMap<String, Semaphore>?
Когда мы создаем новую сущность – мы можем создать новый Semaphore и положить его в hash map.
Для защиты от race condition, когда две операции одновременно пытаются создать одну и ту же сущность, можно использовать concurrentMap.computeIfAbsent():
- Метод обеспечит критическую секцию на создание нового лока, что защитит от race conditions.
- Метод вернет уже существующий лок, если таковой имеется.
Помним, что computeIfAbsent() синхронизирует доступ только в рамках ConcurrentHashMap. При вызове computeIfAbsent() обычной HashMap – синхронизации нет.
Не будем забывать одно из требований: сущности могут создаваться и удаляться в большом количестве.
Это значит, что у нас утечка – если создать 10 миллионов сущностей, удалить 9 из них, затем создать еще 10 млн – то в хранилище с локами окажется 20 млн локов, а не 10 млн.
Шаг 6. Что насчет WeakHashMap?
На ум приходит WeakHashMap. Но что же должно быть ключом, а что – значением?
WeakHashMap<String, Semaphore>WeakHashMap<SemaphoreWrapper, Object>WeakHashMap<SemaphoreWrapper, SemaphoreWrapper>
где SemaphoreWrapper это пара Semaphore и ID сущности, с методами переопределенными hashCode и equals, основанными на ID.
Суть WeakHashMap в том, что в случае запуска GC – JVM почистит те записи, на ключи которых нет прямых ссылок в рамках остальной части программы.
WeakHashMap<String, Semaphore>
Предположим, ThreadA видит что карта пуская, и создает новый лок (псевдокод):
locks.put(entityName, new Semaphore(1));
Далее запускается другой ThreadB, и берет этот же лок по entityName:
locks.get(entityName);
Но кто будет держать strong reference на лок?
Ловушка в том, что пользователь может запустить любую операцию в любой момент времени, поэтому никто не держит strong reference на entityName.
entityName конструируется на лету. Это значит, что ThreadB хоть и получит тот же лок, что и ThreadA – JVM все равно может удалить его из карты в любой момент времени.
Это значит, что для внезапно появившегося ThreadC карта locks может быть уже пустой – тогда он создаст новый инстанс лока одновременно с тем, как ThreadB держит “старый” лок.
Ок, а что если интернировать идентификаторы? Для первого доступа:
String key = entityName.intern();
locks.put(key, new Semaphore(1));
Для последующих:
String key = entityName.intern();
Semaphore lock = locks.get(key);
Но как мы уже помним – интернировать миллионы строк это не самая хорошая идея.
WeakHashMap<SemaphoreWrapper, Object>
В такой конфигурации WeakHashMap – лок будет оставаться в хранилище, пока у какого-либо потока есть ссылка на объект EntityLock.
public synchronized SemaphoreWrapper getLock(String entityName) {
SemaphoreWrapper key = new SemaphoreWrapper(entityName);
return locks.putIfAbsent(probe, key -> new Object());
}
Но мы в ловушке – если внимательно приглядеться, есть несколько проблем, но основная из них заключается в том, что putIfAbsent() возвращает значение, а не ключ, то есть пустой Object, а не наш лок.
- Если всегда возвращать key – то пропадает смысл в хранилище – лок каждый раз будет новым.
- Нельзя “извлечь” оригинальный объект ключа из
WeakHashMap– можно только “сравниваться” с ним через equals() и hashCode().
WeakHashMap<SemaphoreWrapper, SemaphoreWrapper>
А что если положить лок не только в key, но и в и entry value? Есть два варианта.
Первый вариант:
public synchronized SemaphoreWrapper getLock(String entityName) {
SemaphoreWrapper probe = new SemaphoreWrapper(entityName);
return locks.computeIfAbsent(probe, en -> new SemaphoreWrapper(entityName));
}
Внимательный читатель снова заметит подвох – никто не держит strong reference на ключ (probe). А значит, он может быть удален в любой момент, даже если какой-то из потоков активно использует SemaphoreWrapper (тот, что в value).
Второй вариант:
public synchronized SemaphoreWrapper getLock(String entityName) {
SemaphoreWrapper probe = new SemaphoreWrapper(entityName);
return locks.computeIfAbsent(probe, en -> probe);
}
Но это тоже ловушка. Ключи WeakHashMap обернуты в WeakReference, но значения остаются strong reference.
Это значит, что если и ключ, и значение – один и тот же объект, то в JVM всегда будет одна strong referene на такой ключ – сама карта. В итоге такая запись никогда не будет удалена.
Шаг 7. ConcurrentHashMap<String, WeakReference>
А что если вернутся к ConcurrentHashMap, но в качестве значений использовать WeakReference?
Такой лок удалится JVM в том случае, если его никто не использует.
Это верно, но ключ в такой карте – ID сущности – останется, что все так же является утечкой.
Шаг 8. ConcurrentHashMap с очисткой “по удалению”
Хорошо, а что если вручную удалять локи при удалении сущности?
Если кто-то вызвал removeOperation(entityName), то можно сделать следующее:
ConcurrentHashMap<String, ConcurrentHashMap> locks;
public void removeOperation(String entityName) {
// Упрощенно
Semaphore lock = locks.get(entityName);
lock.acquire();
try {
removeEntityInternalOperation(entityName);
} finally {
lock.release();
locks.remove(entityName);
}
}
Кроме очевидных проблем есть еще и следующая: прямо в момент выполнения операции removeEntityInternalOperation может прилететь другая операция по этой же сущности!
Для этого мы и разрабатываем систему координации доступа – чтобы предотвратить такие случаи.
Это значит, что может быть такой сценарий:
[delete start] --------------> [delete finish]
[create start] --------------> [create finish]
Это значит, что операция delete удалит лок из хранилища прямо в то время, как он будет находиться у операции create. То есть для еще одной операции (третья по счету) – хранилище локов будет пустое, несмотря на то что операция create держит лок.
Шаг 9. ConcurrentHashMap с очисткой “по счетчику ссылок”
Мы приближаемся к ключевой концепции – подсчету ссылок.
Идея в том, что при получении лока операция должна “регистрироваться в журнале”, а при завершении – отменять регистрацию.
Схематично это будет выглядеть следующим образом.
Лок с поддержкой ownership change и подсчета ссылок:
class RefCountingSemaphore {
AtomicInteger refCount = new AtomicInteger(0);
Semaphore semaphore = new Semaphore(1);
String entityName;
}
Координатор блокировок:
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.decrementAndGet();
if (refs == 0) {
locks.remove(lock.entityName);
}
}
}
Как видно, каждая операция должна:
- “Зарегистрироваться” перед началом работы – то есть получить лок с помощью
getRef(). - “Отменить регистрацию” после выполнения – то есть вызвать
releaseRef()с ранее полученным локом.
Схематично это будет выглядеть так:
public void createEntity(String entityName) {
RefCountingSemaphore lock = lockManager.getRef(entityName);
lock.acquire();
createEntityStageOne(entityName);
// По условию – операция может заканчиваться в другом потоке
new Thread(() -> {
try {
createEntityStageTwo(entityName);
} finally {
lock.release();
lockManager.releaseRef(lock);
}
}).start();
}
Шаг 10. Механизм bulk-locking
Итак, на данный момент мы решили задачу синхронизации операций по динамическому набору сущностей в конфигурации “не больше 1 сущности на операцию”.
Но среди требований также есть пункт, что бывают bulk-операции, требующие блокировок по нескольким сущностям одновременно.
Например – bulkDelete(entityName1, entityName2).
Можно ли разрулить ситуацию следующим образом?
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();
}
// ...
}
Будет ли работать?
Внимательный читатель заметит – снова ловушка. Если запустить несколько операций bulkDelete в параллельных потоках в конфигурации вроде такой:
new Thread(() -> {
bulkDelete("entity1", "entity2");
}).start();
new Thread(() -> {
bulkDelete("entity2", "entity1");
}).start();
То можно очень быстро наткнуться на дедлок. ThreadA получит лок на entity1 и будет ждать лока на entity2. А ThreadB получит лок на entity2 и будет ждать лока на entity1.
Решение следующее: bulk-операции, требующие эксклюзивного доступа к нескольким сущностям, должны быть синхронизированны между собой. Такой компромисс.
Вводим дополнительный мьютекс для таких операций. При этом, он не отменяет необходимость получать индивидуальный локи:
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();
}
// ...
}
}
Координация bulk-операций готова.
Шаг 11. Глобальная блокировка
Осталось учесть последнее требование – мы должны поддерживать глобальные операции, которые требуют эксклюзивного доступа на все сущности, даже еще не созданные.
Уже разработанный нами механизм bulk-locking здесь не подходит, потому что может работать только с фиксированным набором сущностей. А глобальная блокировка подразумевает собой запрет даже на создание новых сущностей.
Ввести еще один мьютекс по типу globalMutex и поставить synchronized(globalMutex) в каждую обычную операцию и в каждую глобальную операцию?
Очевидно – это не вариант, так как в итоге уровень параллелизма будет равняться одному.
Здесь нам нужен ReadWriteLock. Обычные операции должны будут получать globalLock.readLock(), таким образом не блокируя друг друга, а вот глобальные операции должны получать globalLock.writeLock().
Не забудем выставить параметр fair в true:
final ReadWriteLock globalLock = new ReentrantReadWriteLock(true);
Пример глобальной операции:
public void resetSystem() {
globalLock.writeLock().lock();
try {
doGlobalWork();
} finally {
globalLock.writeLock().unlock();
}
}
Пример обычной операции:
public void createEntity(String entityName) {
RefCountingSemaphore lock = lockManager.getRef(entityName);
globalLock.readLock().lock(); // Сначала получаем глобальный лок
lock.acquire(); // Потом получаем обычный лок
// ...
}
Решено? Нет! Мы не учли, что операции могут заканчиваться в другом потоке:
public void createEntity(String entityName) {
RefCountingSemaphore lock = lockManager.getRef(entityName);
globalLock.readLock().lock();
lock.acquire();
// По условию – операция может заканчиваться в другом потоке
new Thread(() -> {
try {
createEntityStageTwo(entityName);
} finally {
lock.release();
lockManager.releaseRef(lock);
globalLock.readLock().unlock();
}
}).start();
}
Это значит, что при попытке выполнить unlock() для глобального лока JVM начнет кидаться в нас уже знакомым IllegalMonitorStateException.
Ок, скажем мы, плавали – знаем, можно заменить ReadWriteLock на Semaphore.
Или нет? В данном случае – нельзя; а все потому, что Semaphore не поддерживает нужное нам поведение: multiple readers, но single writer.
Как бы “передать” блокировку из одного потока в другой? Локи в Java так делать не имеют.
Одно из оригинальных решений можно описать как “two-phase concurrency control with handoff”.
Суть в том, что мы заводим глобальный счетчик операций:
final AtomicInteger runningOperations = new AtomicInteger(0);
Обычные операции все так же должны получать readLock(), но уже только для регистрации в журнале:
public void createEntity(String entityName) {
RefCountingSemaphore lock = lockManager.getRef(entityName);
globalLock.readLock().lock(); // Сначала получаем глобальный лок
runningOperations.incrementAndGet(); // Регистрируемся в журнале
globalLock.readLock().unlock(); // Отпускаем глобальный лок
lock.acquire(); // Получаем обычный лок
// По условию – операция может заканчиваться в другом потоке
new Thread(() -> {
try {
createEntityStageTwo(entityName);
} finally {
lock.release();
lockManager.releaseRef(lock); // Отмена регистрации в рамках сущности
runningOperations.decrementAndGet(); // Отмена глобальной регистрации
}
}).start();
}
Именно в момент, когда мы регистрируемся в журнале и отпускаем глобальный лок – по-сути происходит “передача” блокировки из рук лока (globalLock) в ведение счетчика (runningOperations).
Глобальные же операции, даже после получается writeLock() должны подождать, пока журнал окажется пуст:
public void resetSystem() {
globalLock.writeLock().lock();
try {
while (runningOperations.get() != 0) {
Thread.sleep(100);
}
doGlobalWork();
} finally {
globalLock.writeLock().unlock();
}
}
На этом все требования выполнены.
Итоговая система координации операций
Соберем все наработки воедино.
Необходимые объекты:
final Object bulkMutex = new Object();
final ReadWriteLock globalLock = new ReentrantReadWriteLock(true);
final AtomicInteger runningOperations = new AtomicInteger(0);
final LockManager lockManager = new LockManager();
Лок с поддержкой ownership change и подсчета ссылок:
class RefCountingSemaphore {
AtomicInteger refCount = new AtomicInteger(0);
Semaphore semaphore = new Semaphore(1);
String entityName;
}
Координатор блокировок:
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.decrementAndGet();
if (refs == 0) {
locks.remove(lock.entityName);
}
}
}
Пример обычной операции:
public void createEntity(String entityName) {
RefCountingSemaphore lock = lockManager.getRef(entityName);
globalLock.readLock().lock(); // Сначала получаем глобальный лок
runningOperations.incrementAndGet(); // Регистрируемся в журнале
globalLock.readLock().unlock(); // Отпускаем глобальный лок
lock.acquire(); // Получаем обычный лок
// По условию – операция может заканчиваться в другом потоке
new Thread(() -> {
try {
createEntityStageTwo(entityName);
} finally {
lock.release();
lockManager.releaseRef(lock); // Отмена регистрации в рамках сущности
runningOperations.decrementAndGet(); // Отмена глобальной регистрации
}
}).start();
}
Пример bulk-операции:
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();
}
// ...
}
}
Пример глобальной операции:
public void resetSystem() {
globalLock.writeLock().lock();
try {
while (runningOperations.get() != 0) {
Thread.sleep(100);
}
doGlobalWork();
} finally {
globalLock.writeLock().unlock();
}
}
На этом все.