Al cerrar la lección anterior, BiblioTech ya era correcto: el catálogo soporta dos empleados a la vez, el registro de préstamos mantiene su invariante y el contador no pierde incrementos. Pero el código para llegar hasta ahí era artesanal: candados que hay que soltar en un finally, órdenes de adquisición que hay que respetar por convención, y un hilo creado a mano por cada tarea. Los doscientos avisos de vencimiento siguen necesitando doscientos hilos, doscientos megabytes de pilas y una espera de join() por cada uno.
Esta lección cambia el nivel de abstracción. java.util.concurrent —diseñado por Doug Lea e incorporado en Java 5— ofrece piezas ya construidas y probadas para todo lo que en las lecciones anteriores hacías a mano: pools de hilos que se reutilizan, tareas que devuelven resultados y propagan excepciones, programadores periódicos, y sincronizadores que sustituyen wait/notify por algo que se puede razonar.
El cambio mental es este: dejas de pensar en hilos y empiezas a pensar en tareas. Tú describes qué hay que hacer; el ejecutor decide dónde y cuándo se hace. Es la misma diferencia que hay entre gestionar la memoria a mano y tener un recolector de basura.
Al terminar, los 200 avisos de BiblioTech se enviarán con un pool de ocho hilos en tres segundos, con barra de progreso, con límite de concurrencia sobre el recurso compartido y con cancelación limpia.
Contenido
- Por qué
java.util.concurrentsustituye la gestión manual Executor,ExecutorServiceyExecutors- Las fábricas de
Executorsy cuándo usar cada una - La advertencia sobre las colas ilimitadas
ThreadPoolExecutora medida- Elegir el tamaño del pool
- El ciclo de vida del ejecutor
CallableyFuture- Cancelación de tareas
invokeAlleinvokeAnyScheduledExecutorServiceCountDownLatchCyclicBarrierSemaphoreExchangery tabla comparativa de sincronizadoresForkJoinPooly divide y vencerás- BiblioTech: 200 avisos con pool, progreso y cancelación
- Errores Comunes y Consejos
- Ejercicios
- Por qué
java.util.concurrent sustituye la gestión manual
java.util.concurrent sustituye la gestión manualRecuerda el código de 08-02 para enviar avisos:
// LO QUE HACIAS: un hilo por tarea, gestion manual completa.
Thread[] hilos = new Thread[200];
for (int i = 0; i < 200; i++) {
hilos[i] = new Thread(new TareaAviso(avisos.get(i)), "bibliotech-avisos-" + i);
hilos[i].start();
}
for (Thread h : hilos) {
h.join();
}
// ¿Y el resultado de cada tarea? Campos volatile y leerlos a mano.
// ¿Y si una falla? Un manejador de excepciones no capturadas.
// ¿Y si quiero cancelar todo? interrupt() a los 200, uno por uno.
// ¿Y si son 200.000 avisos? OutOfMemoryError.Seis problemas, todos reales:
| Problema | Consecuencia |
|---|---|
| Un hilo por tarea | ~1 MB de pila y decenas de µs por hilo; no escala |
| Sin límite de hilos | Con muchas tareas, OutOfMemoryError: unable to create native thread |
| Los hilos no se reutilizan | Se paga la creación una y otra vez |
Runnable no devuelve nada |
Campos volatile y convenciones implícitas |
| Excepciones que no llegan al llamante | Fallos silenciosos (08-02) |
| Cancelación manual | Recorrer un array llamando a interrupt() |
java.util.concurrent resuelve los seis. El mismo código pasa a ser:
// LO QUE HARAS: describir tareas, el ejecutor las coloca.
ExecutorService ejecutor = Executors.newFixedThreadPool(8);
List<Future<ResultadoAviso>> futuros = new ArrayList<>();
for (Aviso a : avisos) {
futuros.add(ejecutor.submit(new TareaAviso(a))); // Callable: devuelve valor
}
for (Future<ResultadoAviso> f : futuros) {
ResultadoAviso r = f.get(); // espera, y RELANZA la excepcion si la hubo
}
ejecutor.shutdown();Ocho hilos en lugar de doscientos, resultados tipados, excepciones que llegan, y cancelación con una llamada.
Executor, ExecutorService y Executors
Executor, ExecutorService y ExecutorsTres nombres parecidos que conviene no confundir:
| Nombre | Qué es | Papel |
|---|---|---|
Executor |
Interfaz con un método: void execute(Runnable) |
La abstracción mínima: "ejecuta esto" |
ExecutorService |
Interfaz que extiende Executor |
Añade resultados (submit), apagado y espera |
Executors |
Clase de utilidad con métodos estáticos | Fábrica de implementaciones ya configuradas |
Executor es deliberadamente minimalista, y esa es su virtud: separa qué se ejecuta de cómo se ejecuta.
// La interfaz completa. Un solo metodo.
public interface Executor {
void execute(Runnable comando);
}
// Tres implementaciones validas de la MISMA interfaz:
Executor enElMismoHilo = tarea -> tarea.run();
Executor unHiloNuevo = tarea -> new Thread(tarea).start();
Executor unPool = Executors.newFixedThreadPool(8);
// El codigo cliente no cambia:
enElMismoHilo.execute(() -> System.out.println("hola"));
unPool.execute(() -> System.out.println("hola"));ExecutorService añade lo que hace falta en la práctica:
public interface ExecutorService extends Executor {
// Enviar tareas y obtener un Future
Future<?> submit(Runnable tarea);
<T> Future<T> submit(Callable<T> tarea);
<T> Future<T> submit(Runnable tarea, T resultado);
// Enviar varias
<T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tareas);
<T> T invokeAny(Collection<? extends Callable<T>> tareas);
// Apagado
void shutdown();
List<Runnable> shutdownNow();
boolean isShutdown();
boolean isTerminated();
boolean awaitTermination(long tiempo, TimeUnit unidad);
}Diferencia clave entre execute y submit, y es una de las trampas del tema:
// execute(Runnable): no devuelve nada. Si la tarea lanza una excepcion,
// va al manejador de excepciones no capturadas del hilo (08-02) y
// normalmente se imprime en la consola.
ejecutor.execute(() -> { throw new RuntimeException("fallo visible"); });
// submit(...): devuelve un Future. Si la tarea lanza una excepcion,
// esta se GUARDA en el Future y NO se imprime en ninguna parte.
// Si nadie llama a get(), el fallo desaparece SIN DEJAR RASTRO.
ejecutor.submit(() -> { throw new RuntimeException("fallo INVISIBLE"); });Ese segundo caso es una de las causas más frecuentes de "mi tarea no hace nada y no sale ningún error". La regla: si usas submit, llama a get() o comprueba el resultado de alguna forma.
- Las fábricas de
Executors y cuándo usar cada una
Executors y cuándo usar cada unaimport java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
// 1. Pool FIJO: N hilos, siempre los mismos, cola ilimitada.
ExecutorService fijo = Executors.newFixedThreadPool(8);
// 2. Pool ELASTICO: crea hilos segun necesidad, los recicla a los 60 s.
ExecutorService elastico = Executors.newCachedThreadPool();
// 3. UN SOLO HILO: garantiza orden de ejecucion.
ExecutorService unico = Executors.newSingleThreadExecutor();
// 4. PROGRAMADO: tareas retardadas y periodicas.
ScheduledExecutorService programado = Executors.newScheduledThreadPool(2);
// 5. ROBO DE TRABAJO: ForkJoinPool con paralelismo = nucleos.
ExecutorService robo = Executors.newWorkStealingPool();| Fábrica | Hilos | Cola | Cuándo usarla | Riesgo |
|---|---|---|---|---|
newFixedThreadPool(n) |
Exactamente n |
Ilimitada | Carga estable, cálculo | La cola crece sin límite |
newCachedThreadPool() |
0 … Integer.MAX_VALUE |
SynchronousQueue (sin capacidad) |
Muchas tareas cortas de E/S | Crea hilos sin límite |
newSingleThreadExecutor() |
1 | Ilimitada | Tareas que deben ir en orden | La cola crece sin límite |
newScheduledThreadPool(n) |
n (crece si hace falta) |
Cola con retardo | Trabajo periódico | Una tarea larga retrasa las demás |
newWorkStealingPool() |
= núcleos | Colas por hilo | Cálculo divisible | No garantiza orden |
newVirtualThreadPerTaskExecutor() |
Uno virtual por tarea (Java 21) | — | Miles de tareas de E/S | Ver 10-06 |
Casos de uso en BiblioTech:
// Envio de avisos: E/S, carga acotada y conocida -> pool fijo.
ExecutorService avisos = Executors.newFixedThreadPool(8);
// Escritura del log de auditoria: DEBE conservar el orden -> un solo hilo.
// Ademas, al ser un unico hilo, no hace falta sincronizar el fichero.
ExecutorService auditoria = Executors.newSingleThreadExecutor();
// Aviso de vencimientos cada hora y copia de seguridad cada noche.
ScheduledExecutorService mantenimiento = Executors.newScheduledThreadPool(2);El caso de newSingleThreadExecutor merece un comentario: un ejecutor de un solo hilo es confinamiento (08-04, apartado 17) implementado como servicio. Todo lo que se ejecuta en él está serializado por construcción, así que su estado no necesita ningún candado. Es una técnica muy potente y poco usada.
- La advertencia sobre las colas ilimitadas
Las fábricas de Executors son cómodas y tienen una trampa seria que hay que conocer antes de usarlas en producción.
newFixedThreadPool y newSingleThreadExecutor usan una LinkedBlockingQueue sin límite de capacidad.
Si las tareas llegan más rápido de lo que el pool las consume, la cola crece indefinidamente. No hay ningún mecanismo que lo frene. El resultado, tras minutos u horas de funcionamiento aparentemente normal, es un OutOfMemoryError.
// BOMBA DE RELOJERIA
ExecutorService pool = Executors.newFixedThreadPool(4);
// Cada tarea tarda 100 ms; llegan 1.000 por segundo.
// El pool procesa 40/s. Se acumulan 960 tareas por segundo en memoria.
// En 10 minutos: ~576.000 objetos en cola. OutOfMemoryError.
while (hayPeticiones()) {
pool.submit(new TareaLenta(siguientePeticion()));
}newCachedThreadPool tiene el problema simétrico y peor: su cola es una SynchronousQueue, que no almacena nada, así que cada tarea que llega cuando todos los hilos están ocupados provoca la creación de un hilo nuevo, sin límite (hasta Integer.MAX_VALUE). Con una ráfaga de diez mil tareas lentas, la JVM intenta crear diez mil hilos y muere con OutOfMemoryError: unable to create native thread.
La solución: una cola acotada y una política de rechazo explícita, que es lo que se construye en el apartado siguiente.
Regla profesional. Las fábricas de
Executorsestán bien para código de ejemplo, herramientas internas y tareas de arranque. Para un servicio que recibe carga externa, construye tuThreadPoolExecutorcon cola acotada. Es la recomendación explícita de las guías de estilo de Google y de la mayoría de manuales de la industria.
ThreadPoolExecutor a medida
ThreadPoolExecutor a medidaTodas las fábricas anteriores devuelven, por debajo, un ThreadPoolExecutor. Construirlo directamente da control sobre los seis parámetros que importan.
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
public class EjecutorBiblioTech {
/**
* ThreadFactory que nombra los hilos. Sin ella, los hilos del pool
* se llaman "pool-1-thread-1", que en un volcado (08-03) no dice nada.
*/
static class FabricaHilosNombrados implements ThreadFactory {
private final String prefijo;
private final boolean demonio;
private final AtomicInteger contador = new AtomicInteger(1);
FabricaHilosNombrados(String prefijo, boolean demonio) {
this.prefijo = prefijo;
this.demonio = demonio;
}
@Override
public Thread newThread(Runnable r) {
Thread t = new Thread(r, prefijo + "-" + contador.getAndIncrement());
t.setDaemon(demonio);
// Red de seguridad: cualquier excepcion que escape de una tarea
// enviada con execute() queda registrada (06-07, 08-02).
t.setUncaughtExceptionHandler((hilo, error) ->
System.err.println("[" + hilo.getName() + "] fallo no capturado: " + error));
return t;
}
}
public static ThreadPoolExecutor crear() {
return new ThreadPoolExecutor(
4, // 1. tamano del nucleo
8, // 2. tamano maximo
60L, TimeUnit.SECONDS, // 3. vida de los hilos extra
new ArrayBlockingQueue<>(100), // 4. cola ACOTADA
new FabricaHilosNombrados("bibliotech-avisos", false), // 5.
new ThreadPoolExecutor.CallerRunsPolicy() // 6.
);
}
}Los seis parámetros, uno a uno:
1. corePoolSize — hilos que se mantienen vivos aunque estén ociosos. El pool los crea según llegan tareas y no los destruye (salvo con allowCoreThreadTimeOut(true)).
2. maximumPoolSize — tope absoluto de hilos.
3. keepAliveTime — cuánto sobrevive un hilo por encima del núcleo estando ocioso.
4. workQueue — dónde esperan las tareas. La decisión más importante.
5. threadFactory — cómo se crean los hilos: nombre, demonio, prioridad, manejador de excepciones.
6. RejectedExecutionHandler — qué hacer cuando la cola está llena y no se pueden crear más hilos.
El algoritmo que sigue el pool ante una tarea nueva —y que es donde casi todo el mundo se equivoca al razonar:
flowchart TD
A["Llega una tarea"] --> B{"¿hilos activos<br/>menor que corePoolSize?"}
B -->|Sí| C["Crear un hilo nuevo<br/>y ejecutar"]
B -->|No| D{"¿cabe en la cola?"}
D -->|Sí| E["Encolar<br/>ESPERA su turno"]
D -->|No| F{"¿hilos activos<br/>menor que maximumPoolSize?"}
F -->|Sí| G["Crear un hilo extra<br/>y ejecutar"]
F -->|No| H["RECHAZAR<br/>RejectedExecutionHandler"]
La consecuencia contraintuitiva: el pool prefiere encolar antes que crear hilos extra. Con core=4, max=100 y una cola ilimitada, nunca se crearán más de 4 hilos: la cola nunca se llena, así que el paso que crea hilos extra no se alcanza jamás. Es exactamente lo que ocurre con newFixedThreadPool, y explica por qué maximumPoolSize solo sirve de algo si la cola está acotada.
Las cuatro políticas de rechazo:
| Política | Comportamiento | Cuándo usarla |
|---|---|---|
AbortPolicy (por defecto) |
Lanza RejectedExecutionException |
Cuando perder una tarea es inaceptable y quieres enterarte |
CallerRunsPolicy |
El hilo que envió ejecuta la tarea | Excelente: crea contrapresión natural |
DiscardPolicy |
Descarta en silencio | Casi nunca: los fallos silenciosos son veneno |
DiscardOldestPolicy |
Descarta la más antigua en cola | Datos donde lo reciente vale más (métricas) |
CallerRunsPolicy merece una explicación porque es la más útil y la menos evidente: cuando el pool está saturado, la tarea la ejecuta el propio hilo que llamó a submit. Ese hilo se queda ocupado y, por tanto, deja de producir tareas nuevas durante ese tiempo. El resultado es un mecanismo de contrapresión automático: el productor se ralentiza al ritmo del consumidor, sin colas infinitas y sin descartar nada.
- Elegir el tamaño del pool
Retomando la distinción de 08-01 entre tareas limitadas por CPU y por E/S:
| Tipo de tarea | Fórmula | Ejemplo en BiblioTech |
|---|---|---|
| Limitada por CPU | N o N + 1 |
Recalcular todas las multas |
| Limitada por E/S | N × (1 + espera/cálculo) |
Enviar 200 avisos |
| Mixta | Medir, o separar en dos pools | Importar (leer + validar) |
public class DimensionadoPool {
private static final int NUCLEOS = Runtime.getRuntime().availableProcessors();
/** Calculo puro: mas hilos que nucleos solo anade cambios de contexto. */
public static int paraCalculo() {
return NUCLEOS;
}
/**
* E/S: la formula de Brian Goetz.
* hilos = nucleos * utilizacionObjetivo * (1 + espera / calculo)
* Con utilizacion 1.0, 300 ms de espera y 20 ms de calculo en 8 nucleos:
* 8 * 1.0 * (1 + 15) = 128 hilos.
* En la practica se acota a un maximo razonable: 128 hilos son 128 MB
* de pilas, y el recurso remoto probablemente no aguante 128 peticiones.
*/
public static int paraEs(double msEspera, double msCalculo) {
double ratio = msEspera / msCalculo;
int sugerido = (int) Math.ceil(NUCLEOS * (1 + ratio));
return Math.min(sugerido, 64); // tope de cordura
}
public static void main(String[] args) {
System.out.println("Nucleos : " + NUCLEOS);
System.out.println("Pool de calculo : " + paraCalculo());
System.out.println("Pool de avisos : " + paraEs(300, 20));
System.out.println("Pool de importacion : " + paraEs(50, 30));
}
}El consejo más importante sobre dimensionado: separa los pools por tipo de trabajo. Un único pool compartido entre tareas de cálculo y de E/S es lo peor de los dos mundos: las tareas de E/S ocupan hilos sin usar CPU, y las de cálculo bloquean a las de E/S.
// BIEN: tres pools con propositos y tamanos distintos, y aislados
// entre si: una saturacion del pool de avisos no afecta al de calculo.
private final ExecutorService poolCalculo = Executors.newFixedThreadPool(NUCLEOS);
private final ExecutorService poolEs = Executors.newFixedThreadPool(32);
private final ScheduledExecutorService poolProgramado =
Executors.newScheduledThreadPool(2);Esta idea —aislamiento por mamparos (bulkhead)— viene de los compartimentos estancos de un barco: si uno se inunda, el resto sigue a flote. Se retoma en 12-07.
- El ciclo de vida del ejecutor
Un ExecutorService tiene tres estados: en marcha, apagándose y terminado.
| Método | Qué hace |
|---|---|
shutdown() |
Deja de aceptar tareas nuevas; termina las pendientes. No bloquea |
shutdownNow() |
Deja de aceptar; interrumpe las que corren; devuelve las no empezadas |
isShutdown() |
¿Se ha llamado a alguno de los dos? |
isTerminated() |
¿Han terminado ya todas las tareas? |
awaitTermination(t, u) |
Bloquea hasta que termine todo o se agote el plazo |
El error crítico: si no apagas el ejecutor, la JVM no termina. Los hilos de un pool son no demonio por defecto, así que siguen vivos esperando trabajo indefinidamente (08-01, apartado 8). Tu main acaba y el proceso se queda ahí.
El patrón de apagado correcto, tal como lo recomienda la documentación de ExecutorService:
import java.util.concurrent.ExecutorService;
import java.util.concurrent.TimeUnit;
public final class ApagadoOrdenado {
/**
* Apagado en dos fases:
* 1. shutdown() + espera cortes: dejar terminar lo que esta en curso.
* 2. Si no terminan, shutdownNow() (interrupcion) + espera corta.
* 3. Si aun asi no terminan, registrarlo: hay tareas que no
* responden a la interrupcion, y eso es un BUG de esas tareas
* (el catch vacio de InterruptedException de 08-02).
*/
public static void apagar(ExecutorService ejecutor, long segundos) {
ejecutor.shutdown(); // fase 1: no acepta mas
try {
if (!ejecutor.awaitTermination(segundos, TimeUnit.SECONDS)) {
ejecutor.shutdownNow(); // fase 2: interrumpir
if (!ejecutor.awaitTermination(segundos, TimeUnit.SECONDS)) {
System.err.println("El ejecutor no ha terminado: "
+ "hay tareas que ignoran la interrupcion");
}
}
} catch (InterruptedException e) {
// Nos han interrumpido a NOSOTROS mientras esperabamos.
ejecutor.shutdownNow();
Thread.currentThread().interrupt(); // restaurar la bandera (08-02)
}
}
}Detalle importante de shutdownNow(): interrumpe, no mata. Una tarea que se traga la InterruptedException sigue corriendo tan tranquila. Todo el protocolo de 08-02 se cobra aquí: un catch vacío hace que tu aplicación no se pueda apagar.
Java 19+:
ExecutorServiceesAutoCloseable. Desde Java 19 se puede usar contry-with-resources(06-06), yclose()hace unshutdown()y espera a que terminen las tareas:try (ExecutorService ejecutor = Executors.newFixedThreadPool(8)) { for (Aviso a : avisos) { ejecutor.submit(new TareaAviso(a)); } } // close(): shutdown() + espera indefinidaEs mucho más limpio, pero espera indefinidamente: si una tarea no termina, tu
tryno sale nunca. Para código con plazos, el patrón de dos fases sigue siendo el correcto.
Callable y Future
Callable y FutureEn 08-02 quedó pendiente el problema: Runnable.run() devuelve void y no puede lanzar excepciones comprobadas. Callable lo resuelve:
@FunctionalInterface
public interface Callable<V> {
V call() throws Exception; // devuelve V y PUEDE lanzar
}Runnable |
Callable<V> |
|
|---|---|---|
| Método | void run() |
V call() throws Exception |
| Devuelve valor | No | Sí |
| Excepciones comprobadas | No | Sí |
| Se envía con | execute o submit |
submit |
Future<V> es el objeto que representa "un resultado que llegará". El <V> es simplemente el tipo del resultado: un Future<Informe> promete un Informe (los genéricos, a fondo, en 10-01).
public interface Future<V> {
V get() throws InterruptedException, ExecutionException;
V get(long tiempo, TimeUnit unidad) throws ..., TimeoutException;
boolean cancel(boolean interrumpirSiEstaCorriendo);
boolean isCancelled();
boolean isDone();
}Ejemplo completo con BiblioTech:
package com.nexussoftware.bibliotech.persistencia;
import java.nio.file.Path;
import java.util.concurrent.*;
public class ImportacionConFuture {
public static void main(String[] args) {
ExecutorService ejecutor = Executors.newFixedThreadPool(2);
try {
// Callable<Informe>: devuelve un Informe y puede lanzar
// excepciones comprobadas. Ambas cosas imposibles con Runnable.
Callable<Informe> tarea = () -> {
ImportadorCatalogo imp = new ImportadorCatalogo();
return imp.importar(Path.of("datos/inventario.csv")); // throws IOException
};
Future<Informe> futuro = ejecutor.submit(tarea);
// El hilo principal NO esta bloqueado: puede seguir trabajando.
System.out.println("[main] importacion lanzada, atiendo el menu");
mostrarMenu();
// isDone() no bloquea: permite sondear sin esperar.
while (!futuro.isDone()) {
System.out.println("[main] importando...");
TimeUnit.MILLISECONDS.sleep(500);
}
// get() con PLAZO: nunca uses get() sin plazo en produccion.
Informe informe = futuro.get(30, TimeUnit.SECONDS);
System.out.println("[main] importados: " + informe.importados());
} catch (TimeoutException e) {
System.err.println("[main] la importacion excedio el plazo");
} catch (ExecutionException e) {
// La tarea lanzo una excepcion. get() la ENVUELVE en
// ExecutionException; la original esta en getCause().
// Es exactamente el encadenamiento de causas de 06-03.
Throwable causa = e.getCause();
System.err.println("[main] la importacion fallo: " + causa.getMessage());
if (causa instanceof java.io.IOException) {
System.err.println("[main] problema de fichero; se conserva el catalogo anterior");
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
ApagadoOrdenado.apagar(ejecutor, 10);
}
}
static void mostrarMenu() { /* ... */ }
}Los tres puntos que hay que retener:
1. get() sin plazo bloquea indefinidamente. Si la tarea se cuelga, tu hilo se cuelga con ella. En producción, siempre get(tiempo, unidad).
2. ExecutionException envuelve la causa. La excepción real está en getCause(). Es el encadenamiento de 06-03, y hay que desenvolverla para tomar decisiones:
try {
Informe i = futuro.get(30, TimeUnit.SECONDS);
} catch (ExecutionException e) {
Throwable causa = e.getCause();
// Recuperar el tipo original para aplicar la politica de 06-07
if (causa instanceof FormatoInvalidoException fie) {
registrarYDegradar(fie);
} else if (causa instanceof java.io.IOException ioe) {
abortarConCodigo(ioe);
} else {
throw new BiblioTechException("Fallo inesperado en la importacion", causa);
}
}3. Las tres excepciones de get() significan cosas distintas:
| Excepción | Significado | Reacción |
|---|---|---|
ExecutionException |
La tarea falló | Desenvolver con getCause() y aplicar la política |
TimeoutException |
La tarea sigue corriendo | Decidir: esperar más o cancel(true) |
InterruptedException |
Tú has sido interrumpido | Restaurar la bandera; la tarea sigue viva |
Ojo con la última fila: una TimeoutException no cancela nada. La tarea sigue ejecutándose. Si quieres que pare, hay que llamar a cancel(true) explícitamente.
- Cancelación de tareas
Future.cancel(boolean interrumpirSiEstaCorriendo) es la cancelación de 08-02 expuesta como método.
Future<Informe> futuro = ejecutor.submit(tarea);
// cancel(false): solo evita que EMPIECE si aun esta en cola.
// Si ya esta corriendo, la deja terminar.
futuro.cancel(false);
// cancel(true): ademas, INTERRUMPE el hilo que la ejecuta.
// Es literalmente un hilo.interrupt() sobre el hilo del pool.
futuro.cancel(true);cancel(true) funciona solo si la tarea coopera. Todo el protocolo de 08-02 se aplica igual: comprobar isInterrupted(), capturar InterruptedException, no tragársela. Una tarea que ignora la interrupción es incancelable, esté en un pool o no.
/**
* Tarea cancelable dentro de un pool. Identica al patron de 08-02:
* cancel(true) del Future se traduce en interrupt() sobre este hilo.
*/
public class TareaAvisoCancelable implements Callable<ResultadoAviso> {
private final Aviso aviso;
public TareaAvisoCancelable(Aviso aviso) { this.aviso = aviso; }
@Override
public ResultadoAviso call() throws Exception {
for (int intento = 1; intento <= 3; intento++) {
// Punto de cancelacion en cada vuelta.
if (Thread.currentThread().isInterrupted()) {
throw new InterruptedException("aviso cancelado antes del intento " + intento);
}
try {
enviar(aviso); // puede lanzar InterruptedException
return ResultadoAviso.exito(aviso);
} catch (EnvioFallidoException e) {
if (intento == 3) {
return ResultadoAviso.fallo(aviso, e.getMessage());
}
TimeUnit.MILLISECONDS.sleep(100L * intento); // retroceso
}
}
return ResultadoAviso.fallo(aviso, "agotados los reintentos");
}
private void enviar(Aviso a) throws Exception { /* ... */ }
}Comprobación del estado tras cancelar:
Future<Informe> f = ejecutor.submit(tarea);
TimeUnit.SECONDS.sleep(2);
boolean cancelada = f.cancel(true);
System.out.println("cancel() devolvio : " + cancelada); // false si ya habia terminado
System.out.println("isCancelled() : " + f.isCancelled());
System.out.println("isDone() : " + f.isDone()); // true: cancelada tambien es "hecha"
try {
f.get();
} catch (CancellationException e) {
// get() sobre un Future CANCELADO lanza CancellationException,
// que es NO comprobada. Es facil olvidarse de capturarla.
System.out.println("confirmado: la tarea fue cancelada");
}
invokeAll e invokeAny
invokeAll e invokeAnyDos métodos para enviar una colección de tareas de golpe.
invokeAll: ejecuta todas y bloquea hasta que todas terminen. Devuelve la lista de Future, todos ya completados.
List<Callable<ResultadoAviso>> tareas = new ArrayList<>();
for (Aviso a : avisos) {
tareas.add(new TareaAvisoCancelable(a));
}
// Bloquea hasta que TODAS terminen. Los Future devueltos
// ya estan completos: get() no bloquea.
List<Future<ResultadoAviso>> resultados = ejecutor.invokeAll(tareas);
int exitos = 0, fallos = 0;
for (Future<ResultadoAviso> f : resultados) {
try {
if (f.get().correcto()) exitos++; else fallos++;
} catch (ExecutionException e) {
fallos++;
LOG.log(Level.WARNING, "aviso fallido", e.getCause());
}
}Con plazo global, que es lo recomendable:
// Plazo GLOBAL para el conjunto. Las que no terminen a tiempo
// se CANCELAN automaticamente, y su get() lanzara CancellationException.
List<Future<ResultadoAviso>> resultados =
ejecutor.invokeAll(tareas, 30, TimeUnit.SECONDS);invokeAny: devuelve el resultado de la primera que termine con éxito y cancela las demás. Útil cuando hay varias formas de obtener lo mismo y quieres la más rápida.
// Tres fuentes para el mismo catalogo; nos vale la primera que responda.
List<Callable<Catalogo>> fuentes = List.of(
() -> cargarDesdeCacheLocal(),
() -> cargarDesdeFicheroPrincipal(),
() -> cargarDesdeCopiaSeguridad()
);
// Devuelve el primer resultado correcto; cancela los otros dos.
// Si TODAS fallan, lanza ExecutionException con la ultima causa.
Catalogo c = ejecutor.invokeAny(fuentes, 10, TimeUnit.SECONDS);invokeAll |
invokeAny |
|
|---|---|---|
| Espera a | Todas | La primera con éxito |
| Devuelve | List<Future<T>> |
T (el valor, no un Future) |
| Con las demás | Nada | Las cancela |
| Si alguna falla | Su Future lanza en get() |
Se ignora, salvo que fallen todas |
| Caso de uso | Trabajo en lote | Fuentes redundantes, la más rápida |
ScheduledExecutorService
ScheduledExecutorServicePara trabajo retardado o periódico. Sustituye a Timer/TimerTask, que son de Java 1.3 y tienen dos defectos graves: un solo hilo para todas las tareas, y una excepción no capturada mata el temporizador entero y las demás tareas nunca se vuelven a ejecutar.
import java.util.concurrent.*;
public class MantenimientoBiblioTech {
private final ScheduledExecutorService programador =
Executors.newScheduledThreadPool(2, r -> {
Thread t = new Thread(r, "bibliotech-programador");
t.setDaemon(true); // mantenimiento: no debe impedir el cierre
return t;
});
public void arrancar() {
// 1. UNA vez, dentro de 5 segundos.
programador.schedule(
() -> System.out.println("[mantenimiento] arranque completado"),
5, TimeUnit.SECONDS);
// 2. PERIODICA a ritmo fijo: cada hora desde el instante de inicio.
programador.scheduleAtFixedRate(
this::avisarVencimientos,
0, 1, TimeUnit.HOURS);
// 3. PERIODICA con retardo fijo: 24 h desde que TERMINA la anterior.
programador.scheduleWithFixedDelay(
this::copiaSeguridad,
1, 24, TimeUnit.HOURS);
}
/**
* REGLA DE ORO: una tarea periodica debe capturar TODA excepcion.
* Si deja escapar una, la JVM CANCELA la planificacion y la tarea
* NO SE VUELVE A EJECUTAR JAMAS, en silencio, sin ningun aviso.
* Es el fallo mas traicionero de esta API.
*/
private void avisarVencimientos() {
try {
int enviados = servicioAvisos.enviarVencimientos();
LOG.log(Level.INFO, "Avisos de vencimiento enviados: {0}", enviados);
} catch (Exception e) {
LOG.log(Level.SEVERE, "Fallo al avisar vencimientos", e);
// NO relanzar: si sale, se acaba la planificacion para siempre.
}
}
private void copiaSeguridad() {
try {
gestorCopias.copiar();
} catch (Exception e) {
LOG.log(Level.SEVERE, "Fallo en la copia de seguridad", e);
}
}
}scheduleAtFixedRate frente a scheduleWithFixedDelay — la diferencia importa:
scheduleAtFixedRate(t, i, p, u) |
scheduleWithFixedDelay(t, i, p, u) |
|
|---|---|---|
| El periodo se mide desde | El inicio de la ejecución anterior | El fin de la ejecución anterior |
| Si la tarea dura más que el periodo | Las siguientes se retrasan y encadenan sin solaparse | Siempre hay p de separación real |
| Ritmo | Constante (intenta cumplir el horario) | Variable (depende de lo que dure) |
| Usar para | Trabajo con horario: informes horarios, métricas | Trabajo cuya duración varía: copias, limpieza |
Tarea que dura 3 s, periodo 5 s:
scheduleAtFixedRate: [---3s---]__2s__[---3s---]__2s__[---3s---]
0 3 5 8 10
inicio cada 5 s exactos
scheduleWithFixedDelay:[---3s---]____5s____[---3s---]____5s____[--
0 3 8 11 16
5 s DESPUES de terminarEl detalle que arruina sistemas de producción: si una tarea de scheduleAtFixedRate lanza una excepción no capturada, la planificación se cancela silenciosamente. La tarea no se vuelve a ejecutar nunca, no hay error en el log salvo el que tú pongas, y nadie se entera hasta que alguien pregunta por qué llevan tres semanas sin llegar los avisos. Envuelve siempre el cuerpo en un try/catch (Exception).
CountDownLatch
CountDownLatchUn cierre de cuenta atrás: un contador que solo baja. Los hilos que esperan en await() se desbloquean cuando llega a cero. Es de un solo uso: no se puede reiniciar.
CountDownLatch cierre = new CountDownLatch(5); // 5 eventos pendientes
cierre.countDown(); // resta 1 (nunca por debajo de 0)
cierre.await(); // bloquea hasta llegar a 0
cierre.await(10, TimeUnit.SECONDS); // con plazo; devuelve false si expira
cierre.getCount(); // cuenta actual (para mostrar progreso)Uso 1: esperar a que terminen N tareas.
import java.util.concurrent.*;
public class EsperarNTareas {
public static void main(String[] args) throws InterruptedException {
final int TAREAS = 200;
ExecutorService pool = Executors.newFixedThreadPool(8);
CountDownLatch terminadas = new CountDownLatch(TAREAS);
for (int i = 0; i < TAREAS; i++) {
final int n = i;
pool.execute(() -> {
try {
enviarAviso(n);
} catch (Exception e) {
LOG.log(Level.WARNING, "aviso " + n + " fallido", e);
} finally {
// OBLIGATORIO en finally: si una tarea falla y no
// hace countDown, el await() de main no vuelve JAMAS.
terminadas.countDown();
}
});
}
// Progreso mientras esperamos, sin sondear con sleep a ciegas.
while (!terminadas.await(500, TimeUnit.MILLISECONDS)) {
long pendientes = terminadas.getCount();
System.out.printf("\rProgreso: %d/%d (%.0f%%)",
TAREAS - pendientes, TAREAS,
100.0 * (TAREAS - pendientes) / TAREAS);
}
System.out.println("\nTodos los avisos procesados");
ApagadoOrdenado.apagar(pool, 10);
}
}El countDown() va SIEMPRE en un finally. Si una tarea falla antes de llamarlo, el contador nunca llega a cero y quien espera se queda bloqueado para siempre. Es el error número uno de CountDownLatch.
Uso 2: puerta de salida — arrancar N hilos exactamente a la vez. Muy útil para pruebas de concurrencia, porque maximiza el solapamiento:
/**
* Dos cierres: uno para dar la salida a todos a la vez, otro para
* esperar a que todos acaben. Es el esqueleto de una prueba de estres
* de concurrencia bien hecha: sin la puerta, los primeros hilos
* terminarian antes de que arrancaran los ultimos, y no habria
* solapamiento real (el efecto que viste en el ejercicio 2 de 08-01).
*/
public static long pruebaDeEstres(int hilos, Runnable tarea) throws InterruptedException {
CountDownLatch salida = new CountDownLatch(1); // la pistola
CountDownLatch terminados = new CountDownLatch(hilos); // la meta
for (int i = 0; i < hilos; i++) {
new Thread(() -> {
try {
salida.await(); // todos esperan aqui
tarea.run();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
terminados.countDown();
}
}, "estres-" + i).start();
}
long inicio = System.nanoTime();
salida.countDown(); // ¡ya! los N arrancan a la vez
terminados.await();
return System.nanoTime() - inicio;
}
CyclicBarrier
CyclicBarrierUna barrera cíclica: N hilos se esperan mutuamente en un punto; cuando llega el último, todos continúan y la barrera se reinicia. Esa última palabra es la diferencia con CountDownLatch.
import java.util.concurrent.*;
public class ProcesamientoPorFases {
public static void main(String[] args) {
final int TRABAJADORES = 4;
// La accion de barrera la ejecuta el ULTIMO hilo en llegar,
// antes de liberar a los demas. Ideal para consolidar la fase.
CyclicBarrier barrera = new CyclicBarrier(TRABAJADORES, () ->
System.out.println("--- fase completada por los 4, consolidando ---"));
for (int i = 0; i < TRABAJADORES; i++) {
final int id = i;
new Thread(() -> {
try {
for (int fase = 1; fase <= 3; fase++) {
System.out.printf(" trabajador %d procesa la fase %d%n", id, fase);
TimeUnit.MILLISECONDS.sleep(100 + id * 50);
barrera.await(); // espera a los otros tres
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} catch (BrokenBarrierException e) {
// Si un hilo se interrumpe o expira, la barrera se ROMPE
// y TODOS los demas reciben esta excepcion. Es correcto:
// sin todos los participantes, la fase no tiene sentido.
System.err.println("barrera rota: " + e);
}
}, "trabajador-" + i).start();
}
}
}CountDownLatch |
CyclicBarrier |
|
|---|---|---|
| Reutilizable | No (un solo uso) | Sí (se reinicia) |
| Quién espera | Unos esperan, otros cuentan | Todos los participantes esperan |
| El contador | Baja con countDown() |
Sube al llegar a await() |
| Acción al completarse | No tiene | Sí, Runnable opcional |
| Si un participante falla | Los demás siguen esperando | Se rompe para todos |
| Caso típico | Esperar a que acabe un lote | Simulaciones y cálculos por fases |
En BiblioTech, CountDownLatch es lo que necesitas casi siempre: los 200 avisos son independientes y nadie tiene que esperar a nadie a mitad de camino. CyclicBarrier encaja en cálculos iterativos por fases, donde cada iteración necesita que todas las anteriores hayan terminado.
Semaphore
SemaphoreUn semáforo mantiene un número de permisos. acquire() toma uno (y bloquea si no hay); release() devuelve uno. Es la herramienta para limitar la concurrencia sobre un recurso.
package com.nexussoftware.bibliotech.servicio;
import java.util.concurrent.*;
/**
* Limita a 3 las exportaciones simultaneas.
*
* Motivo: cada exportacion abre el fichero completo del catalogo, lo
* formatea en memoria y lo escribe. Tres a la vez saturan el disco y
* multiplican el uso de memoria; con veinte, el servidor se arrastra.
* El semaforo pone un tope duro sin limitar el resto del pool.
*/
public class ServicioExportacion {
private static final int MAX_SIMULTANEAS = 3;
// 3 permisos. El 'true' activa la equidad: los que llevan mas tiempo
// esperando pasan antes, lo que evita la inanicion (08-04).
private final Semaphore permisos = new Semaphore(MAX_SIMULTANEAS, true);
private final ExecutorService pool = Executors.newFixedThreadPool(16);
public Future<Path> exportar(String formato) {
return pool.submit(() -> {
// Espera un permiso. Si hay 3 exportaciones en curso, este
// hilo se bloquea aqui hasta que una termine.
permisos.acquire();
try {
System.out.printf("[%s] exportando (%d permisos libres)%n",
Thread.currentThread().getName(),
permisos.availablePermits());
return generarFichero(formato); // trabajo pesado
} finally {
// OBLIGATORIO en finally: un permiso no devuelto se
// pierde PARA SIEMPRE. Tras tres fallos, el semaforo
// queda a 0 y nadie vuelve a exportar jamas.
permisos.release();
}
});
}
/** Variante que no espera: rechaza en vez de encolar. */
public Future<Path> exportarSinEsperar(String formato) {
return pool.submit(() -> {
if (!permisos.tryAcquire(5, TimeUnit.SECONDS)) {
throw new ServicioSaturadoException(
"Hay " + MAX_SIMULTANEAS + " exportaciones en curso; reintente");
}
try {
return generarFichero(formato);
} finally {
permisos.release();
}
});
}
private Path generarFichero(String formato) throws Exception {
TimeUnit.SECONDS.sleep(2);
return Path.of("informes/catalogo." + formato);
}
}Usos habituales del semáforo:
| Uso | Cómo |
|---|---|
| Limitar concurrencia sobre un recurso | new Semaphore(n), acquire/release |
| Convertir una colección en acotada | Un permiso por hueco disponible |
| Exclusión mutua simple | new Semaphore(1) — pero prefiere un candado |
| Señalizar entre hilos | Un hilo hace release, otro acquire |
Diferencia con un candado, que se pregunta siempre: un candado es propiedad del hilo que lo tomó y solo él puede soltarlo; un semáforo no tiene dueño, así que un hilo puede adquirir un permiso y otro liberarlo. Eso lo hace más flexible y también más fácil de usar mal —de ahí el finally obligatorio—.
Exchanger y tabla comparativa de sincronizadores
Exchanger y tabla comparativa de sincronizadoresExchanger<V> permite que dos hilos intercambien objetos en un punto de encuentro. Ambos llaman a exchange(objeto) y cada uno recibe el del otro.
Exchanger<List<Material>> intercambio = new Exchanger<>();
// Hilo lector: llena un bufer y lo intercambia por uno vacio.
List<Material> miBufer = new ArrayList<>();
// ... llenar ...
miBufer = intercambio.exchange(miBufer); // recibe el bufer vacio del otro
// Hilo procesador: procesa un bufer lleno y devuelve uno vacio.
List<Material> aProcesar = intercambio.exchange(new ArrayList<>());Es un nicho estrecho —el patrón de doble búfer entre exactamente dos hilos— y en la práctica casi siempre se resuelve mejor con una BlockingQueue (08-06). Conviene conocerlo por si aparece.
Tabla comparativa completa:
| Sincronizador | Participantes | Reutilizable | Para qué |
|---|---|---|---|
CountDownLatch |
Unos cuentan, otros esperan | No | Esperar a que ocurran N eventos |
CyclicBarrier |
N, todos iguales | Sí | Sincronizar fases entre N hilos |
Semaphore |
Cualquiera | Sí | Limitar accesos concurrentes |
Exchanger |
Exactamente 2 | Sí | Intercambio de objetos por parejas |
Phaser (Java 7) |
Variable, dinámico | Sí | Barrera con participantes que entran y salen |
Phaser es una CyclicBarrier más flexible que permite registrar y desregistrar participantes sobre la marcha. Es potente y poco frecuente; menciónalo y no lo uses salvo que la necesidad sea evidente.
ForkJoinPool y divide y vencerás
ForkJoinPool y divide y vencerásForkJoinPool (Java 7) está pensado para tareas que se pueden partir recursivamente en subtareas independientes: el modelo divide y vencerás.
Su característica distintiva es el robo de trabajo (work stealing): cada hilo tiene su propia cola de subtareas, y cuando se queda sin trabajo, roba de la cola de otro hilo por el extremo opuesto. Esto reparte la carga automáticamente sin un coordinador central y con muy poca contención.
flowchart TD
A["Calcular multas de 100.000 prestamos"] --> B["Trozo 1: 0-50.000"]
A --> C["Trozo 2: 50.000-100.000"]
B --> D["0-25.000"]
B --> E["25.000-50.000"]
C --> F["50.000-75.000"]
C --> G["75.000-100.000"]
D --> H["... hasta el umbral<br/>calculo directo"]
E --> H
F --> H
G --> H
H --> I["Combinar resultados<br/>hacia arriba"]
RecursiveTask<V> para tareas que devuelven valor; RecursiveAction para las que no.
package com.nexussoftware.bibliotech.servicio;
import java.util.List;
import java.util.concurrent.ForkJoinPool;
import java.util.concurrent.RecursiveTask;
/**
* Calcula la suma de multas de una lista de prestamos dividiendo
* el trabajo recursivamente.
*
* RecursiveTask<Double>: el <Double> es el tipo del resultado
* que produce esta tarea (10-01 para los genericos).
*/
public class CalculoMultasParalelo extends RecursiveTask<Double> {
/**
* UMBRAL: por debajo de este tamano se calcula directamente.
* Elegirlo bien es lo mas importante de todo el patron:
* - demasiado bajo: el coste de crear y coordinar subtareas
* supera al del calculo, y va MAS LENTO que en secuencial;
* - demasiado alto: no hay bastantes trozos para todos los nucleos.
* Regla practica: entre 100 y 10.000 elementos, y MEDIR.
*/
private static final int UMBRAL = 1_000;
private final List<Prestamo> prestamos;
private final int desde;
private final int hasta;
private final CalculadoraMultas calculadora;
public CalculoMultasParalelo(List<Prestamo> prestamos, int desde, int hasta,
CalculadoraMultas calculadora) {
this.prestamos = prestamos;
this.desde = desde;
this.hasta = hasta;
this.calculadora = calculadora;
}
@Override
protected Double compute() {
int tamano = hasta - desde;
// CASO BASE: trozo pequeno, calculo secuencial directo.
if (tamano <= UMBRAL) {
double suma = 0;
for (int i = desde; i < hasta; i++) {
suma += calculadora.calcular(prestamos.get(i));
}
return suma;
}
// CASO RECURSIVO: partir por la mitad.
int medio = desde + tamano / 2;
CalculoMultasParalelo izq =
new CalculoMultasParalelo(prestamos, desde, medio, calculadora);
CalculoMultasParalelo der =
new CalculoMultasParalelo(prestamos, medio, hasta, calculadora);
// PATRON CORRECTO: bifurcar UNA y calcular la otra en ESTE hilo.
// Hacer izq.fork() + der.fork() + izq.join() + der.join()
// desaprovecha el hilo actual, que se quedaria esperando.
izq.fork(); // a la cola: otro hilo puede robarla
double resultadoDer = der.compute(); // este hilo trabaja
double resultadoIzq = izq.join(); // recoger la bifurcada
return resultadoIzq + resultadoDer;
}
// --- Uso ---
public static double calcularTodas(List<Prestamo> prestamos,
CalculadoraMultas calculadora) {
// El pool COMUN: compartido por toda la JVM, con
// (nucleos - 1) hilos. Lo usan tambien los streams paralelos (10-04).
ForkJoinPool comun = ForkJoinPool.commonPool();
return comun.invoke(new CalculoMultasParalelo(
prestamos, 0, prestamos.size(), calculadora));
}
}El pool común (ForkJoinPool.commonPool()) es una instancia compartida por toda la JVM, con núcleos - 1 hilos, que se usa por defecto para las tareas fork/join y para CompletableFuture (08-07) y los streams paralelos (10-04).
Su gran peligro: como es compartido por toda la aplicación, una tarea que bloquee en él —una lectura de fichero lenta, un get() que espera— roba un hilo a todos los demás usuarios del pool. Con núcleos - 1 hilos, bastan unas pocas tareas bloqueantes para dejarlo inservible.
// PROHIBIDO en el pool comun: operaciones bloqueantes.
ForkJoinPool.commonPool().submit(() -> {
return Files.readAllLines(rutaEnorme); // bloquea un hilo compartido
});
// CORRECTO: pool propio para el trabajo bloqueante.
ForkJoinPool poolPropio = new ForkJoinPool(4);
try {
poolPropio.invoke(new CalculoMultasParalelo(...));
} finally {
poolPropio.shutdown();
}Nota sobre
parallelStream(). Los streams paralelos de la API de colecciones usan exactamente esteForkJoinPool.commonPool()por debajo, y convierten el patrón anterior en una sola línea. No son materia de este módulo: se explican en 10-04, junto con toda la API de Streams, donde también se discute cuándo compensan y cuándo son contraproducentes.
Nota sobre hilos virtuales. Toda la aritmética de dimensionado de pools de esta lección —
Nhilos para CPU,N × (1 + espera/cálculo)para E/S, colas acotadas, políticas de rechazo— existe porque un hilo de plataforma es caro. Los hilos virtuales de Java 21 cuestan nanosegundos y unos cientos de bytes, y con ellos el patrón recomendado para tareas de E/S pasa a serExecutors.newVirtualThreadPerTaskExecutor(): un hilo por tarea, sin pool y sin dimensionar nada. Para tareas de CPU, en cambio, los pools acotados siguen siendo lo correcto. Se explican en 10-06.
- BiblioTech: 200 avisos con pool, progreso y cancelación
Todo junto. Es el caso B de 08-01, resuelto de verdad.
package com.nexussoftware.bibliotech.servicio;
import com.nexussoftware.bibliotech.dominio.Prestamo;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.logging.Level;
import java.util.logging.Logger;
/**
* Envio masivo de avisos de vencimiento.
*
* Propiedades:
* - Pool ACOTADO con hilos nombrados y cola acotada.
* - SEMAFORO que limita a 5 los accesos simultaneos al fichero de avisos.
* - COUNTDOWNLATCH para mostrar progreso y esperar el final.
* - Cancelacion limpia mediante shutdownNow() -> interrupcion cooperativa.
* - Resultado agregado con contadores atomicos (detalle en 08-06).
*/
public class ServicioAvisos implements AutoCloseable {
private static final Logger LOG = Logger.getLogger(ServicioAvisos.class.getName());
private static final int HILOS = 8;
private static final int MAX_ESCRITURAS_SIMULTANEAS = 5;
private final ThreadPoolExecutor pool;
private final Semaphore accesoFichero =
new Semaphore(MAX_ESCRITURAS_SIMULTANEAS, true);
public ServicioAvisos() {
AtomicInteger n = new AtomicInteger(1);
this.pool = new ThreadPoolExecutor(
HILOS, HILOS,
0L, TimeUnit.MILLISECONDS,
new ArrayBlockingQueue<>(500), // cola ACOTADA
r -> {
Thread t = new Thread(r, "bibliotech-avisos-" + n.getAndIncrement());
t.setUncaughtExceptionHandler((h, e) ->
LOG.log(Level.SEVERE, "Fallo no capturado en " + h.getName(), e));
return t;
},
new ThreadPoolExecutor.CallerRunsPolicy()); // contrapresion
}
/** Resultado agregado del envio. Inmutable (08-04). */
public record ResumenEnvio(int total, int enviados, int fallidos,
int cancelados, long milisegundos) {
public double porcentajeExito() {
return total == 0 ? 0 : 100.0 * enviados / total;
}
}
/**
* Envia todos los avisos en paralelo.
*
* @param plazoSegundos plazo global; al agotarse se cancela lo pendiente
*/
public ResumenEnvio enviarTodos(List<Prestamo> vencidos, long plazoSegundos)
throws InterruptedException {
final int total = vencidos.size();
final CountDownLatch terminados = new CountDownLatch(total);
final AtomicInteger enviados = new AtomicInteger();
final AtomicInteger fallidos = new AtomicInteger();
long inicio = System.nanoTime();
List<Future<?>> futuros = new ArrayList<>(total);
// 1. Enviar todas las tareas al pool.
for (Prestamo p : vencidos) {
futuros.add(pool.submit(() -> {
try {
// Punto de cancelacion antes de empezar.
if (Thread.currentThread().isInterrupted()) {
return;
}
// El semaforo limita el acceso concurrente al fichero.
accesoFichero.acquire();
try {
escribirAviso(p);
enviados.incrementAndGet();
} finally {
accesoFichero.release(); // SIEMPRE
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt(); // restaurar (08-02)
} catch (Exception e) {
fallidos.incrementAndGet();
LOG.log(Level.WARNING, "Aviso fallido para " + p.id(), e);
} finally {
// SIEMPRE, o el await() no vuelve nunca.
terminados.countDown();
}
}));
}
// 2. Progreso mientras se espera, con plazo global.
long limiteNs = TimeUnit.SECONDS.toNanos(plazoSegundos);
boolean completado = false;
while (!(completado = terminados.await(300, TimeUnit.MILLISECONDS))) {
long hechos = total - terminados.getCount();
System.out.printf("\r [%s] %d/%d avisos (%.0f%%)",
barra(hechos, total), hechos, total, 100.0 * hechos / total);
if (System.nanoTime() - inicio > limiteNs) {
System.out.println("\n Plazo agotado: cancelando lo pendiente");
for (Future<?> f : futuros) {
f.cancel(true); // interrumpe o descarta de la cola
}
break;
}
}
long ms = (System.nanoTime() - inicio) / 1_000_000;
int cancelados = total - enviados.get() - fallidos.get();
System.out.printf("\r [%s] %d/%d avisos (100%%)%n",
barra(total, total), total, total);
return new ResumenEnvio(total, enviados.get(), fallidos.get(), cancelados, ms);
}
private static String barra(long hechos, long total) {
int ancho = 30;
int llenos = total == 0 ? 0 : (int) (ancho * hechos / total);
return "#".repeat(llenos) + "-".repeat(ancho - llenos);
}
/** Simula la escritura del aviso: 250 ms de E/S. */
private void escribirAviso(Prestamo p) throws InterruptedException {
TimeUnit.MILLISECONDS.sleep(250);
if (p.id().hashCode() % 37 == 0) {
throw new IllegalStateException("empleado sin direccion de contacto");
}
}
/** AutoCloseable: apagado en dos fases (06-06 + apartado 7). */
@Override
public void close() {
pool.shutdown();
try {
if (!pool.awaitTermination(30, TimeUnit.SECONDS)) {
pool.shutdownNow();
if (!pool.awaitTermination(10, TimeUnit.SECONDS)) {
LOG.severe("Hay tareas de aviso que ignoran la interrupcion");
}
}
} catch (InterruptedException e) {
pool.shutdownNow();
Thread.currentThread().interrupt();
}
}
}Uso:
package com.nexussoftware.bibliotech.presentacion;
import java.util.List;
public class DemostracionAvisos {
public static void main(String[] args) throws InterruptedException {
List<Prestamo> vencidos = generarVencidos(200);
System.out.println("=== ENVIO DE AVISOS DE VENCIMIENTO ===");
System.out.println("Avisos pendientes: " + vencidos.size());
System.out.println("Tiempo secuencial estimado: "
+ (vencidos.size() * 250 / 1000) + " s");
System.out.println();
// try-with-resources: close() apaga el pool ordenadamente.
try (ServicioAvisos servicio = new ServicioAvisos()) {
ServicioAvisos.ResumenEnvio r = servicio.enviarTodos(vencidos, 60);
System.out.println();
System.out.println("=== RESUMEN ===");
System.out.println("Total : " + r.total());
System.out.println("Enviados : " + r.enviados());
System.out.println("Fallidos : " + r.fallidos());
System.out.println("Cancelados : " + r.cancelados());
System.out.println("Tiempo : " + r.milisegundos() + " ms");
System.out.printf ("Exito : %.1f%%%n", r.porcentajeExito());
System.out.printf ("Aceleracion: %.1fx sobre el envio secuencial%n",
(vencidos.size() * 250.0) / r.milisegundos());
}
}
}Salida:
=== ENVIO DE AVISOS DE VENCIMIENTO ===
Avisos pendientes: 200
Tiempo secuencial estimado: 50 s
[##############################] 200/200 avisos (100%)
=== RESUMEN ===
Total : 200
Enviados : 194
Fallidos : 6
Cancelados : 0
Tiempo : 10247 ms
Exito : 97.0%
Aceleracion: 4.9x sobre el envio secuencialAnálisis honesto del resultado, que es donde está el aprendizaje:
- La aceleración es 4,9× y no 8×, aunque el pool tenga 8 hilos. La causa es el semáforo de 5 permisos: solo 5 avisos pueden escribir a la vez, así que el paralelismo real está limitado por ese 5, no por los 8 hilos. 200 × 250 ms / 5 ≈ 10 s, exactamente lo medido. El cuello de botella no es el pool: es el recurso protegido, y esto es la norma, no la excepción, en sistemas reales.
- Si subes el semáforo a 8, la aceleración sube a ~8× y el tiempo baja a 6,3 s. Si lo subes a 20 sin subir los hilos, no cambia nada: el límite pasa a ser el pool. Dimensionar es encontrar qué recurso es el escaso.
- Los 6 fallos no detuvieron el envío. Cada tarea captura sus propias excepciones y hace
countDown()en elfinally: la política de degradación de 06-07, ahora en concurrencia. - Ocho hilos, no doscientos. 8 MB de pilas en lugar de 200 MB, y ocho creaciones de hilo en lugar de doscientas.
Errores Comunes y Consejos
Error 1: no apagar el ejecutor. Los hilos del pool son no demonio; la JVM no termina. Es la causa número uno de "mi programa no acaba".
Error 2: usar newFixedThreadPool con carga externa. Cola ilimitada → OutOfMemoryError tras horas de aparente normalidad. Cola acotada y política de rechazo.
Error 3: usar newCachedThreadPool con muchas tareas lentas. Crea hilos sin límite hasta tumbar la JVM.
Error 4: enviar con submit y no llamar a get(). La excepción de la tarea queda guardada en el Future y no aparece en ninguna parte. Fallo silencioso total. Usa execute si no te interesa el resultado, o comprueba el Future.
Error 5: get() sin plazo en producción. Si la tarea se cuelga, el llamante se cuelga. Siempre get(tiempo, unidad).
Error 6: no desenvolver ExecutionException. La excepción real está en getCause(). Registrar la ExecutionException tal cual esconde la información útil.
Error 7: dejar que una tarea de scheduleAtFixedRate lance una excepción. La planificación se cancela para siempre, en silencio. Envuelve el cuerpo en try/catch (Exception).
Error 8: countDown() fuera del finally. Si la tarea falla antes, quien espera se bloquea indefinidamente.
Error 9: release() de un semáforo fuera del finally. Los permisos perdidos no se recuperan; tras unos fallos, el semáforo queda a cero y nadie pasa nunca más.
Error 10: creer que maximumPoolSize se alcanza con cola ilimitada. El pool encola antes de crear hilos extra: con cola ilimitada nunca pasa de corePoolSize.
Error 11: bloquear en el ForkJoinPool.commonPool(). Es compartido por toda la JVM y tiene núcleos - 1 hilos. Unas pocas tareas bloqueantes lo inutilizan para todos.
Error 12: usar un único pool para cálculo y para E/S. Se estorban mutuamente. Pools separados por tipo de trabajo.
Consejo 1: nombra los hilos con una ThreadFactory. pool-1-thread-3 no dice nada en un volcado; bibliotech-avisos-3 sí. Es el consejo de 08-02 aplicado a los pools.
Consejo 2: usa CallerRunsPolicy cuando no puedas perder tareas. Crea contrapresión automática: el productor se frena solo.
Consejo 3: mide antes de dimensionar. Las fórmulas son puntos de partida. El tamaño real depende de tu carga, tu hardware y tus dependencias externas.
Consejo 4: separa pools por tipo de trabajo (mamparos). Aísla los fallos: que la saturación de los avisos no impida generar un informe.
Consejo 5: un ejecutor de un solo hilo es confinamiento gratis. Todo lo que corre en él está serializado, así que su estado no necesita candados. Muy útil para escrituras en fichero y para logs ordenados.
Consejo 6: recuerda que cancel(true) es interrupt(). Sin tareas cooperativas (08-02), la cancelación no funciona, estés en un pool o no.
Ejercicios
Ejercicio 1: Comparativa de pools
Escribe ComparativaPools que ejecute la misma carga —300 tareas que duermen 100 ms cada una, simulando E/S— con cinco estrategias: (a) un hilo por tarea a mano, (b) newFixedThreadPool(4), (c) newFixedThreadPool(32), (d) newCachedThreadPool(), y (e) un ThreadPoolExecutor a medida con 16 hilos, cola de 50 y CallerRunsPolicy. Para cada una mide el tiempo total con System.nanoTime() y cuenta cuántos hilos distintos ejecutaron tareas (usando un Set sincronizado de nombres de hilo). Imprime una tabla y explica los resultados.
Ejercicio 2: Importación con Callable, plazo y cancelación
Reescribe la importación de catálogo de BiblioTech como Callable<Informe> ejecutado en un ExecutorService. El programa debe:
- Lanzar la importación y seguir mostrando un menú simulado.
- Mostrar el progreso mientras
!futuro.isDone(). - Aplicar un plazo de 5 segundos con
get(5, TimeUnit.SECONDS). - Al agotarse el plazo, llamar a
cancel(true)y verificar conisCancelled(). - Distinguir en el
catchlos tres finales:ExecutionException(desenvolviendo la causa),TimeoutExceptioneInterruptedException. - Apagar el ejecutor con el patrón de dos fases.
Ejercicio 3: Sala de lectura con semáforo y cierre
Modela la sala de lectura de BiblioTech: caben 4 empleados a la vez y hay 20 empleados que quieren entrar. Escribe SalaDeLectura con:
- Un
Semaphore(4, true)que controla el aforo, conentrar()/salir()y elreleaseenfinally. - Un
CountDownLatchde "puerta de salida" que haga que los 20 hilos intenten entrar exactamente a la vez, y otro de "meta" para esperar a que todos hayan terminado. - Un contador atómico del aforo instantáneo, y una comprobación que falle ruidosamente si alguna vez supera 4.
- Un
ScheduledExecutorServiceque imprima el aforo actual cada 200 ms mientras dure la simulación, y que se apague al final.
Cada empleado permanece en la sala un tiempo aleatorio entre 200 y 600 ms.
Soluciones
Solución al Ejercicio 1
import java.util.Collections;
import java.util.HashSet;
import java.util.Set;
import java.util.concurrent.*;
public class ComparativaPools {
static final int TAREAS = 300;
static final long DURACION_MS = 100;
/** Tarea de E/S simulada que registra qué hilo la ejecutó. */
static Runnable tarea(Set<String> hilosUsados) {
return () -> {
hilosUsados.add(Thread.currentThread().getName());
try {
TimeUnit.MILLISECONDS.sleep(DURACION_MS);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
};
}
record Resultado(String estrategia, long ms, int hilos) { }
/** (a) Un hilo por tarea, gestionado a mano. */
static Resultado unHiloPorTarea() throws InterruptedException {
Set<String> hilos = Collections.synchronizedSet(new HashSet<>());
Thread[] ts = new Thread[TAREAS];
long inicio = System.nanoTime();
for (int i = 0; i < TAREAS; i++) {
ts[i] = new Thread(tarea(hilos), "manual-" + i);
ts[i].start();
}
for (Thread t : ts) t.join();
return new Resultado("un hilo por tarea",
(System.nanoTime() - inicio) / 1_000_000, hilos.size());
}
/** (b)-(e) Cualquier ExecutorService. */
static Resultado conEjecutor(String nombre, ExecutorService ej)
throws InterruptedException {
Set<String> hilos = Collections.synchronizedSet(new HashSet<>());
long inicio = System.nanoTime();
for (int i = 0; i < TAREAS; i++) {
ej.execute(tarea(hilos));
}
ej.shutdown();
ej.awaitTermination(2, TimeUnit.MINUTES);
return new Resultado(nombre,
(System.nanoTime() - inicio) / 1_000_000, hilos.size());
}
public static void main(String[] args) throws InterruptedException {
System.out.printf("%d tareas de %d ms de E/S simulada%n", TAREAS, DURACION_MS);
System.out.printf("Secuencial seria: %d ms%n%n", TAREAS * DURACION_MS);
Resultado[] resultados = {
unHiloPorTarea(),
conEjecutor("fixed(4)", Executors.newFixedThreadPool(4)),
conEjecutor("fixed(32)", Executors.newFixedThreadPool(32)),
conEjecutor("cached", Executors.newCachedThreadPool()),
conEjecutor("a medida (16, cola 50, CallerRuns)",
new ThreadPoolExecutor(16, 16, 0L, TimeUnit.MILLISECONDS,
new ArrayBlockingQueue<>(50),
new ThreadPoolExecutor.CallerRunsPolicy()))
};
System.out.printf("%-38s | %8s | %8s%n", "Estrategia", "ms", "hilos");
System.out.println("---------------------------------------|----------|---------");
for (Resultado r : resultados) {
System.out.printf("%-38s | %8d | %8d%n", r.estrategia(), r.ms(), r.hilos());
}
}
}Salida orientativa (máquina de 8 núcleos):
300 tareas de 100 ms de E/S simulada
Secuencial seria: 30000 ms
Estrategia | ms | hilos
---------------------------------------|----------|---------
un hilo por tarea | 147 | 300
fixed(4) | 7621 | 4
fixed(32) | 1024 | 32
cached | 163 | 300
a medida (16, cola 50, CallerRuns) | 1938 | 17Análisis, fila a fila:
- Un hilo por tarea (147 ms) es el más rápido. Con 300 tareas puramente de espera, tener 300 hilos esperando a la vez es óptimo en tiempo. Pero son 300 MB de pilas reservadas; con 30.000 tareas,
OutOfMemoryError. Rápido y no escalable. fixed(4)(7,6 s) es el más lento: solo 4 tareas a la vez, 300/4 = 75 tandas × 100 ms. Dimensionar un pool de E/S como si fuera de CPU es el error clásico.fixed(32)(1,0 s) es un buen equilibrio: 32 hilos, 32 MB, y casi 30× de aceleración.cached(163 ms) creó 300 hilos, igual que la versión manual: suSynchronousQueueno almacena, así que cada tarea que llega con todos ocupados provoca un hilo nuevo. Rápido aquí, y una bomba con carga mayor.- El pool a medida (1,9 s) usó 17 hilos: los 16 del pool más el hilo de
main. Ese hilo extra esCallerRunsPolicyen acción: cuando la cola de 50 se llenó,mainejecutó tareas él mismo. Eso es contrapresión funcionando —el productor se frenó solo— y es por lo que la memoria nunca creció.
Solución al Ejercicio 2
package com.nexussoftware.bibliotech.persistencia;
import java.nio.file.Path;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.logging.Level;
import java.util.logging.Logger;
public class ImportacionConPlazo {
private static final Logger LOG = Logger.getLogger(ImportacionConPlazo.class.getName());
/** Informe inmutable del resultado. */
public record Informe(int importados, int descartados, long ms) { }
/**
* Importacion como Callable<Informe>: devuelve valor y puede
* lanzar excepciones comprobadas. Es cancelable: comprueba la
* interrupcion en cada linea (08-02).
*/
static class TareaImportacion implements Callable<Informe> {
private final Path fichero;
private final AtomicInteger progreso = new AtomicInteger();
TareaImportacion(Path fichero) { this.fichero = fichero; }
int progreso() { return progreso.get(); }
@Override
public Informe call() throws Exception {
long inicio = System.nanoTime();
int importados = 0, descartados = 0;
for (int linea = 0; linea < 50_000; linea++) {
// PUNTO DE CANCELACION: cancel(true) del Future se
// traduce en interrupt() sobre este hilo.
if (Thread.currentThread().isInterrupted()) {
throw new InterruptedException(
"importacion cancelada en la linea " + linea);
}
if (linea % 97 == 0) descartados++; else importados++;
progreso.incrementAndGet();
if (linea % 1000 == 0) {
TimeUnit.MILLISECONDS.sleep(30); // simula validacion pesada
}
}
return new Informe(importados, descartados,
(System.nanoTime() - inicio) / 1_000_000);
}
}
public static void main(String[] args) {
ExecutorService ejecutor = Executors.newSingleThreadExecutor(
r -> new Thread(r, "bibliotech-importador"));
TareaImportacion tarea = new TareaImportacion(Path.of("datos/inventario.csv"));
Future<Informe> futuro = ejecutor.submit(tarea);
try {
// 1-2. El hilo principal sigue vivo: atiende el menu y pinta progreso.
System.out.println("[main] importacion lanzada; el menu sigue activo");
while (!futuro.isDone()) {
System.out.printf("\r[main] progreso: %d lineas", tarea.progreso());
TimeUnit.MILLISECONDS.sleep(300);
}
// 3. Plazo de 5 segundos.
Informe informe = futuro.get(5, TimeUnit.SECONDS);
System.out.printf("%n[main] COMPLETADA: %d importados, %d descartados, %d ms%n",
informe.importados(), informe.descartados(), informe.ms());
} catch (TimeoutException e) {
// 4. El plazo NO cancela por si solo: hay que pedirlo.
System.out.printf("%n[main] plazo agotado tras %d lineas; cancelando%n",
tarea.progreso());
boolean pedida = futuro.cancel(true);
System.out.println("[main] cancel() devolvio : " + pedida);
System.out.println("[main] isCancelled() : " + futuro.isCancelled());
System.out.println("[main] el catalogo conserva su contenido anterior");
} catch (ExecutionException e) {
// 5. La tarea fallo: la causa real esta en getCause() (06-03).
Throwable causa = e.getCause();
if (causa instanceof InterruptedException) {
System.out.println("\n[main] la importacion se detuvo por cancelacion");
} else {
LOG.log(Level.SEVERE, "La importacion fallo", causa);
System.out.println("\n[main] fallo: " + causa.getMessage());
}
} catch (InterruptedException e) {
// Nos han interrumpido a NOSOTROS; la tarea sigue viva.
System.out.println("\n[main] espera interrumpida");
futuro.cancel(true);
Thread.currentThread().interrupt();
} finally {
// 6. Apagado en dos fases.
ejecutor.shutdown();
try {
if (!ejecutor.awaitTermination(5, TimeUnit.SECONDS)) {
ejecutor.shutdownNow();
if (!ejecutor.awaitTermination(5, TimeUnit.SECONDS)) {
LOG.severe("El importador ignora la interrupcion");
}
}
} catch (InterruptedException e) {
ejecutor.shutdownNow();
Thread.currentThread().interrupt();
}
System.out.println("[main] ejecutor apagado; la JVM puede terminar");
}
}
}Salida (cuando se agota el plazo):
[main] importacion lanzada; el menu sigue activo
[main] progreso: 41893 lineas
[main] plazo agotado tras 42017 lineas; cancelando
[main] cancel() devolvio : true
[main] isCancelled() : true
[main] el catalogo conserva su contenido anterior
[main] ejecutor apagado; la JVM puede terminarTres puntos:
TimeoutExceptionno cancela nada por sí sola. La tarea seguía corriendo tan feliz; hubo que llamar acancel(true). Este es el error de expectativa más frecuente conFuture.cancel(true)funcionó porque la tarea coopera. Sin la comprobación deisInterrupted()en el bucle, la importación habría seguido hasta el final y el ejecutor no se habría podido apagar.- El apagado en dos fases al final garantiza que la JVM pueda terminar. Sin él, el hilo
bibliotech-importadorseguiría vivo.
Solución al Ejercicio 3
package com.nexussoftware.bibliotech.servicio;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
public class SalaDeLectura {
private static final int AFORO = 4;
private static final int EMPLEADOS = 20;
// Semaforo EQUITATIVO: los que llevan mas esperando entran antes,
// lo que evita la inanicion (08-04, apartado 12).
private final Semaphore aforo = new Semaphore(AFORO, true);
// Aforo instantaneo, para comprobar el invariante.
private final AtomicInteger dentro = new AtomicInteger();
private final AtomicInteger maximoObservado = new AtomicInteger();
private volatile boolean invarianteRoto = false;
/** Un empleado entra, permanece un rato y sale. */
public void visitar(String empleado) throws InterruptedException {
aforo.acquire(); // bloquea si ya hay 4 dentro
try {
int actual = dentro.incrementAndGet();
// Registro del maximo observado, con bucle CAS (08-06).
maximoObservado.updateAndGet(m -> Math.max(m, actual));
if (actual > AFORO) {
invarianteRoto = true;
System.err.println("!!! INVARIANTE ROTO: " + actual + " dentro !!!");
}
System.out.printf(" [+] %-12s entra (dentro=%d, esperando=%d)%n",
empleado, actual, aforo.getQueueLength());
TimeUnit.MILLISECONDS.sleep(
ThreadLocalRandom.current().nextInt(200, 601));
System.out.printf(" [-] %-12s sale (dentro=%d)%n",
empleado, dentro.get() - 1);
} finally {
dentro.decrementAndGet();
// OBLIGATORIO en finally: un permiso perdido no vuelve nunca,
// y tras cuatro perdidas la sala quedaria cerrada para siempre.
aforo.release();
}
}
public static void main(String[] args) throws InterruptedException {
SalaDeLectura sala = new SalaDeLectura();
// Dos cierres: la "pistola de salida" y la "meta".
CountDownLatch salida = new CountDownLatch(1);
CountDownLatch meta = new CountDownLatch(EMPLEADOS);
// Monitor periodico del aforo.
ScheduledExecutorService monitor = Executors.newSingleThreadScheduledExecutor(
r -> {
Thread t = new Thread(r, "bibliotech-monitor-aforo");
t.setDaemon(true);
return t;
});
monitor.scheduleAtFixedRate(() -> {
try {
System.out.printf(" [monitor] aforo=%d/%d cola=%d%n",
sala.dentro.get(), AFORO, sala.aforo.getQueueLength());
} catch (Exception e) {
// Una excepcion no capturada CANCELARIA la planificacion
// para siempre y en silencio.
e.printStackTrace();
}
}, 200, 200, TimeUnit.MILLISECONDS);
// 20 hilos que esperan la salida y arrancan a la vez.
for (int i = 1; i <= EMPLEADOS; i++) {
final String nombre = "empleado-" + String.format("%02d", i);
new Thread(() -> {
try {
salida.await(); // todos esperan aqui
sala.visitar(nombre);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
meta.countDown(); // SIEMPRE
}
}, nombre).start();
}
System.out.println("=== SALA DE LECTURA (aforo " + AFORO + ") ===");
TimeUnit.MILLISECONDS.sleep(200); // que todos lleguen a await()
long inicio = System.nanoTime();
salida.countDown(); // ¡a la vez!
meta.await();
long ms = (System.nanoTime() - inicio) / 1_000_000;
monitor.shutdown();
monitor.awaitTermination(2, TimeUnit.SECONDS);
System.out.println();
System.out.println("=== RESUMEN ===");
System.out.println("Empleados atendidos : " + EMPLEADOS);
System.out.println("Aforo maximo : " + sala.maximoObservado.get()
+ " (limite " + AFORO + ")");
System.out.println("Invariante respetado : " + !sala.invarianteRoto);
System.out.println("Tiempo total : " + ms + " ms");
System.out.println("Dentro al final : " + sala.dentro.get());
System.out.println("Permisos disponibles : " + sala.aforo.availablePermits()
+ " (debe ser " + AFORO + ")");
}
}Salida (fragmento):
=== SALA DE LECTURA (aforo 4) ===
[+] empleado-03 entra (dentro=1, esperando=16)
[+] empleado-07 entra (dentro=2, esperando=15)
[+] empleado-01 entra (dentro=3, esperando=14)
[+] empleado-12 entra (dentro=4, esperando=13)
[monitor] aforo=4/4 cola=16
[-] empleado-03 sale (dentro=3)
[+] empleado-05 entra (dentro=4, esperando=15)
...
=== RESUMEN ===
Empleados atendidos : 20
Aforo maximo : 4 (limite 4)
Invariante respetado : true
Tiempo total : 2143 ms
Dentro al final : 0
Permisos disponibles : 4 (debe ser 4)Cuatro observaciones:
- El aforo máximo observado es exactamente 4, nunca 5, en cualquier ejecución. El semáforo cumple su contrato.
- Los permisos disponibles al final vuelven a ser 4. Es la comprobación de que ningún
release()se perdió. Si quitas elfinallyy provocas una excepción dentro, verás ese número bajar y la sala quedarse cerrada. - La puerta de salida con
CountDownLatchmaximiza la contención. Sin ella, los primeros empleados habrían entrado y salido antes de que arrancaran los últimos, y no habría habido cola. Es la técnica correcta para una prueba de concurrencia. - Los 2,1 segundos totales cuadran: 20 empleados × ~400 ms de media / 4 simultáneos ≈ 2 s. El aforo, no el número de hilos, determina el tiempo — la misma lección del semáforo de avisos del apartado 17.
Conclusión
Has cambiado de nivel de abstracción: ya no gestionas hilos, describes tareas.
Sabes por qué java.util.concurrent sustituye a la gestión manual: un hilo por tarea no escala, no hay límite que impida crear diez mil, los hilos no se reutilizan, Runnable no devuelve nada, las excepciones no llegan al llamante y cancelar es recorrer un array. Los seis problemas desaparecen con un ejecutor.
Conoces la trinidad Executor / ExecutorService / Executors y la diferencia entre execute —la excepción se ve— y submit —la excepción se guarda en el Future y desaparece si nadie llama a get()—, que es una de las trampas más caras del tema. Sabes qué hace cada fábrica de Executors, y sobre todo conoces la advertencia de las colas ilimitadas: newFixedThreadPool acumula tareas hasta el OutOfMemoryError, y newCachedThreadPool crea hilos sin límite hasta el mismo final por otro camino. Por eso sabes construir un ThreadPoolExecutor a medida, con sus seis parámetros, su ThreadFactory que nombra los hilos, su cola acotada y su política de rechazo —con CallerRunsPolicy como la más útil, porque convierte la saturación en contrapresión automática—. Y entiendes el algoritmo contraintuitivo del pool: encola antes de crear hilos extra, así que maximumPoolSize no sirve de nada si la cola es ilimitada.
Sabes dimensionar: N para cálculo, N × (1 + espera/cálculo) para E/S, y —el consejo que más rinde— pools separados por tipo de trabajo, el aislamiento por mamparos que impide que la saturación de una parte hunda las demás. Y sabes apagar un ejecutor con el patrón de dos fases —shutdown, esperar, shutdownNow, esperar, registrar—, recordando que si no lo apagas la JVM no termina, y que shutdownNow() interrumpe pero no mata: una tarea que se traga la InterruptedException sigue viva, y todo el protocolo de 08-02 se cobra aquí.
Callable y Future resuelven lo que quedó abierto en 08-02: una tarea que devuelve un valor y propaga excepciones comprobadas. Con las tres excepciones de get() bien distinguidas —ExecutionException que envuelve la causa real en getCause(), TimeoutException que no cancela nada y deja la tarea corriendo, e InterruptedException que te afecta a ti y no a la tarea—, la regla de nunca usar get() sin plazo en producción, y cancel(true) que no es magia sino un interrupt() sobre el hilo del pool. Más invokeAll para lotes con plazo global e invokeAny para quedarse con la fuente más rápida.
Conoces ScheduledExecutorService y la diferencia entre scheduleAtFixedRate —periodo desde el inicio, ritmo constante— y scheduleWithFixedDelay —periodo desde el fin, separación real garantizada—, con la advertencia que arruina sistemas: una excepción no capturada en una tarea periódica cancela la planificación para siempre, en silencio.
Y tienes los sincronizadores: CountDownLatch de un solo uso, para esperar N eventos y para dar la salida a N hilos a la vez, con el countDown() siempre en un finally; CyclicBarrier reutilizable, con su acción de barrera y su ruptura para todos si un participante falla; Semaphore para poner un tope duro a la concurrencia sobre un recurso, con el release() siempre en un finally porque un permiso perdido no vuelve jamás; y Exchanger y Phaser reconocidos aunque rara vez necesarios. Más ForkJoinPool con RecursiveTask, su robo de trabajo, el patrón correcto fork() + compute() + join(), la importancia crítica del umbral, y el aviso sobre el pool común: es compartido por toda la JVM y bloquearlo perjudica a todos.
BiblioTech envía sus 200 avisos en diez segundos con ocho hilos, con barra de progreso mediante CountDownLatch, con un semáforo que limita a cinco los accesos simultáneos al fichero, con plazo global y cancelación limpia, con cada fallo individual registrado sin detener el lote, y con AutoCloseable para que el try-with-resources apague el pool. Y el análisis del resultado enseña la lección más valiosa del apartado: la aceleración fue 4,9× y no 8× porque el cuello de botella no era el pool, sino el semáforo de cinco permisos. Dimensionar bien es encontrar cuál es el recurso escaso — y casi nunca es el que crees.
Pero queda una grieta. ServicioAvisos usa AtomicInteger para los contadores sin haberlos explicado, y el Catalogo que hiciste seguro en 08-04 sigue basándose en un HashMap protegido con candados: cada lectura paga un lock(), aunque leer no estorbe a nadie. Y la ColaReservas del módulo 5 sigue siendo un ArrayDeque que dos hilos corromperían, con el patrón productor-consumidor prometido desde entonces y todavía sin implementar.
En la próxima lección, Colecciones Concurrentes y Variables Atómicas, se cierra esa grieta. Verás una demostración real de un HashMap corrompido por dos hilos —con el bucle infinito que puede dejar un núcleo al 100 % para siempre— y la ConcurrentModificationException del fail-fast de 05-02; las tres generaciones de solución, desde Collections.synchronizedMap con su trampa —iterar sigue requiriendo sincronización manual y las operaciones compuestas siguen sin ser atómicas— hasta ConcurrentHashMap con sus operaciones atómicas compuestas putIfAbsent, computeIfAbsent, compute y merge, que son la forma correcta de resolver comprobar-luego-actuar; CopyOnWriteArrayList para los oyentes de eventos; las BlockingQueue en todas sus variantes y, por fin, la implementación completa del patrón productor-consumidor prometido en el módulo 5, aplicado a la cola de reservas y con píldora venenosa para el apagado; y las variables atómicas con la operación compare-and-swap explicada por dentro, su bucle de reintento, el problema ABA y LongAdder para contadores bajo mucha contención. Al terminar tendrás la tabla de decisión definitiva: cuándo synchronized, cuándo un Lock, cuándo un atómico, cuándo una colección concurrente y cuándo simplemente un objeto inmutable.
Curso de Programación en Java
Módulo 1: Introducción a Java
- Introducción a Java
- Configuración del Entorno de Desarrollo
- Sintaxis y Estructura Básica
- Variables y Tipos de Datos
- Operadores
- Entrada y Salida por Consola
- Tu Primer Programa Completo: BiblioTech
Módulo 2: Flujo de Control
- Sentencias Condicionales
- Bucles
- Sentencias Switch
- Break y Continue
- Depuración y Trazas de Ejecución
- Proyecto: Menú Interactivo de BiblioTech
Módulo 3: Programación Orientada a Objetos
- Introducción a la POO
- Clases y Objetos
- Métodos
- Constructores
- Herencia
- Polimorfismo
- Encapsulamiento
- Abstracción
- La Clase Object: equals, hashCode y toString
Módulo 4: Programación Orientada a Objetos Avanzada
- Interfaces
- Clases Abstractas
- Clases Internas
- Clases Anónimas
- Expresiones Lambda
- Interfaces Funcionales y Referencias a Métodos
- Enumeraciones y Registros
Módulo 5: Estructuras de Datos y Colecciones
- Arreglos
- El Framework de Colecciones
- ArrayList
- LinkedList
- HashMap
- HashSet
- Cola y Deque
- Pila
- Ordenación y Búsqueda en Colecciones
Módulo 6: Manejo de Excepciones
- Introducción a las Excepciones
- Bloque Try-Catch
- Throw y Throws
- Excepciones Personalizadas
- Bloque Finally
- Try-with-resources y AutoCloseable
- Estrategias de Manejo de Errores y Logging
Módulo 7: Entrada/Salida de Archivos
- Lectura de Archivos
- Escritura de Archivos
- Flujos de Archivos
- BufferedReader y BufferedWriter
- Serialización
- La API NIO.2: Path y Files
- Formatos de Intercambio: CSV y Properties
Módulo 8: Multihilo y Concurrencia
- Introducción al Multihilo
- Creación de Hilos
- Ciclo de Vida de un Hilo
- Sincronización
- Utilidades de Concurrencia
- Colecciones Concurrentes y Variables Atómicas
- Tareas Asíncronas con CompletableFuture
Módulo 9: Redes
- Introducción a las Redes
- Sockets
- ServerSocket
- DatagramSocket y DatagramPacket
- URL y HttpURLConnection
- El Cliente HTTP Moderno
Módulo 10: Temas Avanzados
- Genéricos
- Anotaciones
- Reflexión
- Características de Java 8: Streams y Optional
- Fechas y Horas con java.time
- Java 9 y Más Allá
- Memoria, Recolección de Basura y Rendimiento
Módulo 11: Frameworks y Librerías de Java
- Introducción a los Frameworks de Java
- Spring Framework
- Hibernate
- JUnit
- Maven
- Pruebas Avanzadas con Mockito
- Librerías Esenciales del Ecosistema
