Los executors de Java ejecutan tus tareas en un pool fijo de hilos y te devuelven un Future. Aprende submit, invokeAll, la cancelación, las cadenas de CompletableFuture, los timeouts, los semáforos y las colas bloqueantes, con una salida que no varía.
Arrancar un hilo a mano para cada trabajo funciona en un programa pequeño. En un servidor no, porque cada trabajo paga por un hilo nuevo y nada impide que una ráfaga de trabajo cree miles. El paquete java.util.concurrent te da executors en su lugar: un conjunto fijo de hilos que toman tareas de una cola y te devuelven un resultado que puedes esperar.
Este post cubre ExecutorService, Future, las excepciones y la cancelación en las tareas, el tamaño del pool, CompletableFuture, CountDownLatch, Semaphore y dos colecciones concurrentes. Cada programa de abajo se ejecutó en Java 25, y su salida está copiada de esa ejecución. Para ejecutar uno tú mismo, guárdalo como Main.java y ejecuta java Main.java. Igual que en la parte sobre hilos y estado compartido, los programas concurrentes fuerzan su orden con latches, así que imprimen lo mismo en cada ejecución.
Por qué executors: un pool de hilos y una cola
Un hilo por tarea tiene dos costos. Cada hilo de plataforma es un hilo del sistema operativo con su propio stack, que en Linux de 64 bits reserva 1 MB por defecto, y crear uno cuesta trabajo de verdad. Y la cantidad no tiene límite: diez mil peticiones que llegan a la vez significan diez mil hilos.
Un executor resuelve las dos cosas. Executors.newFixedThreadPool(4) crea cuatro hilos y los conserva. Las tareas que le entregas esperan en una cola hasta que uno de los cuatro queda libre. Este programa envía diez tareas que se bloquean en una compuerta, y luego mira dentro del pool:
void main() throws InterruptedException {
var running = new AtomicInteger();
var mostAtOnce = new AtomicInteger();
var fourStarted = new CountDownLatch(4);
var gate = new CountDownLatch(1);
var done = new AtomicInteger();
try (ExecutorService pool = Executors.newFixedThreadPool(4)) {
for (int i = 0; i < 10; i++) {
pool.execute(() -> {
mostAtOnce.accumulateAndGet(running.incrementAndGet(), Math::max);
fourStarted.countDown();
awaitQuietly(gate);
running.decrementAndGet();
done.incrementAndGet();
});
}
fourStarted.await();
var queue = ((ThreadPoolExecutor) pool).getQueue();
IO.println("tasks running: " + running.get());
IO.println("tasks waiting in the queue: " + queue.size());
gate.countDown();
} // close() waits for every task to finish
IO.println("tasks finished: " + done.get());
IO.println("most tasks running at once: " + mostAtOnce.get());
}
void awaitQuietly(CountDownLatch latch) {
try {
latch.await();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
Imprime:
tasks running: 4
tasks waiting in the queue: 6
tasks finished: 10
most tasks running at once: 4
Cuatro tareas arrancaron y se quedaron atascadas en la compuerta. Las otras seis esperaron en la cola, porque no había ningún hilo libre para ejecutarlas. Cuando main abrió la compuerta, los cuatro hilos fueron vaciando la cola, y en ningún momento se ejecutaron más de cuatro tareas.
execute le entrega al pool un Runnable y no devuelve nada. El cast a ThreadPoolExecutor solo está ahí para espiar la cola. No lo vas a necesitar en código real.
El bloque try es un try-with-resources. ExecutorService es AutoCloseable desde Java 19 (comprobado: javac --release 18 rechaza este programa), y close() hace que el pool deje de aceptar tareas nuevas y espera a que terminen las que están en la cola. Sin él, los hilos del pool mantienen viva la JVM después de que main retorna.
Un pool fijo de cuatro. Cuatro tareas se ejecutan, seis esperan en la cola, y quien llama tiene un Future por cada tarea que envió. Los hilos se crean una sola vez y se reutilizan para todas las tareas.
Explicado como si tuvieras diez años
Un pool de hilos es la cocina de un restaurante con una cantidad fija de cocineros y un riel de comandas. El restaurante no contrata a un cocinero nuevo por cada pedido. Tú entregas una comanda, y va al riel. Cuando un cocinero queda libre, toma la siguiente comanda del riel y prepara ese plato.
A cambio de tu comanda recibes un localizador de esos que vibran. Ese es el Future. Puedes sentarte y hacer otra cosa, y el localizador se enciende cuando la comida está lista. Si prefieres quedarte parado en el mostrador esperando, eso es get().
A la hora de cerrar, el gerente puede dejar de recibir comandas y dejar que los cocineros terminen las del riel. O puede quitar todas las comandas del riel y decirles a los cocineros que paren ya.
La versión precisa
newFixedThreadPool(n) devuelve un ThreadPoolExecutor con n hilos trabajadores y una LinkedBlockingQueue de tareas. Cada hilo repite un ciclo: toma una tarea de la cola, la ejecuta y vuelve a empezar. submit envuelve tu tarea en una FutureTask, que es a la vez lo que ejecuta el hilo y el Future que tú tienes. Cuando la tarea retorna o lanza una excepción, la FutureTask guarda el resultado o la excepción, y cualquier hilo bloqueado en get() se despierta.
Dónde falla la analogía: el riel de newFixedThreadPool no tiene límite de largo, así que una avalancha de tareas se acumula en memoria en lugar de rechazarse. El localizador además guarda el resultado, incluido un fallo, no solo una señal. Y un cocinero al que le dicen que pare puede ignorarlo: en Java, detener una tarea en ejecución es una petición que la tarea tiene que respetar, como muestra la sección sobre cancelación.
submit, invokeAll e invokeAny
submit recibe un Callable que devuelve un valor, o un Runnable, y te da un Future. invokeAll envía una lista de tareas y espera a todas. invokeAny espera a la primera que tenga éxito:
void main() throws Exception {
try (var pool = Executors.newFixedThreadPool(3)) {
Future<Integer> length = pool.submit(() -> "executor".length());
IO.println("submit: get() returned " + length.get());
var thirdDone = new CountDownLatch(1);
List<Callable<String>> jobs = List.of(
() -> {
thirdDone.await(); // finish last on purpose
return "first";
},
() -> "second",
() -> {
thirdDone.countDown();
return "third";
});
var results = new ArrayList<String>();
for (Future<String> f : pool.invokeAll(jobs)) {
results.add(f.get());
}
IO.println("invokeAll: " + results);
List<Callable<String>> mirrors = List.of(
() -> { throw new IOException("mirror A is down"); },
() -> "downloaded from mirror B",
() -> { throw new IOException("mirror C is down"); });
IO.println("invokeAny: " + pool.invokeAny(mirrors));
}
}
Imprime:
submit: get() returned 8
invokeAll: [first, second, third]
invokeAny: downloaded from mirror B
get() se bloquea hasta que la tarea termina y después devuelve su valor.
El primer trabajo espera a que termine el tercero, así que los trabajos terminan desordenados. Aun así, invokeAll devuelve los futures en el orden en que enviaste las tareas, así que leerlos en un bucle da un orden fijo, sin importar cómo corrieron los hilos.
invokeAny devolvió el resultado del mirror B, porque las otras dos tareas lanzaron una excepción. Cuando una tarea tiene éxito, cancela las demás. Si todas fallan, lanza ExecutionException con uno de los fallos dentro.
Detener un pool: shutdown y shutdownNow
Un pool se detiene de una de dos formas, y se diferencian en lo que pasa con la cola. Aquí cada pool tiene un hilo, una tarea en ejecución bloqueada en una compuerta y tres tareas en cola detrás de ella:
void main() throws InterruptedException {
IO.println("shutdown(): " + finishedTasks(false));
IO.println("shutdownNow(): " + finishedTasks(true));
}
String finishedTasks(boolean now) throws InterruptedException {
var pool = Executors.newFixedThreadPool(1);
var started = new CountDownLatch(1);
var gate = new CountDownLatch(1);
var finished = new AtomicInteger();
pool.execute(() -> {
started.countDown();
try {
gate.await();
finished.incrementAndGet();
} catch (InterruptedException e) {
Thread.currentThread().interrupt(); // interrupted: give up
}
});
for (int i = 0; i < 3; i++) {
pool.execute(finished::incrementAndGet);
}
started.await(); // one task running, three in the queue
int neverStarted = 0;
if (now) {
neverStarted = pool.shutdownNow().size();
} else {
pool.shutdown();
}
String rejected = "";
try {
pool.execute(finished::incrementAndGet);
} catch (RejectedExecutionException e) {
rejected = "new task rejected, ";
}
gate.countDown();
boolean terminated = pool.awaitTermination(5, TimeUnit.SECONDS);
return rejected + "finished " + finished.get() + ", never started " + neverStarted
+ ", terminated " + terminated;
}
Imprime:
shutdown(): new task rejected, finished 4, never started 0, terminated true
shutdownNow(): new task rejected, finished 0, never started 3, terminated true
Los dos métodos hacen que el pool rechace tareas nuevas con RejectedExecutionException.
shutdown() deja que terminen la tarea en ejecución y toda la cola, así que terminaron las cuatro. shutdownNow() quitó las tres tareas de la cola, las devolvió como una lista e interrumpió la que estaba en ejecución. Esa tarea capturó la interrupción y se rindió, así que no terminó ninguna.
awaitTermination espera, hasta un tiempo límite, a que terminen los hilos del pool, y devuelve si terminaron. close() es más o menos shutdown() seguido de awaitTermination en un bucle.
Excepciones dentro de las tareas
Una excepción lanzada dentro de una tarea no llega enseguida al código que la envió. A dónde va depende de cómo entregaste la tarea. Este programa le da a su pool una thread factory que instala un manejador de excepciones no capturadas, para que veamos qué le llega:
void main() throws InterruptedException {
var handlerSaw = new LinkedBlockingQueue<String>();
var threadsMade = new AtomicInteger();
ThreadFactory withHandler = task -> {
threadsMade.incrementAndGet();
var thread = new Thread(task);
thread.setUncaughtExceptionHandler((t, e) -> handlerSaw.add(e.getMessage()));
return thread;
};
try (var pool = Executors.newFixedThreadPool(1, withHandler)) {
Future<Integer> parsed = pool.submit(() -> Integer.parseInt("forty-two"));
try {
parsed.get();
} catch (ExecutionException e) {
IO.println("submit: get() threw " + e.getClass().getSimpleName());
IO.println(" cause: " + e.getCause());
}
pool.execute(() -> Integer.parseInt("seven"));
IO.println("execute: the handler got " + handlerSaw.take());
IO.println(" threads the pool has made: " + threadsMade.get());
Future<?> forgotten = pool.submit(() -> Integer.parseInt("nine"));
while (!forgotten.isDone()) {
Thread.sleep(1);
}
IO.println("submit, no get(): state is " + forgotten.state());
IO.println(" and the handler got " + handlerSaw.poll());
}
}
Imprime:
submit: get() threw ExecutionException
cause: java.lang.NumberFormatException: For input string: "forty-two"
execute: the handler got For input string: "seven"
threads the pool has made: 2
submit, no get(): state is FAILED
and the handler got null
Con submit, el Future captura la excepción y la guarda. get() lanza ExecutionException, y getCause() es la NumberFormatException original.
Con execute, no hay ningún Future que la guarde. La excepción se escapa de la tarea y mata al hilo trabajador, así que la recibe el manejador de excepciones no capturadas del hilo, y el pool crea un hilo de reemplazo. Sin un manejador propio, el de por defecto imprime el stack trace en la salida de error estándar, y main nunca se entera.
El tercer caso es el que pierde errores. Una tarea enviada con submit falló, y nadie llamó a get(). El manejador no recibió nada, y no se imprimió nada en ningún lado. Future.state(), agregado en Java 19, muestra FAILED, pero solo si alguien pregunta. Cuando envíes una tarea, guarda su Future y llama a get() sobre él.
Timeouts y cancelación
get tiene una versión con timeout (tiempo límite de espera), y cancel(true) interrumpe una tarea que ya está en ejecución. Aquí la tarea espera en un latch que nunca se abre, así que el timeout de 50 ms vence con seguridad:
void main() throws Exception {
var neverOpens = new CountDownLatch(1);
var started = new CountDownLatch(1);
var sawInterrupt = new CountDownLatch(1);
try (var pool = Executors.newFixedThreadPool(1)) {
Future<String> slow = pool.submit(() -> {
started.countDown();
try {
neverOpens.await();
return "finished";
} catch (InterruptedException e) {
sawInterrupt.countDown();
throw e;
}
});
started.await();
try {
slow.get(50, TimeUnit.MILLISECONDS);
} catch (TimeoutException e) {
IO.println("get(50 ms) threw TimeoutException");
}
IO.println("the task is still running: " + !slow.isDone());
boolean cancelled = slow.cancel(true);
sawInterrupt.await();
IO.println("cancel(true) returned " + cancelled);
IO.println("the task saw the interrupt and stopped waiting");
IO.println("state: " + slow.state());
try {
slow.get();
} catch (CancellationException e) {
IO.println("get() now throws CancellationException");
}
}
}
Imprime:
get(50 ms) threw TimeoutException
the task is still running: true
cancel(true) returned true
the task saw the interrupt and stopped waiting
state: CANCELLED
get() now throws CancellationException
TimeoutException solo significa que tú dejaste de esperar. La tarea siguió en ejecución. Después, cancel(true) interrumpió su hilo, await() lanzó InterruptedException y la tarea terminó. A partir de ahí, get() lanza CancellationException.
cancel(false) no interrumpe. Igual marca el future como cancelado, pero una tarea que ya está en ejecución sigue hasta el final. Cualquiera de los dos tipos de cancel impide que una tarea en cola llegue a arrancar.
La interrupción es una petición
Una interrupción no detiene un hilo. Activa un indicador (flag) en el hilo, y los métodos bloqueantes como await, sleep y BlockingQueue.take notan el indicador y lanzan InterruptedException. Cuando lo lanzan, limpian el indicador. Así que una tarea que captura la excepción y sigue como si nada borró la petición.
Estas dos tareas son idénticas salvo por una línea en el catch:
void main() throws InterruptedException {
var neverOpens = new CountDownLatch(1);
var bothStarted = new CountDownLatch(2);
var politeStopped = new CountDownLatch(1);
var rudeIgnoredIt = new CountDownLatch(1);
var pool = Executors.newFixedThreadPool(2, Thread.ofPlatform().daemon().factory());
pool.execute(() -> {
bothStarted.countDown();
while (!Thread.currentThread().isInterrupted()) {
try {
neverOpens.await();
} catch (InterruptedException e) {
Thread.currentThread().interrupt(); // put the flag back
}
}
politeStopped.countDown();
});
pool.execute(() -> {
bothStarted.countDown();
while (!Thread.currentThread().isInterrupted()) {
try {
neverOpens.await();
} catch (InterruptedException e) {
rudeIgnoredIt.countDown(); // swallowed: the flag stays clear
}
}
});
bothStarted.await();
pool.shutdownNow(); // interrupts both tasks
politeStopped.await();
rudeIgnoredIt.await();
IO.println("the task that restored the flag stopped");
IO.println("the task that swallowed the interrupt went back to waiting");
IO.println("pool terminated within 200 ms: "
+ pool.awaitTermination(200, TimeUnit.MILLISECONDS));
}
Imprime:
the task that restored the flag stopped
the task that swallowed the interrupt went back to waiting
pool terminated within 200 ms: false
La primera tarea volvió a activar el indicador con Thread.currentThread().interrupt(), así que la condición de su bucle lo vio y la tarea terminó. La segunda se tragó la excepción, el indicador quedó limpio y el bucle volvió directo a esperar. shutdownNow envió una interrupción y no tiene una segunda para enviar, así que el pool no puede terminar. Sus hilos son daemon aquí solo para que el programa pueda salir.
La regla: o dejas que InterruptedException se propague, o la capturas y llamas a Thread.currentThread().interrupt(). Un bucle largo sin llamadas bloqueantes debería revisar Thread.currentThread().isInterrupted() por su cuenta. El método estático Thread.interrupted() también lee el indicador, pero además lo limpia, así que úsalo solo cuando manejas la interrupción ahí mismo.
¿Cuántos hilos?
El tamaño correcto del pool depende de en qué pasan el tiempo las tareas, y no hay una respuesta exacta que puedas buscar.
Las tareas limitadas por CPU (CPU-bound), como parsear, comprimir o hacer cálculos numéricos, necesitan un núcleo todo el tiempo. Más hilos que núcleos significa que los hilos se turnan, y cambiar entre ellos cuesta tiempo. Empieza cerca de Runtime.getRuntime().availableProcessors().
Las tareas limitadas por IO (IO-bound), como las llamadas a una base de datos o a otro servicio, pasan la mayor parte del tiempo esperando. Un hilo que espera no usa CPU, así que más hilos que núcleos pueden ayudar.
La regla general clásica, de Java Concurrency in Practice, es:
hilos = núcleos × (1 + tiempo de espera / tiempo de cómputo)
Una tarea que espera 90 ms y calcula 10 ms en una máquina de 8 núcleos da 8 × (1 + 9) = 80 hilos. Tómalo como una primera estimación, no como una respuesta. Puede que la base de datos permita solo 20 conexiones, y entonces 80 hilos solo hacen fila frente a ella. Mide con carga real y ajusta.
Para el trabajo limitado por IO, los hilos virtuales de Java 21 suelen eliminar la pregunta. La próxima parte los cubre.
CompletableFuture: encadenar pasos
Un Future te da una sola forma de usar el resultado, que es bloquearte en get(). Un CompletableFuture te deja adjuntar el siguiente paso, que se ejecuta cuando llega el valor:
record User(int id, String name) {}
void main() {
try (var pool = Executors.newFixedThreadPool(4)) {
CompletableFuture<User> user = CompletableFuture.supplyAsync(() -> findUser(7), pool);
CompletableFuture<String> name = user.thenApply(User::name);
IO.println("thenApply: " + name.join());
CompletableFuture<CompletableFuture<Integer>> nested =
user.thenApply(u -> cartTotal(u, pool));
CompletableFuture<Integer> total = user.thenCompose(u -> cartTotal(u, pool));
IO.println("thenApply, nested: " + nested.join().join());
IO.println("thenCompose: " + total.join());
CompletableFuture<Integer> shipping = CompletableFuture.supplyAsync(() -> 5, pool);
CompletableFuture<Integer> toPay = total.thenCombine(shipping, Integer::sum);
IO.println("thenCombine: " + toPay.join());
}
}
User findUser(int id) {
return new User(id, "Ana"); // stands in for a database call
}
CompletableFuture<Integer> cartTotal(User u, Executor pool) {
return CompletableFuture.supplyAsync(() -> 37, pool); // stands in for another service
}
Imprime:
thenApply: Ana
thenApply, nested: 37
thenCompose: 37
thenCombine: 42
supplyAsync ejecuta el supplier en el executor que le pasas. thenApply transforma el valor cuando está listo, como map en un stream.
thenApply y thenCompose se diferencian cuando el siguiente paso devuelve a su vez un future. Con thenApply, obtienes un future de un future, CompletableFuture<CompletableFuture<Integer>>, y necesitas dos llamadas a join(). thenCompose lo aplana en un solo CompletableFuture<Integer>, como flatMap. thenCombine espera a dos futures independientes y combina sus valores.
join() espera igual que get(), pero no lanza excepciones comprobadas, así que encaja dentro de las lambdas.
Pásale un executor. Sin uno, supplyAsync usa ForkJoinPool.commonPool(). Ese pool se comparte con los parallel streams y con todo lo demás en la JVM, y su tamaño está pensado para trabajo de CPU: un hilo menos que los núcleos de la máquina. Las llamadas bloqueantes ahí pueden dejar sin hilos a código que no tiene nada que ver. Sus hilos además son daemon (comprobado: Thread.currentThread().isDaemon() devolvió true dentro de uno), así que la JVM puede salir mientras tu tarea todavía se está ejecutando.
Esperar a varios con allOf
CompletableFuture.allOf devuelve un future que se completa cuando se completaron todos los futures que le pasas. No lleva ningún valor, así que lees cada uno después. Aquí cada pronóstico espera al anterior, así que terminan en orden inverso:
void main() {
var finishOrder = new ConcurrentLinkedQueue<String>();
try (var pool = Executors.newFixedThreadPool(3)) {
CompletableFuture<?> nothing = CompletableFuture.completedFuture(null);
// each forecast waits for the one before it, so they finish in reverse
var beijing = forecast("Beijing", nothing, finishOrder, pool);
var madrid = forecast("Madrid", beijing, finishOrder, pool);
var lisbon = forecast("Lisbon", madrid, finishOrder, pool);
CompletableFuture.allOf(lisbon, madrid, beijing).join();
IO.println("finished: " + finishOrder);
for (var f : List.of(lisbon, madrid, beijing)) {
IO.println(f.join());
}
}
}
CompletableFuture<String> forecast(
String city, CompletableFuture<?> after, Queue<String> finishOrder, Executor pool) {
return CompletableFuture.supplyAsync(() -> {
after.join();
finishOrder.add(city);
int celsius = switch (city) {
case "Lisbon" -> 24;
case "Madrid" -> 31;
default -> 28;
};
return city + ": " + celsius + "C";
}, pool);
}
Imprime:
finished: [Beijing, Madrid, Lisbon]
Lisbon: 24C
Madrid: 31C
Beijing: 28C
El primer pronóstico en terminar fue el de Beijing. Después de allOf(...).join(), todos los futures están listos, así que cada join() del bucle retorna de inmediato. El bucle los lee en el orden que tú elegiste, así que el orden de la salida es fijo.
Cuando falla una etapa
Un CompletableFuture envuelve un fallo de forma distinta según cómo lo esperes:
void main() throws InterruptedException {
try (var pool = Executors.newFixedThreadPool(2)) {
CompletableFuture<Integer> port =
CompletableFuture.supplyAsync(() -> Integer.parseInt("eighty"), pool);
try {
port.join();
} catch (CompletionException e) {
IO.println("join() threw CompletionException");
IO.println(" cause: " + e.getCause());
}
try {
port.get();
} catch (ExecutionException e) {
IO.println("get() threw ExecutionException");
IO.println(" cause: " + e.getCause());
}
int withDefault = port.exceptionally(ex -> 8080).join();
IO.println("exceptionally: " + withDefault);
String report = port
.thenApply(p -> "listening on " + p)
.handle((value, ex) -> {
if (ex == null) {
return value;
}
Throwable real = ex instanceof CompletionException ? ex.getCause() : ex;
return "handle got " + ex.getClass().getSimpleName()
+ ", real problem: " + real.getMessage();
})
.join();
IO.println(report);
}
}
Imprime:
join() threw CompletionException
cause: java.lang.NumberFormatException: For input string: "eighty"
get() threw ExecutionException
cause: java.lang.NumberFormatException: For input string: "eighty"
exceptionally: 8080
handle got CompletionException, real problem: For input string: "eighty"
join() envuelve el fallo en una CompletionException no comprobada. get() lo envuelve en una ExecutionException comprobada, igual que Future.get(). En los dos casos, getCause() es la excepción real.
exceptionally se ejecuta solo si hay un fallo y aporta un valor de reemplazo. handle se ejecuta en los dos casos y recibe el valor o la excepción; uno de los dos es null.
Esto nos sorprendió: handle recibió una CompletionException, no la NumberFormatException. Un fallo que viene de una etapa anterior llega envuelto. También vimos que exceptionally, llamado directamente sobre el future de supplyAsync, también recibió el envoltorio. Así que desenvuélvelo antes de mirarlo, como hace el programa.
Timeouts en un CompletableFuture
orTimeout hace fallar el future después de un tiempo, y completeOnTimeout lo completa con un valor de respaldo. Las dos búsquedas de aquí esperan en un latch que nunca se abre:
void main() {
var neverOpens = new CountDownLatch(1);
var slowTaskEnded = new CountDownLatch(2);
try (var pool = Executors.newFixedThreadPool(2)) {
Supplier<String> slowLookup = () -> {
try {
neverOpens.await();
return "live price";
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return "interrupted";
} finally {
slowTaskEnded.countDown();
}
};
CompletableFuture<String> strict = CompletableFuture
.supplyAsync(slowLookup, pool)
.orTimeout(50, TimeUnit.MILLISECONDS);
try {
strict.join();
} catch (CompletionException e) {
IO.println("orTimeout: " + e.getCause().getClass().getSimpleName());
}
String price = CompletableFuture
.supplyAsync(slowLookup, pool)
.completeOnTimeout("cached price", 50, TimeUnit.MILLISECONDS)
.join();
IO.println("completeOnTimeout: " + price);
IO.println("slow lookups that have ended: " + (2 - slowTaskEnded.getCount()));
neverOpens.countDown(); // let them finish, or close() would wait forever
}
IO.println("slow lookups that have ended: " + (2 - slowTaskEnded.getCount()));
}
Imprime:
orTimeout: TimeoutException
completeOnTimeout: cached price
slow lookups that have ended: 0
slow lookups that have ended: 2
orTimeout hizo fallar el future con TimeoutException, que join() envolvió en CompletionException. completeOnTimeout dio el precio en caché.
Ninguno de los dos detuvo el trabajo. Después de los dos timeouts, no había terminado ninguna búsqueda lenta: sus hilos seguían bloqueados. main tuvo que abrir el latch, o close() habría esperado para siempre. cancel(true) tampoco ayuda. Su documentación dice que el indicador de interrupción no tiene efecto, y lo comprobamos: la tarea nunca vio una interrupción. Cuando main abrió el latch, las dos búsquedas terminaron, que es la última línea. Un timeout en un CompletableFuture detiene la espera, no la tarea.
Herramientas de coordinación
java.util.concurrent también tiene herramientas pequeñas para que los hilos se esperen entre sí. CountDownLatch, que usó la parte sobre hilos, deja que un hilo espere hasta que otros terminen algo:
void main() throws InterruptedException {
var services = List.of("search", "cache", "database");
var allReady = new CountDownLatch(services.size());
var ready = new ConcurrentSkipListSet<String>();
try (var pool = Executors.newFixedThreadPool(3)) {
for (String service : services) {
pool.execute(() -> {
ready.add(service); // stands in for slow start-up work
allReady.countDown();
});
}
allReady.await();
IO.println("every service is up: " + ready);
}
}
Imprime:
every service is up: [cache, database, search]
El latch empieza en 3, cada servicio hace una cuenta regresiva y await() retorna cuando llega a 0. Un ConcurrentSkipListSet mantiene sus elementos ordenados, así que el orden impreso es fijo.
Un Semaphore guarda una cantidad de permisos. acquire() toma uno, o espera si no queda ninguno, y release() lo devuelve. Limita cuántos hilos hacen algo a la vez, como llamar a un servicio que acepta dos conexiones. Aquí seis tareas comparten dos permisos, y un latch se asegura de que dos estén realmente adentro al mismo tiempo:
void main() throws InterruptedException {
var permits = new Semaphore(2);
var inside = new AtomicInteger();
var mostInside = new AtomicInteger();
var twoInside = new CountDownLatch(2);
var gate = new CountDownLatch(1);
try (var pool = Executors.newFixedThreadPool(6)) {
for (int i = 0; i < 6; i++) {
pool.execute(() -> {
try {
permits.acquire();
try {
mostInside.accumulateAndGet(inside.incrementAndGet(), Math::max);
twoInside.countDown();
gate.await();
inside.decrementAndGet();
} finally {
permits.release();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
});
}
twoInside.await();
while (permits.getQueueLength() < 4) {
Thread.sleep(1);
}
IO.println("holding a permit: " + inside.get());
IO.println("waiting for one: " + permits.getQueueLength());
gate.countDown();
}
IO.println("most inside at once: " + mostInside.get());
}
Imprime:
holding a permit: 2
waiting for one: 4
most inside at once: 2
Las dos primeras tareas tomaron los permisos y esperaron en la compuerta. Las otras cuatro esperaron en acquire(). El máximo es 2. El latch twoInside obligó a dos poseedores a estar adentro juntos, y el semáforo dejó afuera a un tercero. Libera el permiso en finally, igual que con un lock, o una tarea que falla se queda con su permiso para siempre.
Un CyclicBarrier hace que una cantidad fija de hilos espere hasta que todos lleguen al mismo punto, y después los deja seguir a todos juntos.
Colecciones concurrentes, otra vez
La parte sobre hilos mostró que una llamada a ConcurrentHashMap es segura por sí sola, pero comprobar y luego actuar con dos llamadas no lo es. Los métodos compute meten toda la actualización dentro de una sola llamada:
void main() {
var stock = new ConcurrentHashMap<String, Integer>();
stock.put("tea", 20);
stock.put("cake", 3);
var requests = new ConcurrentHashMap<String, Integer>();
var basket = List.of("tea", "cake", "tea", "soup", "cake", "tea");
try (var pool = Executors.newFixedThreadPool(4)) {
for (int i = 0; i < 4; i++) {
pool.execute(() -> {
for (String item : basket) {
requests.merge(item, 1, Integer::sum);
// take one; the last one removes the entry
stock.computeIfPresent(item, (name, left) -> left == 1 ? null : left - 1);
}
});
}
}
IO.println("requests: " + new TreeMap<>(requests));
IO.println("stock: " + new TreeMap<>(stock));
}
Imprime:
requests: {cake=8, soup=4, tea=12}
stock: {tea=8}
Cuatro tareas hicieron 24 pedidos. merge los contó sin perder ninguna actualización. computeIfPresent tomó un artículo a la vez, y devolver null eliminó la entrada, así que el pastel se agotó después de tres y desapareció. La sopa nunca estuvo en stock, así que la función nunca se ejecutó para ella.
Una BlockingQueue pasa trabajo de un hilo a otro. put espera mientras la cola está llena, y take espera mientras está vacía:
void main() throws Exception {
var queue = new ArrayBlockingQueue<String>(2);
var received = new ArrayList<String>();
try (var pool = Executors.newFixedThreadPool(2)) {
Future<?> producer = pool.submit(() -> {
for (int i = 1; i <= 5; i++) {
queue.put("order-" + i); // waits while the queue is full
}
queue.put("DONE");
return null;
});
Future<?> consumer = pool.submit(() -> {
String order = queue.take(); // waits while the queue is empty
while (!order.equals("DONE")) {
received.add(order);
order = queue.take();
}
return null;
});
producer.get();
consumer.get();
}
IO.println("consumer received: " + received);
}
Imprime:
consumer received: [order-1, order-2, order-3, order-4, order-5]
Con una capacidad de 2, el productor no puede adelantarse al consumidor por más de dos órdenes. El valor "DONE" le indica al consumidor que pare. Con un productor y un consumidor, las órdenes salen en el orden en que entraron. Con varios de cualquiera de los dos, cada cola sigue siendo segura, pero el orden entre hilos no es fijo.
Las dos tareas son lambdas Callable, ya que terminan con return null. Eso permite que put y take lancen InterruptedException hacia el Future sin un try.
Qué recordar
- Un executor reutiliza un conjunto fijo de hilos y pone en cola las tareas que sobran. Ábrelo con try-with-resources, para que
close()espere a que termine el trabajo. submitdevuelve unFuture. Llama aget()sobre él, o la excepción de una tarea que falla se pierde en silencio.invokeAlldevuelve los futures en el orden de envío.get()lanzaExecutionException, yjoin()lanzaCompletionException. La excepción real es la causa.- Un timeout detiene la espera, no la tarea. La interrupción es una petición, así que restaura el indicador cuando captures
InterruptedException. - Dimensiona un pool para trabajo de CPU cerca de la cantidad de núcleos. Para trabajo de IO, usa la regla general de espera y cómputo como estimación inicial, y después mide.
- Dale a
CompletableFuturetu propio executor. UsathenComposecuando el siguiente paso devuelve un future, yallOfy luegojoinen un orden fijo. - Un
Semaphorelimita cuántos hilos hacen algo a la vez. UnaBlockingQueuepasa trabajo entre hilos en orden.
Entrega las tareas a un pool, guarda cada future y decide de antemano qué pasa cuando una tarea falla o tarda demasiado.