Система координации операций


Что содержится в статье ⌄
  • 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 и будет прав.

Важно учесть, что семафор не обеспечивает согласованность пермитов сам. Рассмотрим такой сценарий:

  1. Мы создаем семафор с одним пермитом для некой сущности.
  2. ThreadA хочет поработать с сущностью, и захватывает один пермит.
  3. Состояние семафора – 0 пермитов.
  4. ThreadC по ошибке высвобождает 2 пермита.
  5. 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();
    }
}

На этом все.