Sistema de coordinación de operaciones
- Sincronización de acceso a entidades en java a través de hilos
- Java semaphore vs reentrantlock para coordinación multi-hilo
- Prevenir fugas de memoria con bloqueos de entidades en java
- Control de acceso concurrente para entidades creadas dinámicamente en java
- Estrategias de bloqueo global y a nivel de entidad en aplicaciones java multihilo
- Problemas de weakhashmap para la gestión de bloqueos de entidades en java
- Cambio de propiedad de reentrantlock
- Control de concurrencia en dos fases con transferencia
- IllegalMonitorStateException
Consideremos la tarea de diseñar un sistema de coordinación de operaciones sobre entidades en una aplicación Java multihilo.
Requisitos y restricciones
- Cada operación debe obtener acceso exclusivo a la entidad.
- La operación puede comenzar en un hilo y terminar en otro.
- Las entidades pueden agregarse y eliminarse en gran cantidad, esto no debe conducir a fugas de bloqueos / mutex / otros objetos.
- Hay operaciones que requieren acceso exclusivo a todas las entidades, incluso a las que aún no se han creado (en esencia, bloqueo en la creación de entidades).
- Hay operaciones masivas sobre varias entidades.
Paso 1. Sincronización por string internada
Si escuchamos la frase “sincronización de acceso a entidades por ID”, podríamos querer usar la sincronización por ID de entidad internado.
Recordemos que string.intern() es una operación de internación de strings. Tales strings serán deduplicadas y colocadas en el String Pool.
Como resultado, dos objetos String diferentes se convertirán en uno, siempre que strings tengan el mismo contenido:
String s1 = new String("abc").intern();
String s2 = new String("abc").intern();
System.out.println(s1 == s2); // true!
En teoría, podemos sincronizarnos con estas strings internadas:
synchronized (entityName.intern()) {
doWork();
}
Pero por supuesto no haremos esto porque:
- En Java < 8 strings internadas no se limpian en absoluto.
- En 8 <= Java <= 17 hay peculiaridades con la limpieza del String Pool, a menudo depende de la implementación de JVM.
Además, si hay otro bloqueo en la JVM por strings internadas, tales bloqueos se cruzarán (las direcciones de strings con el mismo contenido son equivalentes).
Paso 2. Intento de usar ReentrantLock
Lo primero que viene a la mente es la primitiva de sincronización – ReentrantLock.
- Asignaremos a cada entidad su propio bloqueo.
- Antes de realizar la operación, cada una debe adquirir el bloqueo.
- Después de una ejecución exitosa o en caso de error, cada operación debe liberar el bloqueo.
Pero el problema con ReentrantLock es que el propietario del bloqueo (lock owner) no puede cambiarse “sobre la marcha”.
Uno de los requisitos es que la operación puede comenzar a ejecutarse en un hilo y terminar en otro.
Esto significa que lock() se llamará en un hilo y unlock() en otro, a lo que JVM dirá “no” y al intentar llamar a unlock() en otro hilo lanzará IllegalMonitorStateException con el comentario “thread doesn’t own the lock”.
Paso 3. Reemplazamos Lock por Semaphore
Para el problema del cambio de propietario del bloqueo hay una solución – podemos usar Semaphore(1) en lugar de ReentrantLock.
ThreadA puede obtener un permiso llamando a acquire(), y ThreadB puede devolverlo con release().
Un lector atento dirá: wait, Semaphore(1) != Lock y tendrá razón.
Es importante tener en cuenta que el semáforo no garantiza la consistencia de los permisos por sí mismo. Consideremos este escenario:
- Creamos un semáforo con un permiso para alguna entidad.
- ThreadA quiere trabajar con la entidad y adquiere un permiso.
- Estado del semáforo – 0 permisos.
- ThreadC por error libera 2 permisos.
- ThreadA libera el permiso adquirido anteriormente.
Es decir, en el cuarto paso nada impide que el “externo” ThreadC libere 2 permisos, aunque ni siquiera los obtuvo. Como resultado, el estado del semáforo después de ejecutar el escenario anterior – 3 permisos disponibles.
Esto significa que la tarea de garantizar la consistencia de los permisos recae en nosotros:
- Debemos estar absolutamente seguros – cuántos permisos (en nuestro caso – no más de uno) y quién los tiene.
- Quién, cuándo obtiene permisos y los libera.
Esto no significa que al usar ReentrantLock no haya que pensar. La peculiaridad es que si se cometen errores en el ciclo de vida de los bloqueos:
- Con ReentrantLock obtenemos una Exception, vemos el problema inmediatamente y lo corregimos.
- Al usar Semaphore – el sistema continuará funcionando con un problema enmascarado que se manifestará en un lugar completamente diferente.
Así que usaremos Semaphore(1) para cada entidad y seremos extremadamente cuidadosos.
Paso 4. Almacén de bloqueos – new Semaphore[size]
Puede haber muchas entidades, por lo que necesitamos un almacén de bloqueos.
Una de las posibles soluciones es almacenar semáforos en un array de tamaño fijo. La idempotencia se implementa a través de “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];
}
La solución es simple y generalmente viable, pero hay un gran problema – colisiones. Si hay muchas operaciones, incluso un pequeño porcentaje de colisiones puede ralentizar significativamente el funcionamiento de la aplicación.
Paso 5. ConcurrentHashMap<String, Semaphore>
Qué hay de ConcurrentHashMap<String, Semaphore>?
Cuando creamos una nueva entidad, podemos crear un nuevo Semaphore y ponerlo en el hash map.
Para proteger contra condiciones de carrera, cuando dos operaciones intentan crear la misma entidad simultáneamente, podemos usar concurrentMap.computeIfAbsent():
- El método proporcionará una sección crítica para la creación de un nuevo bloqueo, lo que protegerá contra condiciones de carrera.
- El método devolverá un bloqueo ya existente, si lo hay.
Recordemos que computeIfAbsent() sincroniza el acceso solo dentro de ConcurrentHashMap. Al llamar a computeIfAbsent() con un HashMap normal, no hay sincronización.
No olvidemos uno de los requisitos: las entidades pueden crearse y eliminarse en gran cantidad.
Esto significa que tenemos una fuga – si creamos 10 millones de entidades, eliminamos 9 de ellas, luego creamos otros 10 millones – en el almacén con bloqueos habrá 20 millones de bloqueos, no 10 millones.
Paso 6. Qué hay de WeakHashMap?
Viene a la mente WeakHashMap. Pero qué debe ser la clave y qué el valor?
WeakHashMap<String, Semaphore>WeakHashMap<SemaphoreWrapper, Object>WeakHashMap<SemaphoreWrapper, SemaphoreWrapper>
donde SemaphoreWrapper es un par de Semaphore e ID de entidad, con métodos sobreescritos hashCode y equals, basados en ID.
La esencia de WeakHashMap es que en caso de ejecución de GC, JVM limpiará aquellos registros cuyas claves no tengan referencias directas en el resto del programa.
WeakHashMap<String, Semaphore>
Supongamos que ThreadA ve que el mapa está vacío y crea un nuevo bloqueo (pseudocódigo):
locks.put(entityName, new Semaphore(1));
Luego se inicia otro ThreadB y toma este mismo bloqueo por entityName:
locks.get(entityName);
Pero quién mantendrá la strong reference al bloqueo?
La trampa es que el usuario puede iniciar cualquier operación en cualquier momento, por lo que nadie mantiene una strong reference a entityName.
entityName se construye sobre la marcha. Esto significa que aunque ThreadB obtenga el mismo bloqueo que ThreadA, JVM aún puede eliminarlo del mapa en cualquier momento.
Esto significa que para un ThreadC que aparece repentinamente, el mapa locks ya puede estar vacío – entonces creará una nueva instancia de bloqueo mientras ThreadB mantiene el bloqueo “viejo”.
Bien, y si internamos los identificadores? Para el primer acceso:
String key = entityName.intern();
locks.put(key, new Semaphore(1));
Para los siguientes:
String key = entityName.intern();
Semaphore lock = locks.get(key);
Pero como ya recordamos, internar millones de strings no es la mejor idea.
WeakHashMap<SemaphoreWrapper, Object>
En tal configuración de WeakHashMap, el bloqueo permanecerá en el almacén mientras algún hilo tenga una referencia al objeto EntityLock.
public synchronized SemaphoreWrapper getLock(String entityName) {
SemaphoreWrapper key = new SemaphoreWrapper(entityName);
return locks.putIfAbsent(probe, key -> new Object());
}
Pero estamos atrapados – si miramos atentamente, hay varios problemas, pero el principal es que putIfAbsent() devuelve el valor, no la clave, es decir, un Object vacío, no nuestro bloqueo.
- Si siempre devolvemos key, desaparece el sentido del almacén – el bloqueo será nuevo cada vez.
- No se puede “extraer” el objeto original de la clave de
WeakHashMap– solo se puede “comparar” con él a través de equals() y hashCode().
WeakHashMap<SemaphoreWrapper, SemaphoreWrapper>
Y si ponemos el bloqueo no solo en key, sino también en entry value? Hay dos opciones.
Primera opción:
public synchronized SemaphoreWrapper getLock(String entityName) {
SemaphoreWrapper probe = new SemaphoreWrapper(entityName);
return locks.computeIfAbsent(probe, en -> new SemaphoreWrapper(entityName));
}
Un lector atento notará de nuevo la trampa – nadie mantiene una strong reference a la clave (probe). Lo que significa que puede eliminarse en cualquier momento, incluso si alguno de los hilos está usando activamente SemaphoreWrapper (el que está en value).
Segunda opción:
public synchronized SemaphoreWrapper getLock(String entityName) {
SemaphoreWrapper probe = new SemaphoreWrapper(entityName);
return locks.computeIfAbsent(probe, en -> probe);
}
Pero esto también es una trampa. Las claves de WeakHashMap están envueltas en WeakReference, pero los valores siguen siendo strong reference.
Esto significa que si tanto la clave como el valor son el mismo objeto, entonces en JVM siempre habrá una strong reference a dicha clave – el propio mapa. Como resultado, tal registro nunca será eliminado.
Paso 7. ConcurrentHashMap<String, WeakReference>
Y si volvemos a ConcurrentHashMap, pero usamos WeakReference como valores?
Tal bloqueo será eliminado por JVM en caso de que nadie lo use.
Esto es cierto, pero la clave en tal mapa – ID de entidad – permanecerá, lo que sigue siendo una fuga.
Paso 8. ConcurrentHashMap con limpieza “por eliminación”
Bien, y si eliminamos manualmente los bloqueos al eliminar la entidad?
Si alguien llamó a removeOperation(entityName), podemos hacer lo siguiente:
ConcurrentHashMap<String, ConcurrentHashMap> locks;
public void removeOperation(String entityName) {
// Simplificado
Semaphore lock = locks.get(entityName);
lock.acquire();
try {
removeEntityInternalOperation(entityName);
} finally {
lock.release();
locks.remove(entityName);
}
}
Además de los problemas obvios, hay otro: justo en el momento de ejecutar la operación removeEntityInternalOperation puede llegar otra operación para la misma entidad!
Para esto estamos desarrollando el sistema de coordinación de acceso – para prevenir tales casos.
Esto significa que puede haber un escenario como este:
[delete start] --------------> [delete finish]
[create start] --------------> [create finish]
Esto significa que la operación delete eliminará el bloqueo del almacén justo cuando está siendo utilizado por la operación create. Es decir, para otra operación (la tercera en orden) – el almacén de bloqueos estará vacío, a pesar de que la operación create mantiene el bloqueo.
Paso 9. ConcurrentHashMap con limpieza “por contador de referencias”
Nos acercamos al concepto clave – el conteo de referencias.
La idea es que al obtener un bloqueo, la operación debe “registrarse en el diario”, y al finalizar – cancelar el registro.
Esquemáticamente esto se verá así.
Bloqueo con soporte para cambio de propiedad y conteo de referencias:
class RefCountingSemaphore {
AtomicInteger refCount = new AtomicInteger(0);
Semaphore semaphore = new Semaphore(1);
String entityName;
}
Coordinador de bloqueos:
class LockManager {
ConcurrentHashMap<String, RefCountingSemaphore> locks;
private synchronized RefCountingSemaphore getRef(String entityName) {
RefCountingSemaphore lock = locks.get(entityName);
if (lock == null) {
lock = new RefCountingSemaphore(entityName);
}
// Registro en el diario de la entidad
lock.refCount.incrementAndGet();
return lock;
}
private synchronized void releaseRef(RefCountingSemaphore lock) {
// Cancelación del registro en el diario de la entidad
int refs = lock.decrementAndGet();
if (refs == 0) {
locks.remove(lock.entityName);
}
}
}
Como se puede ver, cada operación debe:
- “Registrarse” antes de comenzar a trabajar – es decir, obtener un bloqueo usando
getRef(). - “Cancelar el registro” después de la ejecución – es decir, llamar a
releaseRef()con el bloqueo obtenido anteriormente.
Esquemáticamente se verá así:
public void createEntity(String entityName) {
RefCountingSemaphore lock = lockManager.getRef(entityName);
lock.acquire();
createEntityStageOne(entityName);
// Según la condición – la operación puede terminar en otro hilo
new Thread(() -> {
try {
createEntityStageTwo(entityName);
} finally {
lock.release();
lockManager.releaseRef(lock);
}
}).start();
}
Paso 10. Mecanismo de bulk-locking
Así, en este momento hemos resuelto la tarea de sincronización de operaciones por un conjunto dinámico de entidades en la configuración “no más de 1 entidad por operación”.
Pero entre los requisitos también hay un punto de que existen operaciones bulk, que requieren bloqueos en varias entidades simultáneamente.
Por ejemplo – bulkDelete(entityName1, entityName2).
Se puede resolver la situación de la siguiente manera?
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();
}
// ...
}
Funcionará?
Un lector atento notará – otra trampa. Si ejecutamos varias operaciones bulkDelete en hilos paralelos en una configuración como esta:
new Thread(() -> {
bulkDelete("entity1", "entity2");
}).start();
new Thread(() -> {
bulkDelete("entity2", "entity1");
}).start();
Entonces podemos encontrarnos rápidamente con un deadlock. ThreadA obtendrá el bloqueo de entity1 y esperará el bloqueo de entity2. Y ThreadB obtendrá el bloqueo de entity2 y esperará el bloqueo de entity1.
La solución es la siguiente: las operaciones bulk, que requieren acceso exclusivo a varias entidades, deben estar sincronizadas entre sí. Un compromiso.
Introducimos un mutex adicional para tales operaciones. Sin embargo, no elimina la necesidad de obtener bloqueos individuales:
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();
}
// ...
}
}
La coordinación de operaciones bulk está lista.
Paso 11. Bloqueo global
Queda considerar el último requisito – debemos soportar operaciones globales, que requieren acceso exclusivo a todas las entidades, incluso a las que aún no se han creado.
El mecanismo de bulk-locking que ya hemos desarrollado no es adecuado aquí, porque solo puede trabajar con un conjunto fijo de entidades. Y el bloqueo global implica la prohibición incluso de crear nuevas entidades.
Introducir otro mutex como globalMutex y poner synchronized(globalMutex) en cada operación normal y en cada operación global?
Obviamente, esta no es una opción, ya que el nivel de paralelismo será igual a uno.
Aquí necesitamos un ReadWriteLock. Las operaciones normales deberán obtener globalLock.readLock(), sin bloquearse entre sí, mientras que las operaciones globales deberán obtener globalLock.writeLock().
No olvidemos establecer el parámetro fair en true:
final ReadWriteLock globalLock = new ReentrantReadWriteLock(true);
Ejemplo de operación global:
public void resetSystem() {
globalLock.writeLock().lock();
try {
doGlobalWork();
} finally {
globalLock.writeLock().unlock();
}
}
Ejemplo de operación normal:
public void createEntity(String entityName) {
RefCountingSemaphore lock = lockManager.getRef(entityName);
globalLock.readLock().lock(); // Primero obtenemos el bloqueo global
lock.acquire(); // Luego obtenemos el bloqueo normal
// ...
}
Resuelto? No! No hemos tenido en cuenta que las operaciones pueden terminar en otro hilo:
public void createEntity(String entityName) {
RefCountingSemaphore lock = lockManager.getRef(entityName);
globalLock.readLock().lock();
lock.acquire();
// Según la condición – la operación puede terminar en otro hilo
new Thread(() -> {
try {
createEntityStageTwo(entityName);
} finally {
lock.release();
lockManager.releaseRef(lock);
globalLock.readLock().unlock();
}
}).start();
}
Esto significa que al intentar ejecutar unlock() para el bloqueo global, JVM comenzará a lanzarnos el ya familiar IllegalMonitorStateException.
Ok, diremos, hemos navegado por esto, podemos reemplazar ReadWriteLock con Semaphore.
O no? En este caso – no; porque Semaphore no soporta el comportamiento que necesitamos: multiple readers, pero single writer.
Cómo “transferir” el bloqueo de un hilo a otro? Los bloqueos en Java no pueden hacerlo.
Una de las soluciones originales se puede describir como “control de concurrencia en dos fases con transferencia”.
Se trata de que creamos un contador global de operaciones:
final AtomicInteger runningOperations = new AtomicInteger(0);
Las operaciones normales siguen obteniendo readLock(), pero ahora solo para el registro en el diario:
public void createEntity(String entityName) {
RefCountingSemaphore lock = lockManager.getRef(entityName);
globalLock.readLock().lock(); // Primero obtenemos el bloqueo global
runningOperations.incrementAndGet(); // Nos registramos en el diario
globalLock.readLock().unlock(); // Liberamos el bloqueo global
lock.acquire(); // Obtenemos el bloqueo normal
// Según la condición – la operación puede terminar en otro hilo
new Thread(() -> {
try {
createEntityStageTwo(entityName);
} finally {
lock.release();
lockManager.releaseRef(lock); // Cancelación del registro en el marco de la entidad
runningOperations.decrementAndGet(); // Cancelación del registro global
}
}).start();
}
Precisamente en el momento en que nos registramos en el diario y liberamos el bloqueo global – en esencia ocurre la “transferencia” del bloqueo de las manos del bloqueo (globalLock) al cuidado del contador (runningOperations).
Las operaciones globales, incluso después de obtener writeLock(), deben esperar hasta que el diario esté vacío:
public void resetSystem() {
globalLock.writeLock().lock();
try {
while (runningOperations.get() != 0) {
Thread.sleep(100);
}
doGlobalWork();
} finally {
globalLock.writeLock().unlock();
}
}
Con esto se cumplen todos los requisitos.
Sistema final de coordinación de operaciones
Reunamos todos los desarrollos en uno.
Objetos necesarios:
final Object bulkMutex = new Object();
final ReadWriteLock globalLock = new ReentrantReadWriteLock(true);
final AtomicInteger runningOperations = new AtomicInteger(0);
final LockManager lockManager = new LockManager();
Bloqueo con soporte para cambio de propiedad y conteo de referencias:
class RefCountingSemaphore {
AtomicInteger refCount = new AtomicInteger(0);
Semaphore semaphore = new Semaphore(1);
String entityName;
}
Coordinador de bloqueos:
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);
}
}
}
Ejemplo de operación normal:
public void createEntity(String entityName) {
RefCountingSemaphore lock = lockManager.getRef(entityName);
globalLock.readLock().lock(); // Primero obtenemos el bloqueo global
runningOperations.incrementAndGet(); // Nos registramos en el diario
globalLock.readLock().unlock(); // Liberamos el bloqueo global
lock.acquire(); // Obtenemos el bloqueo normal
// Según la condición – la operación puede terminar en otro hilo
new Thread(() -> {
try {
createEntityStageTwo(entityName);
} finally {
lock.release();
lockManager.releaseRef(lock); // Cancelación del registro en el marco de la entidad
runningOperations.decrementAndGet(); // Cancelación del registro global
}
}).start();
}
Ejemplo de operación 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();
}
// ...
}
}
Ejemplo de operación global:
public void resetSystem() {
globalLock.writeLock().lock();
try {
while (runningOperations.get() != 0) {
Thread.sleep(100);
}
doGlobalWork();
} finally {
globalLock.writeLock().unlock();
}
}
Eso es todo.