CountDownLatch
- CountDownLatch permite que uno o más threads esperen a que otros threads completen sus operaciones.
- CountDownLatch puede reemplazar la función de join, ofreciendo funcionalidades más avanzadas.
- El método countDown de CountDownLatch decrementa N en 1; el método await bloqueará el thread actual hasta que N se convierta en cero.
- No es posible reinicializar ni modificar el valor del contador interno de un objeto CountDownLatch después de su creación.
- CountDownLatch se implementa internamente usando el lock compartido AQS.
public class SincronizacionFinal {
private static final CountDownLatch LATCH_FINAL = new CountDownLatch(2);
public static void main(String[] args) throws InterruptedException {
new Thread(() -> {
System.out.println(1);
LATCH_FINAL.countDown();
System.out.println(2);
LATCH_FINAL.countDown();
}).start();
LATCH_FINAL.await();
System.out.println("3");
}
}
CyclicBarrier
- CyclicBarrier establece una barrera (también llamada punto de sincronización) que bloquea un grupo de threads hasta que el último thread alcance la barrera, momento en el cual se permite continuar la ejecución de todos los threads bloqueados.
- El constructor por defecto de CyclicBarrier es CyclicBarrier(int parties), donde el parámetro representa la cantidad de threads que la barrera interceptará. Cada thread invoca el método await para indicar que ha llegado a la barrera, y luego es bloqueado.
- CyclicBarrier también proporciona un constructor avanzado CyclicBarrier(int parties, Runnable barrierAction) que ejecuta barrierAction cuando los threads alcanzan la barrera, facilitando el manejo de escenarios de negocio más complejos.
- El método getNumberWaiting permite obtener la cantidad de threads bloqueados por CyclicBarrier; el método isBroken() indica si los threads bloqueados han sido interrumpidos.
- El contador de CyclicBarrier puede ser reiniciado utilizando el método reset(), a diferencia de CountDownLatch cuyo contador solo se puede usar una vez. Esto permite que CyclicBarrier maneje escenarios de negocio más complejos, como reiniciar el contador y permitir que los threads se ejecuten nuevamente si ocurre un error en el cálculo.
- CyclicBarrier es útil en escenarios donde múltiples threads calculan datos que posteriormente se combinan.
- CyclicBarrier se implementa internamente usando el lock reentrante ReentrantLock.
public class ServicioBancario implements Runnable {
// Crea 4 barreras, después de procesar todas se ejecuta el método run de esta clase
private CyclicBarrier barrera = new CyclicBarrier(4, this);
// Suponemos 4 tareas de cálculo, así que iniciamos 4 threads
private Executor ejecutor = Executors.newFixedThreadPool(4);
// Guarda el resultado de cálculo de cada tarea
private ConcurrentHashMap<String, Integer> resultadosCuenta = new ConcurrentHashMap<>();
private AtomicInteger contadorAtomico = new AtomicInteger(1);
private void procesar() {
for (int i = 0; i < 4; i++) {
Thread thread = new Thread(() -> {
// Resultado de cálculo de la tarea actual (proceso omitido)
resultadosCuenta.put(Thread.currentThread().getName(), 1);
// Al terminar el cálculo, inserta una barrera
try {
barrera.await();
} catch (InterruptedException | BrokenBarrierException e) {
e.printStackTrace();
}
}, "Hilo" + contadorAtomico.getAndIncrement());
ejecutor.execute(thread);
}
}
@Override
public void run() {
int resultadoFinal = 0;
// Sumariza los resultados de cada tarea
for (Map.Entry<String, Integer> hoja : resultadosCuenta.entrySet()) {
resultadoFinal += hoja.getValue();
}
// Almacena el resultado final
resultadosCuenta.put("resultadoFinal", resultadoFinal);
System.out.println(resultadoFinal);
}
public static void main(String[] args) {
ServicioBancario servicioBancario = new ServicioBancario();
servicioBancario.procesar();
}
}
Semaphore
- Semaphore (señal) se utiliza para controlar el número de threads que acceden simultáneamente a un recurso específico, coordinando los threads para garantizar un uso adecuado de los recursos públicos.
- Semaphore puede ser utilizado para control de flujo, especialmente en escenarios donde los recursos son limitados, como conexiones a bases de datos.
- El constructor de Semaphore, Semaphore(int permits), acepta un número entero que representa la cantidad de permisos disponibles.
- Los threads utilizan el método acquire() de Semaphore para obtener un permiso, y luego llaman a release() para devolverlo. También se puede usar el método tryAcquire() para intentar obtener un permiso sin bloqueo.
- int availablePermits(): Devuelve el número actual de permisos disponibles en este semáforo.
- int getQueueLength(): Devuelve el número de threads que esperan para obtener un permiso.
- hasQueuedThreads(): Indica si hay threads esperando para obtener un permiso.
- Semaphore se implementa internamente usando el lock compartido AQS.
public class PruebaSemáforo {
private static final int CANTIDAD_THREADS = 30;
private static ExecutorService EJECUTOR = Executors.newFixedThreadPool(CANTIDAD_THREADS);
private static Semaphore SEMAFORO = new Semaphore(10);
private static AtomicInteger CONTADOR = new AtomicInteger(1);
public static void main(String[] args) {
for (int i = 0; i < CANTIDAD_THREADS; i++) {
EJECUTOR.execute(() -> {
try {
SEMAFORO.acquire();
System.out.println("guardar datos" + CONTADOR.getAndIncrement());
SEMAFORO.release();
} catch (InterruptedException e) {
}
});
}
EJECUTOR.shutdown();
}
}
Exchanger
- Exchanger es una herramienta para la colaboración entre threads, utilizada para el intercambio de datos entre ellos. Proporciona un punto de sincronización donde dos threads pueden intercambiar sus datos. Estos threads intercambian datos mediante el método exchange. Si el primer thread ejecuta exchange() primero, esperará hasta que el segundo thread también ejecute este método.
- Se puede entender conceptualmente a un objeto Exchanger como un contenedor con dos compartimentos. A través del método exchange, se pueden llenar ambos compartimentos con información. Cuando ambos compartimentos están llenos, el objeto automáticamente intercambia la información entre ellos y la devuelve a los threads, lográndose así el intercambio de datos.
- Exchanger puede utilizarse en algoritmos genéticos (donde se seleccionan dos individuos para cruzar, intercambiando sus datos y aplicando reglas de cruce para obtener resultados).
- Exchanger es útil para tareas de verificación, como cuando un dato necesita ser verificado por dos personas simultáneamente. Solo después de que ambas verificaciones estén correctas, se puede proceder con el procesamiento posterior.
- Exchanger se implementa internamente usando CAS (Compare-And-Swap) sin bloqueos. Exchange utiliza dos atributos del objeto interno Node: item y match, para almacenar los valores de los dos threads respectivamente.
public class PruebaIntercambio {
private static final Exchanger<String> INTERCAMBIADOR = new Exchanger<>();
private static ExecutorService grupoThreads = Executors.newFixedThreadPool(2);
public static void main(String[] args) {
grupoThreads.execute(() -> {
try {
String resultado = INTERCAMBIADOR.exchange("DatoA");
System.out.println("Resultado del intercambio para A: " + resultado);
} catch (InterruptedException e) {
}
});
grupoThreads.execute(() -> {
try {
String resultado = INTERCAMBIADOR.exchange("DatoB");
System.out.println("Resultado del intercambio para B: " + resultado);
} catch (InterruptedException e) {
}
});
grupoThreads.shutdown();
}
}
Phaser
Phaser es una barrera de múltiples fases que permite estalbecer inicialmente el número de threads participantes, y también puede registrar o dar de baja participantes中途. Cuando el número de participantes que llegan satisface la cantidad establecida en la barrera, se produce un avance de fase. A continuación se presentan algunos conceptos básicos de Phaser.
- phase (fase): En cualquier momento, Phaser solo se encuentra en una fase específica. La fase inicial es 0, pudiendo alcanzar hasta Integer.MAX_VALUE y luego volver a cero. Cuando todos los participantes (parties) han llegado, el valor de phase se incrementa.
- parties (participantes): Se refiere a los threads participantes. Phaser puede especificar la cantidad de participantes en su construcción inicial, o también puede registrar o dar de baja participantes mediante métodos como register, bulkRegister, arriveAndDeregister, etc.
- arrive (llegada): Después de registrar los parties en Phaser, el estado inicial de los participantes es unarrived. Cuando un participante llega (arrive) a la fase actual, su estado cambia a arrived.
- advance (avance): Cuando el número de participantes que han llegado a una fase cumple con la condición (el número registrado es igual al número llegado), la fase avanza (advance), es decir, el valor de phase se incrementa en 1.
- termination (terminación): Indica que el objeto Phaser ha alcanzado su estado terminal.
- Tiering (jerarquización): Una estructura en árbol que permite especificar un nodo padre para el objeto Phaser que se está construyendo. 6.1 La jerarquización se introduce porque cuando un Phaser tiene muchos participantes, las operaciones de sincronización internas causan una disminución drástica del rendimiento, mientras que la jerarquización reduce la competencia y, por lo tanto, el sobrecosto adicional debido a la sincornización. 6.2 En una estructura de árbol de Phasers jerárquicos, el registro y eliminación de sub-Phasers o Phaser padre es gestionado automáticamente. Cuando la cantidad de participantes de un Phaser se convierte en 0, y si tiene un nodo padre, se elimina automáticamente de este.
public class PruebaPhaser implements Runnable {
private final Phaser faseador;
public PruebaPhaser(Phaser faseador) {
this.faseador = faseador;
}
@Override
public void run() {
System.out.println(Thread.currentThread().getName() + ": Iniciando tarea, fase actual =" + faseador.getPhase() + "");
// 1. Espera a que otros threads participantes lleguen
// 2. No responde a interrupciones, es decir, aunque el thread actual sea interrumpido,
// el método arriveAndAwaitAdvance no devolverá ni lanzará excepción, sino que continuará esperando.
// Para respuesta a interrupciones, se puede usar awaitAdvanceInterruptibly.
int avanzado = faseador.arriveAndAwaitAdvance();
System.out.println(Thread.currentThread().getName() + ": Finalizando tarea, fase actual =" + avanzado + "");
}
public static void main(String[] args) {
Phaser faseador = new Phaser();
for (int i = 0; i < 10; i++) {
// Agrega nuevos participantes
faseador.register();
new Thread(new PruebaPhaser(faseador), "Hilo-" + i).start();
}
}
}