Concurrencia, paralelismo y virtual threads
La concurrencia es donde más se separa un programador junior de uno senior, y no por la sintaxis:
synchronized se aprende en diez minutos. Lo difícil es razonar sobre lo que puede
pasar cuando dos hilos miran la misma memoria, y saber qué garantiza exactamente la JVM. Este módulo va de
eso: primero el modelo mental (qué es un hilo, qué cuesta, qué puede reordenar la CPU), después las
herramientas (java.util.concurrent, executors, CompletableFuture), después Loom
(virtual threads y concurrencia estructurada, que han cambiado la forma correcta de escribir servicios en
Java) y por último los patrones, el testing y el diagnóstico que se usan en producción.
ExecutorService y nunca entiende por qué hace falta. Aquí primero
rompemos programas (condiciones de carrera, visibilidad, deadlocks) y luego explicamos el
modelo de memoria que dicta qué está garantizado. Con eso, el resto son herramientas
obvias en lugar de recetas mágicas. Los ejemplos se ejecutan con java Fichero.java en un JDK 21+.
1 · Fundamentos: concurrencia, paralelismo y el coste de un hilo
1.1 Concurrencia no es paralelismo
Son conceptos ortogonales que se confunden constantemente, incluso en entrevistas senior:
- Concurrencia es una propiedad del diseño: el programa está estructurado en tareas independientes que progresan de forma solapada. Puede ocurrir con un solo núcleo.
- Paralelismo es una propiedad de la ejecución: dos instrucciones se ejecutan literalmente en el mismo instante. Requiere más de un núcleo.
UN NÚCLEO — concurrencia sin paralelismo (el SO reparte el tiempo)
Núcleo 0: [A··][B··][A··][C··][B··][A··] A, B y C progresan solapados
DOS NÚCLEOS — concurrencia CON paralelismo
Núcleo 0: [A·············][C·········]
Núcleo 1: [B·········][A·····][B·····]
TAREAS DE E/S — el caso real del backend (el hilo NO consume CPU mientras espera)
Hilo 1: calcula ▓▓ ─── espera respuesta de la BD ─────────── ▓▓ responde
└ 2 ms ┘ └────────── 40 ms dormido ─────────┘ └ 1 ms ┘
43 ms de trabajo total, 3 ms de CPU: el 95% del tiempo el hilo solo OCUPA MEMORIA.
De ahí todo el sentido de los virtual threads (sección 8).
Esta distinción decide la arquitectura. En un backend típico (REST + base de datos + llamadas a otros servicios) el 90 % de las tareas son de espera, no de cómputo. Por eso la solución histórica fue “muchos hilos bloqueados” (que cuesta memoria), luego “pocos hilos y programación asíncrona” (que cuesta legibilidad) y ahora “muchísimos hilos baratísimos”.
1.2 Por qué la concurrencia es difícil de verdad
- Pierdes el determinismo. El entrelazado lo elige el planificador del SO y crece factorialmente: dos hilos con 10 instrucciones tienen 184.756 entrelazados posibles. Tu test explora uno.
- Lo que escribes no es lo que se ejecuta. Compilador, JIT y CPU reordenan siempre que un hilo aislado no note la diferencia. Otro hilo sí la nota (sección 4).
- No hay “memoria compartida” real. Cada núcleo tiene sus cachés L1/L2: escribir en un campo no implica que otro núcleo lo vea.
- Los errores no son locales. Surgen de la interacción entre dos trozos de código correctos por separado, a veces en equipos distintos.
- Son heisenbugs. Añadir un
printlno un depurador cambia el timing y el bug desaparece. Se manifiestan bajo carga, en producción, un viernes.
ConcurrentHashMap, AtomicLong) y solo en último
lugar (5) escribir tu propia sincronización. El 95 % del código concurrente correcto vive en los niveles 1–4.
1.3 Ley de Amdahl, throughput y latencia
1
Aceleración = ───────────────── (Ley de Amdahl)
(1 - p) + p / N
p = fracción paralelizable, N = procesadores
Con p = 0,95 (¡solo un 5% secuencial!):
N = 2 → 1,9x N = 32 → 12,5x
N = 8 → 5,9x N = 256 → 17,4x
N = ∞ → 20,0x ← techo absoluto: 1 / (1 - 0,95)
Con p = 0,75: techo = 4x, por muchos núcleos que compres.
Y en la práctica es PEOR que la fórmula, porque paralelizar añade coste propio:
sincronización, contención, tráfico de coherencia de caché, reparto y recolección.
Existe un N óptimo tras el cual añadir hilos EMPEORA el resultado.
| Métrica | Qué mide | Cómo se mejora | Trampa habitual |
|---|---|---|---|
| Latencia | Tiempo de una operación (p50, p95, p99) | Menos trabajo, cachés, paralelizar dentro de la petición, no esperar en serie | Mirar la media: oculta la cola. Mide percentiles |
| Throughput | Operaciones por segundo del sistema entero | Más paralelismo, lotes, menos contención, pools bien dimensionados | Subirlo a costa del p99: colas llenas = latencia enorme |
1.4 Hilos de plataforma, hilos del SO y cuánto cuesta un hilo
Un java.lang.Thread clásico (hoy llamado hilo de plataforma) es una capa fina
sobre un hilo del sistema operativo: relación 1:1. Todo lo que cuesta un hilo del kernel
lo pagas tú.
| Coste | Magnitud típica (Linux x86-64, HotSpot) | Consecuencia práctica |
|---|---|---|
| Stack reservado | ≈ 1 MB por hilo (-Xss; se reserva memoria virtual y se compromete por páginas) | 10.000 hilos ≈ 10 GB de direcciones reservadas. En un contenedor con límite de memoria, muerte |
| Creación | ≈ 0,5–1 ms (syscall clone + estructuras del kernel) | Nunca un hilo por tarea corta: de ahí los pools |
| Cambio de contexto | ≈ 1–10 µs, más el coste indirecto de vaciar cachés y TLB | Con miles de hilos activos, el SO cambia de contexto más que trabaja |
| Estructuras del kernel | Unos KB por hilo, no reclamables | Límite práctico: unos pocos miles de hilos por JVM |
| Presión sobre el GC | Cada stack es una raíz de GC que hay que escanear | Pausas más largas con muchos hilos |
# Cuántos hilos tiene realmente tu proceso Java (Linux)
ls /proc/$(pgrep -f mi-app.jar)/task | wc -l
# Límites del sistema (si los superas: OutOfMemoryError: unable to create native thread)
ulimit -u
cat /proc/sys/kernel/threads-max
# Tamaño de stack por defecto de tu JVM
java -XX:+PrintFlagsFinal -version | grep ThreadStackSize
# intx ThreadStackSize = 1024 ← en KB, o sea 1 MB
// Demostración del coste: crea hilos de plataforma hasta que la JVM se rinde.
// Ejecuta con -Xmx256m y observa "OutOfMemoryError: unable to create native thread".
public class CuantosHilos {
public static void main(String[] args) {
var contador = new java.util.concurrent.atomic.AtomicInteger();
try {
while (true) {
Thread t = new Thread(() -> {
contador.incrementAndGet();
try { Thread.sleep(Long.MAX_VALUE); } // el hilo se queda esperando
catch (InterruptedException e) { Thread.currentThread().interrupt(); }
});
t.setDaemon(true); // daemon: no impide que la JVM termine
t.start();
}
} catch (Throwable e) {
System.out.println("Murió con " + contador.get() + " hilos: " + e);
// Típicamente entre 2.000 y 30.000 según memoria, -Xss y límites del SO.
}
}
}
1.5 CPU-bound vs IO-bound: dimensionar un pool con números
| CPU-bound | IO-bound | |
|---|---|---|
| Qué hace | Calcula: compresión, cifrado, parsing, imágenes, ML, ordenación | Espera: BD, HTTP, disco, colas, otro microservicio |
| Cómo se detecta | CPU al 100 %, poca espera en los volcados de hilos | CPU baja y latencia alta; hilos en WAITING/TIMED_WAITING o en SocketRead |
| Tamaño de pool | ≈ número de núcleos (o núcleos + 1) | Muchos más que núcleos; con la fórmula de abajo o, mejor, virtual threads |
| ¿Ayudan los virtual threads? | No. No hay CPU que ganar | Sí, muchísimo. Es su caso exacto |
| Cuello de botella real | Núcleos disponibles | El recurso remoto: conexiones de BD, cuota del proveedor, red |
FÓRMULA 1 — Tamaño de pool (Brian Goetz, "Java Concurrency in Practice")
N_hilos = N_núcleos × U_objetivo × (1 + T_espera / T_cómputo)
N_núcleos : Runtime.getRuntime().availableProcessors()
U_objetivo : utilización de CPU deseada, 0..1 (usa 0,8 para dejar aire al GC)
T_espera : tiempo bloqueado por tarea
T_cómputo : tiempo de CPU por tarea
Ejemplo A — endpoint que consulta la BD:
8 núcleos, U = 0,8, 2 ms de CPU y 40 ms de espera
N = 8 × 0,8 × (1 + 40/2) = 8 × 0,8 × 21 ≈ 134 hilos
Ejemplo B — servicio que redimensiona imágenes (CPU pura):
8 núcleos, U = 1, T_espera = 0
N = 8 × 1 × (1 + 0) = 8 hilos ← más hilos = más cambios de contexto, no más trabajo
FÓRMULA 2 — Ley de Little (para validar el número con datos de producción)
L = λ × W
L : trabajos simultáneos en el sistema → la concurrencia que necesitas
λ : tasa de llegada (peticiones/segundo)
W : tiempo de permanencia de cada petición (segundos)
Ejemplo: 500 req/s con latencia media de 200 ms
L = 500 × 0,2 = 100 peticiones simultáneas
→ necesitas ~100 hilos concurrentes (o 100 permisos de semáforo con virtual threads).
Si tu pool tiene 20, las 80 restantes esperan en la cola: la latencia se dispara y la
propia ley te lo devuelve como W creciente. Es la firma inequívoca de un pool pequeño.
Uso inverso, muy útil: con 200 hilos y 50 ms de servicio, el techo teórico de throughput
es λ = L / W = 200 / 0,05 = 4.000 req/s. Si mides 800, el cuello de botella NO son los
hilos: búscalo en la BD, en un lock o en el GC.
2 · Hilos clásicos: Thread, ciclo de vida e interrupción
2.1 Thread, Runnable y formas de crear un hilo
// 1) Implementar Runnable (preferido: separa QUÉ se hace de CÓMO se ejecuta)
Runnable tarea = () -> System.out.println("Hola desde " + Thread.currentThread().getName());
Thread t1 = new Thread(tarea, "trabajador-1"); // ✅ SIEMPRE con nombre: lo agradecerás en un jstack
t1.start(); // start() crea el hilo del SO y ejecuta run()
// t1.run(); // ❌ ERROR CLÁSICO: ejecuta en el hilo ACTUAL
// 2) Extender Thread (evítalo: gastas la única herencia y acoplas tarea y mecanismo)
class MiHilo extends Thread {
@Override public void run() { /* ... */ }
}
// 3) Java 21+: builders explícitos
Thread plataforma = Thread.ofPlatform()
.name("importador-", 0) // importador-0, importador-1, …
.daemon(false)
.priority(Thread.NORM_PRIORITY)
.unstarted(tarea); // creado pero no arrancado
plataforma.start();
Thread virtual = Thread.ofVirtual().name("peticion-", 0).start(tarea); // sección 8
// 4) Lo que harás el 99% de las veces: NO crear hilos a mano
try (ExecutorService pool = Executors.newFixedThreadPool(8)) { // Java 19+: AutoCloseable
Future<Integer> futuro = pool.submit(() -> calcular()); // Callable<Integer>
System.out.println(futuro.get());
} // close() = shutdown() + espera a que terminen las tareas
// Callable vs Runnable: la diferencia que preguntan en toda entrevista
Runnable r = () -> { /* void, no puede lanzar excepciones comprobadas */ };
Callable<String> c = () -> { /* devuelve valor Y puede lanzar Exception */ return "ok"; };
Runnable y no extends Thread: separar la tarea del
mecanismo de ejecución. La misma Runnable la ejecutas en un hilo, en un pool, en un
virtual thread o síncronamente en un test. Es el mismo principio que inyectar interfaces en lugar de
instanciar clases (módulo 01).
2.2 Ciclo de vida y estados
new Thread(tarea)
│
▼
┌─────────────────┐
│ NEW │ creado, aún no arrancado
└────────┬────────┘
start()
▼
┌──────────────────► ┌─────────────────┐ ◄─────────────────────┐
│ │ RUNNABLE │ ejecutándose O listo │
│ │ (la JVM NO │ para ejecutarse │
│ │ distingue │ │
│ │ "en CPU") │ │
│ └──┬───┬───┬──────┘ │
│ adquiere el monitor │ │ │ notifyAll(), unpark(), │
│ │ │ │ fin del join │
│ │ │ └───────────────────────────────┤
┌────┴──────────┐ │ │ wait(), join(), park() ┌────┴─────────┐
│ BLOCKED │◄───────────┘ └──────────────────────────────►│ WAITING │
│ esperando un │ entra en un │ espera │
│ monitor │ synchronized ocupado │ INDEFINIDA │
└───────┬───────┘ └──────┬───────┘
│ sleep(ms), wait(ms), join(ms), │
│ parkNanos(), tryLock(timeout) │
│ ┌──────────────────┐ │
└─────────────────►│ TIMED_WAITING │◄────────────────────────┘
│ espera con plazo │
└────────┬─────────┘
│ vence el plazo o lo despiertan
▼
(vuelve a RUNNABLE)
│
run() termina o lanza una excepción
▼
┌─────────────────┐
│ TERMINATED │ irreversible: volver a llamar a
└─────────────────┘ start() lanza IllegalThreadStateException
| Estado | Qué significa | Cómo se llega | Qué buscar en un volcado |
|---|---|---|---|
NEW | Objeto creado, hilo del SO todavía inexistente | new Thread(...) | No aparece |
RUNNABLE | Ejecutando o esperando CPU. Ojo: también aparece así un hilo bloqueado en E/S de socket | start() | Muchos RUNNABLE en SocketRead = espera de red, no CPU |
BLOCKED | Esperando entrar en un synchronized ocupado | Contención de monitor | waiting to lock <0x…> + locked <0x…> en otro hilo |
WAITING | Espera indefinida a que otro hilo actúe | wait(), join(), park(), latch.await() | parking to wait for / in Object.wait() |
TIMED_WAITING | Igual, con plazo máximo | sleep(n), poll(timeout), tryLock(timeout) | Pools inactivos: normal y sano |
TERMINATED | Terminado (normal o por excepción) | Fin de run() | No aparece |
BLOCKED solo se refiere a monitores de
synchronized. Un hilo esperando un ReentrantLock aparece como
WAITING/TIMED_WAITING (usa LockSupport.park), y uno esperando datos
de un socket aparece como RUNNABLE. Lo que importa es la traza de pila, no la
etiqueta.
2.3 join: esperar a que otro hilo termine
public class JoinDemo {
public static void main(String[] args) throws InterruptedException {
var resultados = new java.util.concurrent.ConcurrentHashMap<String, Integer>();
Thread a = new Thread(() -> resultados.put("a", trabajo(300)), "hilo-a");
Thread b = new Thread(() -> resultados.put("b", trabajo(500)), "hilo-b");
a.start();
b.start();
// ❌ SIN join: main puede leer el mapa antes de que a y b hayan escrito
// System.out.println(resultados); // podría imprimir {} o {a=300}
a.join(); // espera INDEFINIDAMENTE
b.join(2_000); // ✅ mejor: como máximo 2 s
if (b.isAlive()) { // join(ms) NO dice si terminó: hay que preguntarlo
b.interrupt(); // pide la cancelación (sección 2.4)
b.join(500);
}
// join() establece happens-before: todo lo que 'a' escribió ANTES de terminar es
// visible aquí DESPUÉS del join. Sin join no hay ninguna garantía.
System.out.println(resultados);
}
static int trabajo(int ms) {
try { Thread.sleep(ms); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
return ms;
}
}
2.4 Interrupción cooperativa: la parte que casi todos hacen mal
Java no puede matar un hilo. Lo único que existe es un mecanismo cooperativo: se activa un flag en el hilo destino y es su propio código quien decide si lo respeta. Si no lo comprueba, no se cancela nunca.
CÓMO FUNCIONA REALMENTE
otroHilo.interrupt()
│
├── ¿está DENTRO de un método bloqueante que responde a interrupción?
│ (sleep, wait, join, take, put, await, lockInterruptibly, park…)
│ SÍ → el método LIMPIA el flag (queda false) y lanza InterruptedException
│
└── ¿está ejecutando código normal (bucle, cálculo, E/S de socket clásica)?
SÍ → solo se ACTIVA el flag; nada más ocurre. Hay que consultarlo:
Thread.currentThread().isInterrupted() ← consulta sin limpiar
Thread.interrupted() ← ¡CONSULTA Y LIMPIA!
REGLA CENTRAL: si capturas InterruptedException y no la propagas, DEBES restaurar el
flag con Thread.currentThread().interrupt(). Si no, destruyes la información de
cancelación para todo el código que hay por encima en la pila.
// ❌❌❌ EL PECADO CAPITAL: tragarse la interrupción. Aparece en millones de líneas y es
// un bug grave: el hilo sigue vivo, el pool no se puede cerrar y el shutdown se cuelga.
void malisimo() {
try { Thread.sleep(1000); }
catch (InterruptedException e) {
// "no pasa nada" ← el flag se ha limpiado y nadie sabrá que hubo cancelación
}
}
// ❌ Casi igual de malo: registrar y seguir como si nada
void malo() {
try { cola.take(); }
catch (InterruptedException e) { log.error("ups", e); } // sigue el bucle infinito
}
// ✅ OPCIÓN 1 (la mejor): PROPAGAR. Declara la excepción y deja decidir a quien llama.
String leerConEspera(BlockingQueue<String> cola) throws InterruptedException {
return cola.take();
}
// ✅ OPCIÓN 2: si NO puedes propagar (implementas Runnable.run(), firma fija), restaura
// el flag y termina de forma ordenada.
@Override public void run() {
try {
while (!Thread.currentThread().isInterrupted()) {
String mensaje = cola.take(); // punto de cancelación
procesar(mensaje);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt(); // ✅ RESTAURA el flag
log.info("Consumidor interrumpido; cerrando ordenadamente");
} finally {
liberarRecursos(); // el finally SÍ se ejecuta
}
}
// ✅ OPCIÓN 3: bucle de cómputo largo sin llamadas bloqueantes → comprueba el flag tú
long sumaCarisima(long hasta) {
long suma = 0;
for (long i = 0; i < hasta; i++) {
if ((i & 0xFFFF) == 0 && Thread.currentThread().isInterrupted()) {
// cada 65.536 iteraciones: barato y con buena reactividad
throw new CancellationException("Cálculo cancelado en i=" + i);
}
suma += trabajo(i);
}
return suma;
}
// ⚠️ La E/S de socket clásica NO responde a interrupt(): un hilo bloqueado leyendo de un
// InputStream sigue bloqueado. La única salida es CERRAR el socket (o poner timeouts).
// Los InterruptibleChannel de NIO sí responden: lanzan ClosedByInterruptException.
InterruptedException: tu servicio no se apaga.
Kubernetes envía SIGTERM, Spring llama a shutdown() del pool, el pool interrumpe
a sus hilos, tu catch vacío se los come, awaitTermination agota su plazo y a los
30 segundos llega SIGKILL con peticiones a medio procesar. Los despliegues “que pierden
tráfico” tienen esta causa muchísimo más a menudo de lo que la gente cree.
2.5 Hilos daemon, nombres y excepciones no capturadas
// DAEMON: la JVM termina cuando solo quedan hilos daemon. Para tareas de fondo accesorias.
Thread metricas = new Thread(this::publicarMetricas, "metricas");
metricas.setDaemon(true); // ANTES de start(); si no, IllegalThreadStateException
metricas.start();
// ⚠️ Un daemon se mata en seco al terminar la JVM: sus finally NO se ejecutan de forma
// fiable. Nunca pongas en un daemon algo que deba completarse (escribir en BD, flush).
// NOMBRAR HILOS: obligatorio en producción. Un volcado con "pool-1-thread-7" no dice nada;
// con "pedidos-worker-7" localizas el problema en cinco segundos.
ThreadFactory fabrica = Thread.ofPlatform().name("pedidos-worker-", 0).factory();
ExecutorService pool = Executors.newFixedThreadPool(8, fabrica);
// UncaughtExceptionHandler: si run() lanza una excepción no capturada, el hilo MUERE.
// Por defecto se imprime en System.err… y en un contenedor con logs JSON nadie lo verá.
Thread.setDefaultUncaughtExceptionHandler((hilo, error) -> {
log.error("Excepción no capturada en el hilo {}", hilo.getName(), error);
if (error instanceof OutOfMemoryError) Runtime.getRuntime().halt(1); // morir es lo correcto
});
// Handler por hilo concreto (tiene prioridad sobre el global)
Thread t = new Thread(tarea, "importador");
t.setUncaughtExceptionHandler((h, e) -> log.error("El importador ha fallado", e));
t.start();
// ⚠️ TRAMPA MUY PREGUNTADA: el handler NO se aplica a submit() en un ExecutorService.
pool.execute(() -> { throw new RuntimeException("visible: va al handler"); });
pool.submit (() -> { throw new RuntimeException("INVISIBLE: se queda en el Future"); });
| Aspecto | Hilo normal (user) | Hilo daemon |
|---|---|---|
| ¿Impide que la JVM termine? | Sí | No |
¿Se ejecutan sus finally al cerrar la JVM? | Sí (termina él solo) | No garantizado |
| Uso típico | Trabajo de negocio, servidores de peticiones | Métricas, limpieza, watchdogs |
| Virtual threads | Son siempre daemon y su prioridad no se puede cambiar | |
2.6 Por qué Thread.stop() ya no existe
Thread.stop() lanzaba un ThreadDeath asíncrono en un punto arbitrario del hilo
destino: podía interrumpir a mitad de la actualización de dos campos relacionados, liberando los monitores
y dejando el objeto en un estado inconsistente pero visible para los demás. Es
irremediablemente inseguro: no se puede escribir código robusto frente a “me pueden lanzar una excepción en
cualquier bytecode”.
// Lo que podía pasar, y por qué era imposible defenderse:
synchronized void transferir(Cuenta destino, long importe) {
this.saldo -= importe; // ← si te matan AQUÍ, el dinero desaparece del sistema
destino.saldo += importe; // y el monitor se libera dejando el saldo mal
}
// Historia de la API (dilo así en la entrevista):
// Java 1.2 : stop(), suspend(), resume() marcados deprecated
// Java 18 : deprecated for removal
// Java 20 : stop(), suspend(), resume() lanzan UnsupportedOperationException
// Después : degradados/eliminados del API; ThreadDeath también desapareció
//
// La ÚNICA cancelación correcta hoy: interrupción cooperativa comprobando el flag,
// cerrar el recurso que bloquea (socket, canal), Future.cancel(true) o StructuredTaskScope.
Checklist — hilos e interrupción
3 · Los problemas clásicos, con código que puedes romper hoy
3.1 Condición de carrera: contador++ son tres operaciones
public class CarreraContador {
static int contador = 0; // estado mutable COMPARTIDO: aquí empieza el problema
public static void main(String[] args) throws InterruptedException {
int hilos = 4, incrementosPorHilo = 100_000;
var lista = new java.util.ArrayList<Thread>();
for (int i = 0; i < hilos; i++) {
Thread t = new Thread(() -> {
for (int j = 0; j < incrementosPorHilo; j++) contador++; // ❌ NO es atómico
});
lista.add(t);
t.start();
}
for (Thread t : lista) t.join();
System.out.println("Esperado: " + hilos * incrementosPorHilo); // 400000
System.out.println("Real: " + contador); // casi siempre MENOR. Pruébalo 5 veces.
}
}
POR QUÉ FALLA: contador++ se compila a TRES operaciones de bytecode
getstatic contador ← 1. LEER el valor actual
iconst_1 / iadd ← 2. SUMAR uno (sobre una COPIA, en la pila de operandos)
putstatic contador ← 3. ESCRIBIR el resultado
Entrelazado desastroso con contador = 41:
tiempo Hilo A Hilo B contador
─────────────────────────────────────────────────────────────────────────────
t1 lee 41 41
t2 lee 41 41
t3 suma → 42 (en su pila) 41
t4 suma → 42 (en su pila) 41
t5 escribe 42 42
t6 escribe 42 42 ⚠️
─────────────────────────────────────────────────────────────────────────────
Dos incrementos, una sola unidad de progreso: se ha PERDIDO una actualización
("lost update"). Con 400.000 iteraciones esto ocurre miles de veces.
Es una condición de carrera "read-modify-write" o "check-then-act", exactamente
el mismo bug que:
if (!mapa.containsKey(k)) mapa.put(k, v); ← check-then-act
if (instancia == null) instancia = new X(); ← singleton roto
if (saldo >= importe) saldo -= importe; ← descubierto en el banco
Cuatro formas correctas de arreglarlo, de la más rápida a la más general:
// ✅ 1. Atómico con CAS: la mejor opción para un contador (sección 5.1)
static final AtomicInteger contadorAtomico = new AtomicInteger();
contadorAtomico.incrementAndGet();
// ✅ 1b. LongAdder: aún mejor con contención altísima (muchos hilos escribiendo)
static final LongAdder contadorAdder = new LongAdder();
contadorAdder.increment(); // se lee con .sum()
// ✅ 2. synchronized: mutua exclusión + visibilidad
static int contadorSync = 0;
static synchronized void incrementar() { contadorSync++; } // el monitor es la clase
// ✅ 3. Lock explícito: cuando necesitas tryLock, timeout o varias Condition
static final ReentrantLock lock = new ReentrantLock();
static void incrementarConLock() {
lock.lock();
try { contadorSync++; }
finally { lock.unlock(); } // ⚠️ el unlock SIEMPRE en finally
}
// ✅ 4. No compartir: cada hilo lleva su cuenta y se suman al final (la más escalable)
long total = IntStream.range(0, 4).parallel()
.mapToLong(h -> { long local = 0; for (int j = 0; j < 100_000; j++) local++; return local; })
.sum();
// ❌ volatile NO lo arregla: da visibilidad, NO atomicidad. Siguen siendo 3 operaciones.
static volatile int contadorVolatile = 0;
contadorVolatile++; // sigue perdiendo actualizaciones
3.2 Visibilidad: el bucle que nunca termina
public class VisibilidadRota {
// ❌ sin volatile: el hilo trabajador puede NO ver nunca el cambio
static boolean parar = false;
public static void main(String[] args) throws InterruptedException {
Thread trabajador = new Thread(() -> {
int vueltas = 0;
while (!parar) { // el JIT puede "izar" la lectura fuera del bucle:
vueltas++; // if (!parar) while (true) vueltas++;
}
System.out.println("Parado tras " + vueltas + " vueltas");
});
trabajador.start();
Thread.sleep(1000);
parar = true; // main escribe…
System.out.println("He pedido parar");
trabajador.join(3000);
System.out.println("¿Sigue vivo? " + trabajador.isAlive()); // MUY a menudo: true 💥
}
}
// En muchas máquinas el programa NO TERMINA. No es un bug de la JVM: sin volatile no hay
// ninguna garantía de que el trabajador vea la escritura de main.
POR QUÉ OCURRE — dos causas independientes, ambas perfectamente legales
1) CACHÉS DE CPU: cada núcleo trabaja sobre su propia copia
Núcleo 0 (main) Núcleo 1 (trabajador)
┌──────────────┐ ┌──────────────┐
│ L1: parar=T │ │ L1: parar=F │ ← lee su copia, eternamente
└──────┬───────┘ └──────┬───────┘
│ (la escritura puede quedarse en el store buffer)
▼ ▼
╔══════════════════════ RAM: parar = ? ══════════════════════╗
2) OPTIMIZACIÓN DEL JIT (aún más frecuente): el compilador ve que en ese hilo nadie
modifica 'parar', así que hace "hoisting" (saca la lectura del bucle) y lo convierte
en un bucle infinito. Es una transformación CORRECTA según el JMM, porque sin volatile
no existe relación happens-before entre los dos hilos.
SOLUCIONES CORRECTAS
static volatile boolean parar = false; ← lo mínimo, y suficiente aquí
static final AtomicBoolean parar = new AtomicBoolean(); ← si además necesitas CAS
Thread.currentThread().isInterrupted() ← el flag de interrupción YA es visible
escribir y leer dentro del MISMO synchronized ← también da visibilidad
volatile convierte una variable normal en un punto de
sincronización: la escritura hace un flush de lo pendiente y la lectura invalida la caché, de modo
que se establece happens-before entre el hilo que escribe y el que lee. Y además prohíbe al
compilador cachear la variable en un registro.”
3.3 Publicación insegura de objetos
// ❌ PUBLICACIÓN INSEGURA: otro hilo puede ver la referencia YA asignada pero los campos
// del objeto A MEDIO CONSTRUIR.
public class Holder {
private int n;
public Holder(int n) { this.n = n; }
public void comprobar() {
if (n != n) throw new AssertionError("¡n != n!"); // parece imposible… puede pasar
}
}
public class PublicacionInsegura {
public static Holder holder; // ❌ ni volatile, ni final, ni sincronizado
static void publicar() { holder = new Holder(42); } // hilo A
static void usar() {
if (holder != null) holder.comprobar(); // hilo B
}
}
/* El constructor NO es atómico. La secuencia real puede ser:
1. reservar memoria
2. escribir la referencia en 'holder' ← ¡reordenado antes del paso 3!
3. inicializar el campo n = 42
El hilo B, entre 2 y 3, ve un objeto con n = 0. Y como no hay barrera, puede ver n = 0
en una lectura y n = 42 en la siguiente: de ahí el absurdo "n != n". */
// ✅ FORMAS CORRECTAS DE PUBLICAR (publicación segura)
public static volatile Holder seguro1; // volatile
public static final Holder seguro2 = new Holder(42); // final estático
private static final AtomicReference<Holder> seguro3 = new AtomicReference<>();
public static final Map<String, Holder> seguro4 = new ConcurrentHashMap<>(); // colección concurrente
// …o publicar desde dentro de un bloque synchronized que el lector también use.
// ✅ Y la más robusta: hacer el objeto INMUTABLE (todos los campos final). Un objeto
// inmutable correctamente construido se puede publicar de CUALQUIER forma.
public record Configuracion(String url, int intentos) { }
3.4 Deadlock: dos locks, orden distinto
public class DeadlockClasico {
static final Object LOCK_A = new Object();
static final Object LOCK_B = new Object();
public static void main(String[] args) {
new Thread(() -> {
synchronized (LOCK_A) { // 1) coge A
dormir(50);
System.out.println("hilo-1 tiene A, quiere B");
synchronized (LOCK_B) { System.out.println("hilo-1 nunca llega aquí"); }
}
}, "hilo-1").start();
new Thread(() -> {
synchronized (LOCK_B) { // 1) coge B ← ORDEN INVERSO
dormir(50);
System.out.println("hilo-2 tiene B, quiere A");
synchronized (LOCK_A) { System.out.println("hilo-2 nunca llega aquí"); }
}
}, "hilo-2").start();
}
static void dormir(long ms) {
try { Thread.sleep(ms); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
}
}
ESPERA CIRCULAR — las cuatro condiciones de Coffman (deben darse TODAS)
hilo-1 ──tiene──► LOCK_A ◄──quiere── hilo-2
│ ▲
└──quiere──► LOCK_B ──tiene──────────┘
1. Exclusión mutua : el recurso solo lo puede tener uno
2. Retener y esperar : mantienes un lock mientras pides otro
3. Sin preferencia : nadie puede quitarte el lock a la fuerza
4. Espera circular : el ciclo de arriba
Rompe UNA cualquiera y el deadlock es imposible. En la práctica se rompe la 4
(ordenación global de locks) o la 2 (tryLock con timeout).
CASO REAL FRECUENTE: transferencias bancarias
transferir(cuentaA, cuentaB) bloquea A y luego B
transferir(cuentaB, cuentaA) bloquea B y luego A ← simultáneas: deadlock
Solo aparece bajo carga, con dos transferencias cruzadas a la vez.
Detectarlo: jstack lo dice con nombres y apellidos
jps -l # localiza el PID
jstack -l <pid> > hilos.txt # -l incluye los "ownable synchronizers"
jcmd <pid> Thread.print -l # equivalente y preferido en contenedores
grep -A 30 "Found one Java-level deadlock" hilos.txt
Found one Java-level deadlock:
=============================
"hilo-2":
waiting to lock monitor 0x00007f8c1c0062a8 (object 0x000000071ab2c1f8, a java.lang.Object),
which is held by "hilo-1"
"hilo-1":
waiting to lock monitor 0x00007f8c1c008d18 (object 0x000000071ab2c208, a java.lang.Object),
which is held by "hilo-2"
Java stack information for the threads listed above:
"hilo-2":
at DeadlockClasico.lambda$main$1(DeadlockClasico.java:20)
- waiting to lock <0x000000071ab2c1f8> (a java.lang.Object)
- locked <0x000000071ab2c208> (a java.lang.Object) ← tiene B, quiere A
"hilo-1":
at DeadlockClasico.lambda$main$0(DeadlockClasico.java:12)
- waiting to lock <0x000000071ab2c208> (a java.lang.Object)
- locked <0x000000071ab2c1f8> (a java.lang.Object) ← tiene A, quiere B
Found 1 deadlock.
⚠️ La detección automática solo encuentra ciclos de monitores (synchronized) y de Lock
del j.u.c. NO detecta "deadlocks lógicos": un latch que nadie cuenta, un pool agotado
o una espera de red eterna. Para esos hay que leer las trazas a mano.
Evitarlo: ordenación global de la adquisición
// ❌ Vulnerable: el orden depende de los argumentos
void transferirMal(Cuenta origen, Cuenta destino, long importe) {
synchronized (origen) {
synchronized (destino) { origen.restar(importe); destino.sumar(importe); }
}
}
// ✅ SOLUCIÓN 1: ORDEN GLOBAL DETERMINISTA. Todos adquieren en el mismo orden.
void transferir(Cuenta origen, Cuenta destino, long importe) {
if (origen.id() == destino.id()) throw new IllegalArgumentException("misma cuenta");
// Se ordena por una clave estable y única (el id, no el hashCode, que puede repetirse)
Cuenta primera = origen.id() < destino.id() ? origen : destino;
Cuenta segunda = origen.id() < destino.id() ? destino : origen;
synchronized (primera) {
synchronized (segunda) {
if (origen.saldo() < importe) throw new SaldoInsuficienteException(origen.id());
origen.restar(importe);
destino.sumar(importe);
}
}
}
// ✅ SOLUCIÓN 2: tryLock con timeout y reintento (rompe "retener y esperar")
boolean transferirConTimeout(CuentaLock a, CuentaLock b, long importe) throws InterruptedException {
while (true) {
if (a.lock().tryLock(100, TimeUnit.MILLISECONDS)) {
try {
if (b.lock().tryLock(100, TimeUnit.MILLISECONDS)) {
try { a.restar(importe); b.sumar(importe); return true; }
finally { b.lock().unlock(); }
}
} finally { a.lock().unlock(); }
}
// No lo hemos conseguido: soltamos TODO y reintentamos con espera aleatoria.
// El jitter evita que los dos hilos vuelvan a chocar en bucle (livelock).
Thread.sleep(ThreadLocalRandom.current().nextInt(10, 60));
}
}
// ✅ SOLUCIÓN 3 (la mejor): no anidar locks. Un solo lock de grano más grueso, o una
// operación atómica en la base de datos dentro de una transacción, delegando el
// problema al gestor de bloqueos del SGBD.
@Transactional
void transferirEnBd(long idOrigen, long idDestino, BigDecimal importe) {
long primero = Math.min(idOrigen, idDestino), segundo = Math.max(idOrigen, idDestino);
repo.bloquearParaActualizar(primero); // SELECT … FOR UPDATE, en orden estable
repo.bloquearParaActualizar(segundo);
repo.mover(idOrigen, idDestino, importe);
}
ThreadMXBean bean = ManagementFactory.getThreadMXBean(); y
long[] ids = bean.findDeadlockedThreads(); devuelve los hilos en deadlock o
null. Exponerlo como métrica y alertar sobre ella es barato y ha salvado más de un fin de
semana.
3.5 Livelock y starvation
| Patología | Qué ocurre | Síntoma observable | Solución |
|---|---|---|---|
| Deadlock | Todos esperan y nadie avanza | CPU al 0 %, hilos en BLOCKED, todo congelado |
Orden global de locks, tryLock con timeout, no anidar |
| Livelock | Todos actúan, se estorban y ninguno progresa | CPU al 100 % y cero trabajo útil; reintentos infinitos | Backoff aleatorio y exponencial, límite de reintentos |
| Starvation | Un hilo nunca obtiene el recurso porque otros se lo llevan siempre | p99 desastroso con p50 normal; una cola que no se vacía nunca | Locks fair, colas FIFO, no retener locks mucho tiempo |
// LIVELOCK: el "pasillo estrecho". Dos hilos ceden educadamente a la vez, para siempre.
class Livelock {
static volatile boolean turnoDeA = true;
static void hiloA() {
while (true) {
if (!turnoDeA) { Thread.onSpinWait(); continue; } // "cedo el paso"
hacerTrabajo();
turnoDeA = false;
}
}
// Con dos hilos simétricos que se ceden el turno tras cada intento, pueden pasarse la
// vida cediendo. La CPU está al 100%: parece que trabaja. No trabaja.
// ✅ Arreglo: introducir ASIMETRÍA con aleatoriedad
static void reintentarBien(Runnable operacion) throws InterruptedException {
for (int intento = 0; intento < 5; intento++) {
if (intentar(operacion)) return;
long base = (long) (50 * Math.pow(2, intento)); // 50,100,200,400,800
long jitter = ThreadLocalRandom.current().nextLong(base / 2 + 1); // desincroniza
Thread.sleep(base + jitter);
}
throw new IllegalStateException("Agotados los reintentos");
}
}
// STARVATION por lock injusto: por defecto ReentrantLock NO es fair (permite "barging")
ReentrantLock injusto = new ReentrantLock(); // rápido, pero un hilo puede pasar hambre
ReentrantLock justo = new ReentrantLock(true); // FIFO garantizado, 10-100× más lento
// Usa fair=true SOLO con evidencia de inanición real: el coste es enorme.
// STARVATION por prioridades: Thread.setPriority() es una SUGERENCIA que la mayoría de los
// sistemas operativos ignora o interpreta a su manera. No construyas nada sobre ella.
3.6 Thread starvation deadlock: el pool que se bloquea a sí mismo
// ❌ BUG DEVASTADOR: una tarea del pool espera el resultado de OTRA tarea del MISMO pool.
ExecutorService pool = Executors.newFixedThreadPool(2); // solo 2 hilos
Future<String> padre = pool.submit(() -> {
Future<String> hija1 = pool.submit(() -> "parte-1");
Future<String> hija2 = pool.submit(() -> "parte-2");
return hija1.get() + hija2.get(); // ⛔ BLOQUEA su hilo esperando a las hijas
});
// Si dos tareas padre entran a la vez: los 2 hilos están ocupados esperando y las 4 hijas
// están en la cola sin ningún hilo libre. NADIE avanza. La JVM NO lo detecta como
// deadlock (no hay ciclo de monitores): se ve como un "cuelgue" con CPU al 0% y una cola
// que crece. Es dificilísimo de diagnosticar si no conoces el patrón.
// ✅ SOLUCIÓN 1: pools separados por tipo de trabajo (bulkhead)
ExecutorService poolPadres = Executors.newFixedThreadPool(4, fabrica("padre-"));
ExecutorService poolHijas = Executors.newFixedThreadPool(8, fabrica("hija-"));
// ✅ SOLUCIÓN 2: no bloquear; componer de forma asíncrona (sección 7)
CompletableFuture<String> p1 = CompletableFuture.supplyAsync(() -> "parte-1", poolHijas);
CompletableFuture<String> p2 = CompletableFuture.supplyAsync(() -> "parte-2", poolHijas);
CompletableFuture<String> combinado = p1.thenCombine(p2, String::concat); // sin bloquear
// ✅ SOLUCIÓN 3: ForkJoinPool, que sabe hacer esto: cuando una tarea hace join(), su hilo
// roba y ejecuta otras tareas en vez de dormirse (work-stealing + ManagedBlocker).
// ✅ SOLUCIÓN 4 (2026): virtual threads. Bloquear un virtual thread es baratísimo, así que
// "espero a mis subtareas" deja de ser un problema estructural.
try (var exec = Executors.newVirtualThreadPerTaskExecutor()) {
Future<String> f = exec.submit(() -> {
Future<String> h1 = exec.submit(() -> "parte-1");
Future<String> h2 = exec.submit(() -> "parte-2");
return h1.get() + h2.get(); // no hay límite de hilos que agotar
});
}
parallelStream()
dentro de una tarea del commonPool, un join() dentro de un
thenApplyAsync del commonPool, o un @Async que espera a otro
@Async del mismo executor.
4 · El modelo de memoria de Java (JMM)
El Java Memory Model (JLS capítulo 17, definido por la JSR-133 en 2004) es un contrato. No describe cómo funciona una CPU concreta: describe qué está garantizado en cualquier JVM sobre cualquier hardware. Sin él, “funciona en mi x86” no significaría nada al desplegar en ARM, que reordena mucho más agresivamente. Entenderlo es la diferencia entre escribir código concurrente y escribir código concurrente que casualmente funciona.
4.1 Reordenamientos y cachés: lo que escribes no es lo que ocurre
TRES NIVELES DE REORDENAMIENTO, todos legales si un hilo aislado no nota la diferencia
Tu código → 1. COMPILADOR (javac + JIT)
a = 1; reordena, elimina, mete en registros, "iza" lecturas fuera de bucles
b = 2; → 2. PROCESADOR
ejecución fuera de orden, especulación, prefetch
→ 3. SISTEMA DE MEMORIA
store buffers y caché por núcleo: una escritura puede tardar en
hacerse visible a otros núcleos
JERARQUÍA DE MEMORIA (y el porqué de todo esto)
┌──────────┐ ┌──────────┐ registros : < 1 ciclo
│ Núcleo 0 │ │ Núcleo 1 │ caché L1 : ~4 ciclos (32-64 KB, por núcleo)
│ ┌──────┐ │ │ ┌──────┐ │ caché L2 : ~12 ciclos (por núcleo)
│ │ L1 │ │ │ │ L1 │ │ caché L3 : ~40 ciclos (COMPARTIDA)
│ └──┬───┘ │ │ └──┬───┘ │ RAM : ~200 ciclos ← 50× más lenta que L1
│ ┌──┴───┐ │ │ ┌──┴───┐ │
│ │ L2 │ │ │ │ L2 │ │ Por eso el hardware NO va a RAM en cada acceso: sin
│ └──┬───┘ │ │ └──┬───┘ │ barreras explícitas, tus escrituras pueden quedarse
└────┼─────┘ └────┼─────┘ en L1 o en el store buffer un buen rato.
└──── L3 ─────┘
│
RAM
EL EJEMPLO CANÓNICO (test de Dekker)
int x = 0, y = 0; int a = 0, b = 0;
Hilo 1: x = 1; Hilo 2: y = 1;
a = y; b = x;
Resultados "intuitivos": (a=0,b=1) (a=1,b=0) (a=1,b=1)
Resultado REAL y observable en x86: (a=0, b=0) ← ¡ambos leen 0!
Porque cada CPU puede retrasar su escritura (store buffer) y adelantar su lectura.
Declarando x e y volatile, (0,0) pasa a ser imposible.
4.2 happens-before: la única garantía que existe
happens-before no significa “antes en el tiempo”. Significa: si la acción A
happens-before la acción B, entonces todo lo que A escribió en memoria es visible para B, y ningún
reordenamiento puede alterar ese orden observado. Es una relación de orden parcial y
transitiva. Si entre dos accesos a la misma variable (con al menos una escritura)
no hay una cadena happens-before, tienes una data race y el resultado es
indefinido. Punto.
| Regla | Qué establece |
|---|---|
| Orden del programa | Dentro de un hilo, cada acción happens-before las siguientes del código. Por eso un hilo nunca se ve a sí mismo desordenado. |
| Monitor | Un unlock happens-before cualquier lock posterior del mismo monitor. |
volatile | Una escritura volatile happens-before cualquier lectura posterior de esa misma variable. |
Thread.start() | Todo lo hecho antes de t.start() es visible dentro de t. |
Thread.join() | Todo lo hecho por t es visible tras t.join(). |
| Interrupción | t.interrupt() happens-before que t detecte la interrupción. |
Campos final | La inicialización de campos final en el constructor es visible para quien vea la referencia (con el matiz de no publicar this desde el constructor). |
java.util.concurrent | Todas las clases documentan sus garantías: meter en una BlockingQueue hb sacarlo; countDown() hb el retorno de await(); enviar a un executor hb ejecutar la tarea; completar un Future hb el retorno de get(). |
| Transitividad | Si A hb B y B hb C, entonces A hb C. Esta es la regla que usas todo el rato sin darte cuenta. |
// LA TRANSITIVIDAD EN ACCIÓN: el "truco del flag volatile"
class PublicacionPorFlag {
private int a, b, c; // ¡campos NORMALES, sin volatile!
private volatile boolean listo; // un único volatile como "puerta"
// Hilo escritor
void preparar() {
a = 1; b = 2; c = 3; // (1) escrituras normales
listo = true; // (2) escritura volatile ← BARRERA de liberación
}
// Hilo lector
int leer() {
if (!listo) return -1; // (3) lectura volatile ← BARRERA de adquisición
return a + b + c; // (4) ve GARANTIZADAMENTE 1, 2 y 3
}
}
/* ¿Por qué funciona? Cadena happens-before:
(1) hb (2) por orden del programa
(2) hb (3) por la regla de volatile
(3) hb (4) por orden del programa
→ por TRANSITIVIDAD, (1) hb (4). Las escrituras normales quedan cubiertas.
Este es el mecanismo del double-checked locking correcto (4.6) y de la publicación
segura: un solo volatile bien colocado protege muchos campos. */
// ⚠️ El orden importa y no es simétrico:
// listo = true; a = 1; ← ❌ INÚTIL: la escritura de 'a' va DESPUÉS de la barrera
// return a; if (!listo) … ← ❌ INÚTIL: la lectura de 'a' va ANTES de la barrera
4.3 volatile: qué garantiza y qué no
volatile SÍ garantiza | volatile NO garantiza |
|---|---|
| Visibilidad: toda lectura ve la última escritura completada | Atomicidad de read-modify-write: v++, v += 2, v = v * 2 siguen siendo carreras |
| No reordenamiento alrededor del acceso (barreras de memoria) | Mutua exclusión: no bloquea nada, no hay sección crítica |
Atomicidad de long/double: sin él podría haber word tearing en VM de 32 bits |
Coherencia entre varias variables: dos campos volatile pueden verse en estados incompatibles entre sí |
| Publicación segura del objeto referenciado y de todo lo escrito antes | Protección del objeto apuntado: volatile List protege la referencia, jamás el contenido |
// ✅ CASOS DE USO LEGÍTIMOS (son pocos y muy concretos)
// 1. Flag de parada
private volatile boolean cerrado = false;
public void cerrar() { cerrado = true; }
public void comprobar() { if (cerrado) throw new IllegalStateException("cerrado"); }
// 2. Publicación de una referencia inmutable que se reemplaza entera
private volatile Configuracion config = Configuracion.inicial();
public void recargar(Configuracion nueva) { this.config = nueva; } // swap atómico
public Configuracion config() { return config; } // lectores sin bloqueo
// 3. "Barrera" para varios campos normales (el patrón de 4.2)
// 4. Double-checked locking (4.6): aquí volatile es OBLIGATORIO
// ❌ ANTIPATRONES CON volatile
private volatile int contador;
void incrementar() { contador++; } // ❌ carrera: 3 operaciones
private volatile List<String> lista = new ArrayList<>();
void anadir(String s) { lista.add(s); } // ❌ la LISTA no es thread-safe
private volatile long saldo;
void retirar(long x) { if (saldo >= x) saldo -= x; } // ❌ check-then-act
private volatile int min, max;
void rango(int a, int b) { min = a; max = b; } // ❌ un lector puede ver min nuevo y max viejo
// ✅ Equivalencias correctas
private final AtomicInteger contadorOk = new AtomicInteger();
private final List<String> listaOk = new CopyOnWriteArrayList<>();
private final AtomicLong saldoOk = new AtomicLong(); // + CAS en bucle
private volatile Rango rangoOk = new Rango(0, 0); // un objeto inmutable
record Rango(int min, int max) { } // → coherencia garantizada
4.4 synchronized: exclusión mutua y visibilidad
CADA OBJETO EN JAVA TIENE UN MONITOR (intrinsic lock) asociado
synchronized (obj) { → monitorenter obj ← adquirir; si está ocupado: BLOCKED
... ... (barrera de adquisición)
} → monitorexit obj ← liberar (barrera de liberación);
el compilador lo garantiza incluso
si se lanza una excepción
DOS GARANTÍAS EN UNA (esto es lo que hay que decir en la entrevista):
1. EXCLUSIÓN MUTUA : solo un hilo dentro del bloque a la vez → atomicidad compuesta
2. VISIBILIDAD : al entrar ves todo lo que hizo el hilo que salió del MISMO monitor
⚠️ Solo hay exclusión si TODOS los accesos usan el MISMO objeto como monitor. Sincronizar
la escritura y no la lectura equivale a no haber sincronizado nada.
REENTRANCIA: el monitor cuenta las adquisiciones de su propietario.
synchronized void a() { b(); } ← a() ya tiene el lock…
synchronized void b() { } ← …y b() lo readquiere: contador 2. Correcto.
(Si no fuera reentrante, llamar a un método sincronizado desde otro del mismo objeto
sería un deadlock instantáneo.)
public class VariantesDeSynchronized {
// A) MÉTODO DE INSTANCIA → el monitor es 'this'
public synchronized void metodoInstancia() { }
public void metodoInstanciaExpandido() { synchronized (this) { } } // equivalente exacto
// B) MÉTODO ESTÁTICO → el monitor es VariantesDeSynchronized.class (¡otro monitor!)
public static synchronized void metodoEstatico() { }
public static void metodoEstaticoExpandido() {
synchronized (VariantesDeSynchronized.class) { }
}
// ⚠️ Un método de instancia sincronizado y uno estático sincronizado NO se excluyen
// entre sí: usan monitores diferentes. Fuente inagotable de bugs.
// C) BLOQUE → el monitor lo eliges tú. PREFERIBLE: sección crítica mínima.
private final Map<String, Integer> datos = new HashMap<>();
private final Object candado = new Object(); // ✅ lock PRIVADO y dedicado
public void procesar(String clave) {
Datos externos = llamadaHttpLenta(clave); // ✅ FUERA del lock: 200 ms sin bloquear
synchronized (candado) {
datos.merge(clave, externos.valor(), Integer::sum); // ✅ dentro: lo imprescindible
}
publicarEvento(clave); // ✅ fuera: no llames a código ajeno con un lock
}
// ❌ MALAS ELECCIONES DE MONITOR
// synchronized (this) → this es PÚBLICO: cualquiera puede bloquearlo y
// crear un deadlock que tú no puedes ver ni arreglar
// synchronized ("clave") → los literales String están internados: TODO el
// classpath que use "clave" comparte tu monitor 💥
// synchronized (Boolean.TRUE) → instancia cacheada y global: mismo problema
// synchronized (Integer.valueOf(1)) → caché de -128..127: monitor global de facto
// synchronized (nuevoObjetoCadaVez) → cada hilo bloquea SU objeto: no excluye nada
// synchronized (campoNoFinal) → si la referencia cambia, dos hilos usan monitores
// distintos y la exclusión desaparece
}
synchronized: desde Java 6, un monitor sin contención es
prácticamente gratis (bloqueo ligero con CAS). Lo caro es la contención: cuando hay
competencia real el monitor se “infla” y los hilos van a dormir al kernel. Optimizar concurrencia casi
nunca es “quitar synchronized”; es reducir el tiempo dentro del lock y
partir el estado para que los hilos no coincidan (lo que hacen
ConcurrentHashMap y LongAdder). Nota histórica: el biased locking se
deshabilitó por defecto en Java 15 y se eliminó después; en 2026 no cuentes con él.
4.5 final, atomicidad y publicación segura
// SEMÁNTICA ESPECIAL DE final (JSR-133): si un objeto solo tiene campos final y no publica
// 'this' durante la construcción, cualquier hilo que vea la referencia verá TODOS los
// campos correctamente inicializados, sin volatile ni locks.
public final class PuntoInmutable {
private final int x, y;
public PuntoInmutable(int x, int y) { this.x = x; this.y = y; }
public int x() { return x; }
public int y() { return y; }
}
// Por eso los objetos inmutables son la herramienta número uno de la concurrencia: se
// comparten libremente, sin sincronización de ningún tipo.
// ❌ ROMPER final publicando 'this' en el constructor (error real y sutilísimo)
public class Escuchador {
private final int umbral;
public Escuchador(EventBus bus, int umbral) {
bus.registrar(this); // ⛔ 'this' escapa ANTES de terminar de construirse:
this.umbral = umbral; // otro hilo puede ver umbral == 0
}
// ✅ Solución: factoría estática que construye primero y registra después
public static Escuchador crear(EventBus bus, int umbral) {
Escuchador e = new Escuchador(umbral); // constructor privado, sin fugas
bus.registrar(e);
return e;
}
}
// ❌ Otra fuga clásica: arrancar un hilo desde el constructor
public class ServicioMalo {
private final String nombre;
public ServicioMalo(String nombre) {
new Thread(this::bucle).start(); // ⛔ el hilo puede leer nombre == null
this.nombre = nombre;
}
private void bucle() { System.out.println(nombre); }
}
// ATOMICIDAD POR TIPO — lo que la JVM garantiza sin ayuda
// · Lecturas y escrituras de tipos de 32 bits o menos (int, boolean, char, float y
// referencias) son ATÓMICAS: nunca ves "medio valor".
// · long y double NO están garantizados como atómicos sin volatile (word tearing en VM
// de 32 bits). En las de 64 bits lo son en la práctica, pero el CONTRATO es lo que
// importa: si los compartes, ponles volatile o usa AtomicLong.
// · Atómico ≠ thread-safe: leer un int es atómico, pero contador++ no lo es.
4.6 Inicialización perezosa: DCL correcto y por qué el holder es mejor
// ❌ VERSIÓN 1: perezosa y rota (check-then-act)
class SingletonRoto {
private static SingletonRoto instancia;
static SingletonRoto get() {
if (instancia == null) instancia = new SingletonRoto(); // dos hilos → dos instancias
return instancia;
}
}
// 😐 VERSIÓN 2: correcta, pero con lock en CADA lectura
class SingletonSincronizado {
private static SingletonSincronizado instancia;
static synchronized SingletonSincronizado get() {
if (instancia == null) instancia = new SingletonSincronizado();
return instancia; // se sincroniza un millón de veces
} // aunque ya esté creado
}
// ❌ VERSIÓN 3: double-checked locking SIN volatile — EL BUG MÁS FAMOSO DE JAVA
class DclRoto {
private static DclRoto instancia; // ⛔ falta volatile
static DclRoto get() {
if (instancia == null) { // 1.ª comprobación, sin lock
synchronized (DclRoto.class) {
if (instancia == null) instancia = new DclRoto(); // 2.ª comprobación
}
}
return instancia;
}
}
/* Por qué falla: instancia = new DclRoto() no es una operación, son tres:
(a) reservar memoria
(b) ejecutar el constructor
(c) asignar la referencia al campo
El compilador puede reordenar a (a)(c)(b). Otro hilo que ejecute la PRIMERA
comprobación (fuera del synchronized, sin barrera) puede ver instancia != null y
devolver un objeto A MEDIO CONSTRUIR. Antes de la JSR-133 (Java 5) este patrón era
irreparable; con volatile sí funciona. */
// ✅ VERSIÓN 4: DCL CORRECTO (volatile obligatorio)
class DclCorrecto {
private static volatile DclCorrecto instancia; // ✅ volatile: barreras de memoria
static DclCorrecto get() {
DclCorrecto local = instancia; // ✅ una sola lectura volatile
if (local == null) { // camino rápido, sin lock
synchronized (DclCorrecto.class) {
local = instancia;
if (local == null) instancia = local = new DclCorrecto();
}
}
return local;
}
}
// ✅✅ VERSIÓN 5: HOLDER IDIOM — la mejor. Cero sincronización explícita, cero volatile.
class SingletonHolder {
private SingletonHolder() { /* inicialización costosa */ }
private static final class Holder { // no se carga hasta que se usa
static final SingletonHolder INSTANCIA = new SingletonHolder();
}
static SingletonHolder get() { return Holder.INSTANCIA; }
}
/* Funciona porque la ESPECIFICACIÓN de la JVM garantiza que la inicialización de una clase
(<clinit>) se ejecuta UNA sola vez y es thread-safe: la propia JVM sincroniza la carga
de clases. Y es perezosa: Holder no se inicializa hasta el primer get(). Ventajas sobre
DCL: menos código, imposible de escribir mal y sin coste de lectura volatile. */
// ✅✅ VERSIÓN 6: enum singleton (Effective Java, ítem 3) — sin pereza, pero infalible
enum SingletonEnum {
INSTANCIA;
private final Map<String, String> datos = cargar();
private Map<String, String> cargar() { return Map.of(); }
}
// ✅ Y para memoización perezosa de MUCHOS valores (no un singleton): ConcurrentHashMap
private final ConcurrentHashMap<String, Resultado> cache = new ConcurrentHashMap<>();
Resultado calcular(String clave) {
return cache.computeIfAbsent(clave, this::calculoCostoso); // atómico por clave
}
Checklist — modelo de memoria
5 · Herramientas de java.util.concurrent
Doug Lea escribió este paquete (Java 5, JSR-166) precisamente para que no tengas que escribir tú la
sincronización. La regla es simple: si existe una clase en j.u.c. que hace lo que
necesitas, úsala. Está más probada, es más rápida y documenta sus garantías de memoria.
5.1 Variables atómicas y CAS
COMPARE-AND-SWAP (CAS): la instrucción de hardware sobre la que se construye TODO j.u.c.
En x86 es LOCK CMPXCHG; en ARM, LDXR/STXR. Es atómica a nivel de procesador.
boolean compareAndSet(esperado, nuevo):
"si el valor actual es EXACTAMENTE 'esperado', ponlo a 'nuevo' y devuelve true;
si no, no toques nada y devuelve false"
CÓMO SE IMPLEMENTA incrementAndGet() — bucle optimista, SIN bloqueos:
┌─────────────────────────────────────────────────────┐
│ 1. actual = leer valor │
│ 2. nuevo = actual + 1 (cálculo sin lock) │
│ 3. ¿compareAndSet(actual, nuevo)? │
│ true → listo, devuelve nuevo ───────────────┼──► FIN
│ false → otro hilo se me adelantó │
│ vuelve al paso 1 ──────┐ │
└───────────────────────────────────────┴─────────────┘
Hilo A: lee 41 → calcula 42 → CAS(41,42) ✅ éxito contador = 42
Hilo B: lee 41 → calcula 42 → CAS(41,42) ❌ FALLA (ya vale 42)
→ REINTENTA: lee 42 → calcula 43 → CAS(42,43) ✅ contador = 43
Ningún incremento se pierde y ningún hilo se ha ido a dormir.
BLOQUEANTE (lock) vs NO BLOQUEANTE (CAS)
pesimista: "asumo conflicto" optimista: "asumo que no habrá conflicto"
el hilo perdedor DUERME el hilo perdedor REINTENTA (gira)
cambio de contexto (µs) solo ciclos de CPU (ns)
gana con secciones críticas LARGAS gana con operaciones CORTÍSIMAS
contención alta = colas ordenadas contención alta = reintentos quemando CPU
import java.util.concurrent.atomic.*;
// --- Contadores ---
AtomicInteger contador = new AtomicInteger(0);
contador.incrementAndGet(); // ++c (devuelve el nuevo)
contador.getAndIncrement(); // c++ (devuelve el anterior)
contador.addAndGet(5);
int anterior = contador.getAndSet(0); // leer y reiniciar de forma atómica
boolean cambiado = contador.compareAndSet(10, 20); // CAS explícito
// --- Actualización con una función arbitraria: el bucle CAS ya escrito ---
contador.updateAndGet(v -> Math.min(v + 1, 100)); // ✅ tope atómico
contador.accumulateAndGet(7, Math::max); // ✅ máximo atómico
// ⚠️ La lambda DEBE ser pura y sin efectos laterales: puede ejecutarse VARIAS veces (una
// por reintento). Nunca metas dentro un log, una llamada HTTP ni una escritura.
// --- Referencias atómicas: swap atómico de objetos inmutables ---
record Estado(int intentos, String ultimoError) { }
AtomicReference<Estado> estado = new AtomicReference<>(new Estado(0, null));
estado.updateAndGet(e -> new Estado(e.intentos() + 1, "timeout")); // dos campos coherentes
// --- AtomicBoolean: idempotencia y "hazlo solo una vez" ---
private final AtomicBoolean iniciado = new AtomicBoolean(false);
public void iniciar() {
if (!iniciado.compareAndSet(false, true)) return; // ✅ solo el primer hilo entra
arrancarRecursos();
}
// --- AtomicLong para IDs; arrays atómicos elemento a elemento ---
AtomicLong secuencia = new AtomicLong();
long id = secuencia.incrementAndGet();
AtomicIntegerArray cubetas = new AtomicIntegerArray(16);
cubetas.incrementAndGet(3);
// EL PROBLEMA ABA — pregunta de entrevista senior
//
// Hilo A: lee el valor "A" … se queda suspendido …
// Hilo B: cambia A → B, y después B → A otra vez
// Hilo A: despierta, hace CAS(A, C) y TIENE ÉXITO
// …pero el mundo cambió en medio y A no se ha enterado.
//
// Con un contador de enteros es inofensivo. Es peligroso con estructuras enlazadas
// (pilas/colas lock-free): el nodo "A" puede haber sido eliminado, reciclado y reinsertado;
// el CAS pasa y la estructura queda corrupta.
//
// ✅ Solución: añadir un SELLO (versión) que solo crece. Si el sello cambió, el CAS falla.
AtomicStampedReference<Nodo> cima = new AtomicStampedReference<>(nodoInicial, 0);
int[] selloLeido = new int[1];
Nodo actual = cima.get(selloLeido); // lee valor Y sello
boolean ok = cima.compareAndSet(actual, nuevoNodo, selloLeido[0], selloLeido[0] + 1);
// También existe AtomicMarkableReference (referencia + un bit), útil para marcar nodos
// como "borrados lógicamente" en listas concurrentes.
// VarHandle (Java 9+): el mecanismo moderno de bajo nivel para CAS sobre campos propios.
// Sustituye a sun.misc.Unsafe y a los ...FieldUpdater. Úsalo solo si escribes estructuras
// de datos; en código de negocio, Atomic* es suficiente.
class Nodo {
volatile Nodo siguiente;
private static final VarHandle SIGUIENTE;
static {
try {
SIGUIENTE = MethodHandles.lookup().findVarHandle(Nodo.class, "siguiente", Nodo.class);
} catch (ReflectiveOperationException e) { throw new ExceptionInInitializerError(e); }
}
boolean casSiguiente(Nodo esperado, Nodo nuevo) {
return SIGUIENTE.compareAndSet(this, esperado, nuevo);
}
}
5.2 LongAdder y el false sharing
// Problema de AtomicLong con contención MUY alta (32 hilos contando peticiones): todos
// hacen CAS sobre LA MISMA dirección. La línea de caché rebota entre núcleos y la mayoría
// de los CAS fallan.
// ✅ LongAdder: reparte el contador en varias celdas, una por hilo (aproximadamente).
LongAdder peticiones = new LongAdder();
peticiones.increment(); // escribe en SU celda: casi nunca hay conflicto
peticiones.add(5);
long total = peticiones.sum(); // suma todas las celdas (no es un instante exacto)
peticiones.sumThenReset();
// Cuándo cada uno:
// AtomicLong → necesitas el valor exacto en cada operación (IDs, CAS lógico)
// LongAdder → solo acumulas y lees de vez en cuando: MÉTRICAS. Hasta 10× más rápido
// LongAccumulator → como LongAdder con una función asociativa arbitraria (máx, mín, or)
LongAccumulator maximo = new LongAccumulator(Long::max, Long.MIN_VALUE);
maximo.accumulate(latenciaMs);
FALSE SHARING — por qué dos variables independientes se pelean
Las CPU no mueven bytes: mueven LÍNEAS DE CACHÉ de 64 bytes.
Línea de caché (64 bytes)
┌──────────────────────────────────────────────────────────┐
│ contadorA (8B) │ contadorB (8B) │ … resto de la línea … │
└──────────────────────────────────────────────────────────┘
▲ ▲
Núcleo 0 escribe Núcleo 1 escribe
solo contadorA solo contadorB
Aunque las variables son LÓGICAMENTE independientes, comparten línea. Cada escritura
invalida la línea en el otro núcleo (protocolo MESI): la línea viaja de un lado a otro
y el rendimiento se hunde 5-10× sin que haya ningún lock ni ninguna variable compartida
a la vista. Esto es FALSE SHARING.
CÓMO SE MITIGA
· Separar los datos "calientes" de cada hilo (padding de 64/128 bytes entre ellos).
· Que cada hilo acumule en variables LOCALES y solo publique al final.
· Usar LongAdder, que ya lo hace por dentro.
· La anotación @jdk.internal.vm.annotation.Contended existe pero es INTERNA del JDK
(requiere -XX:-RestrictContended y --add-exports): no la uses en código de negocio.
CÓMO SE DETECTA
· perf c2c en Linux, o sospecha si el escalado con más hilos es NEGATIVO en un
algoritmo que no tiene locks.
5.3 Locks explícitos: ReentrantLock, ReadWriteLock, StampedLock
// ReentrantLock: mismo modelo semántico que synchronized, pero programable.
private final ReentrantLock lock = new ReentrantLock(); // fair = false por defecto
void conLock() {
lock.lock();
try { /* sección crítica */ }
finally { lock.unlock(); } // ⚠️ SIEMPRE en finally. Sin esto, un throw dentro
} // deja el lock cogido para siempre.
// ✅ tryLock: no bloquear indefinidamente (la herramienta anti-deadlock)
if (lock.tryLock()) { // intenta sin esperar nada
try { /* ... */ } finally { lock.unlock(); }
} else {
metricas.incrementar("lock.rechazado");
throw new RecursoOcupadoException(); // degradar en lugar de colgarse
}
// ✅ tryLock con timeout (responde a interrupción)
if (lock.tryLock(200, TimeUnit.MILLISECONDS)) {
try { /* ... */ } finally { lock.unlock(); }
} else {
log.warn("No se obtuvo el lock en 200 ms; se aplica el camino alternativo");
}
// ✅ lockInterruptibly: permite cancelar un hilo que espera el lock
// (lock() NO responde a interrupt(); lockInterruptibly() sí)
lock.lockInterruptibly();
// Diagnóstico integrado, muy útil en logs
lock.isLocked(); lock.isHeldByCurrentThread(); lock.getHoldCount(); lock.getQueueLength();
// FAIRNESS
new ReentrantLock(false); // por defecto: "barging" permitido. Máximo throughput.
new ReentrantLock(true); // FIFO estricto: evita starvation, 10-100× menos throughput.
| Criterio | synchronized | ReentrantLock |
|---|---|---|
| Liberación | Automática al salir del bloque (incluso con excepción) | Manual: finally { unlock(); } obligatorio |
tryLock / timeout | No | Sí — la razón principal para usarlo |
| Interrumpible mientras espera | No | Sí (lockInterruptibly) |
| Fairness configurable | No | Sí |
| Varias condiciones de espera | Una sola (wait/notify) | Varias Condition por lock |
| Adquirir en un método y liberar en otro | Imposible | Posible (y peligroso) |
Visibilidad en jstack | Muy clara (BLOCKED, detección de deadlock) | Como WAITING en park; usa jstack -l |
| Virtual threads (antes de Java 24) | Podía fijar el carrier (pinning) | No fija: era la opción recomendada |
| Recomendación | Por defecto: más simple e imposible de olvidar | Cuando necesites timeout, cancelación, fairness o varias condiciones |
// ReadWriteLock: muchos lectores concurrentes, un escritor exclusivo. Rentable SOLO si las
// lecturas son mayoría abrumadora y duran algo (si son de 20 ns, el coste del lock se come
// la ventaja).
public class CacheConRW<K, V> {
private final Map<K, V> mapa = new HashMap<>();
private final ReentrantReadWriteLock rw = new ReentrantReadWriteLock();
private final Lock lectura = rw.readLock();
private final Lock escritura = rw.writeLock();
public V get(K k) {
lectura.lock(); // varios hilos a la vez, sin exclusión entre ellos
try { return mapa.get(k); } finally { lectura.unlock(); }
}
public void put(K k, V v) {
escritura.lock(); // exclusivo: excluye lectores y escritores
try { mapa.put(k, v); } finally { escritura.unlock(); }
}
// ⚠️ NO se puede SUBIR de lectura a escritura (deadlock instantáneo): hay que soltar la
// lectura, coger la escritura y RE-COMPROBAR el estado. Sí se puede DEGRADAR
// (escritura → lectura) cogiendo lectura antes de soltar escritura.
}
// 👉 En la práctica, para un mapa: usa ConcurrentHashMap y olvídate del ReadWriteLock.
// StampedLock (Java 8): añade LECTURA OPTIMISTA sin escribir en memoria. Es lo más rápido
// para estructuras con lecturas masivísimas. NO es reentrante, NO soporta Condition y NO
// es interrumpible: úsalo con cuidado quirúrgico.
public class PuntoStamped {
private double x, y;
private final StampedLock sl = new StampedLock();
public void mover(double dx, double dy) {
long sello = sl.writeLock();
try { x += dx; y += dy; } finally { sl.unlockWrite(sello); }
}
public double distanciaAlOrigen() {
long sello = sl.tryOptimisticRead(); // 1) NO bloquea: solo toma un "sello"
double cx = x, cy = y; // 2) lee los campos a locales
if (!sl.validate(sello)) { // 3) ¿escribió alguien mientras leía?
sello = sl.readLock(); // sí → reintenta con lock real
try { cx = x; cy = y; } finally { sl.unlockRead(sello); }
}
return Math.hypot(cx, cy);
}
// ⚠️ En la lectura optimista los datos pueden ser INCONSISTENTES hasta que validate()
// lo confirme: cópialos a locales y no actúes sobre ellos antes de validar.
}
5.4 Condition frente a wait/notify
// wait/notify (el mecanismo antiguo): funciona, pero es fácil de usar mal.
class BufferAntiguo<T> {
private final Queue<T> cola = new ArrayDeque<>();
private final int capacidad;
BufferAntiguo(int capacidad) { this.capacidad = capacidad; }
public synchronized void poner(T t) throws InterruptedException {
while (cola.size() == capacidad) { // ✅ SIEMPRE while, NUNCA if: hay despertares
wait(); // espurios y otro hilo puede haber llenado
} // el buffer entre el notify y tu reanudación
cola.add(t);
notifyAll(); // ✅ notifyAll, no notify: con notify puedes
} // despertar al hilo equivocado y colgar todo
public synchronized T tomar() throws InterruptedException {
while (cola.isEmpty()) wait();
T t = cola.poll();
notifyAll();
return t;
}
// Reglas que se incumplen constantemente:
// 1. Solo se pueden llamar CON el monitor adquirido (si no: IllegalMonitorStateException)
// 2. wait() LIBERA el monitor mientras espera; sleep() NO lo libera ← trampa de entrevista
// 3. La condición se comprueba en un WHILE, siempre
// 4. notifyAll() por defecto; notify() solo si todos esperan exactamente lo mismo
}
// ✅ Con Condition: varias colas de espera por lock → despiertas exactamente a quien toca
class BufferModerno<T> {
private final Queue<T> cola = new ArrayDeque<>();
private final int capacidad;
private final ReentrantLock lock = new ReentrantLock();
private final Condition noLleno = lock.newCondition(); // esperan los productores
private final Condition noVacio = lock.newCondition(); // esperan los consumidores
BufferModerno(int capacidad) { this.capacidad = capacidad; }
public void poner(T t) throws InterruptedException {
lock.lock();
try {
while (cola.size() == capacidad) noLleno.await(); // await, no wait
cola.add(t);
noVacio.signal(); // despierta a UN consumidor: preciso y eficiente
} finally { lock.unlock(); }
}
public T tomar() throws InterruptedException {
lock.lock();
try {
while (cola.isEmpty()) noVacio.await();
T t = cola.poll();
noLleno.signal();
return t;
} finally { lock.unlock(); }
}
// Condition ofrece además: await(long, TimeUnit), awaitUntil(Date) y
// awaitUninterruptibly() (este último, con muchísimo cuidado).
}
// 👉 Y LA VERSIÓN QUE DEBES ESCRIBIR EN PRODUCCIÓN: ninguna de las dos.
BlockingQueue<T> buffer = new ArrayBlockingQueue<>(1000); // ya está escrito y probado
buffer.put(t); // bloquea si está lleno
T t = buffer.take(); // bloquea si está vacío
wait/notify es conocimiento de lectura
(para entender código antiguo y aprobar entrevistas), no de escritura. En código nuevo:
BlockingQueue, CountDownLatch, Semaphore,
CompletableFuture o concurrencia estructurada. Si te ves escribiendo
notifyAll(), para y pregúntate qué sincronizador de j.u.c. estás reimplementando.
5.5 Sincronizadores: Semaphore, CountDownLatch, CyclicBarrier, Phaser, Exchanger
| Clase | Idea | ¿Reutilizable? | Caso de uso real |
|---|---|---|---|
Semaphore | N permisos; acquire los coge, release los devuelve | Sí | Limitar concurrencia contra un recurso escaso (API externa, disco, BD) |
CountDownLatch | Cuenta que solo baja; await espera a que llegue a 0 | No (un solo uso) | “Espera a que arranquen los 3 componentes”; sincronizar el inicio de un test |
CyclicBarrier | N hilos se esperan mutuamente en un punto | Sí | Simulaciones y cálculos por fases con número fijo de hilos |
Phaser | Barrera con partes dinámicas y fases numeradas | Sí | Pipelines donde entran y salen participantes; jerarquías |
Exchanger | Dos hilos intercambian objetos en un punto de encuentro | Sí | Doble búfer entre un productor y un consumidor |
// --- SEMAPHORE: la herramienta de throttling. Imprescindible con virtual threads ---
// "Como máximo 10 llamadas simultáneas a la API de pagos, aunque haya 10.000 peticiones"
private final Semaphore permisos = new Semaphore(10);
public Respuesta llamarApiPagos(Peticion p) throws InterruptedException {
if (!permisos.tryAcquire(2, TimeUnit.SECONDS)) { // ✅ con timeout, nunca acquire() a pelo
throw new ServicioSaturadoException("Cola de pagos llena"); // fail fast > cola infinita
}
try {
return clienteHttp.enviar(p);
} finally {
permisos.release(); // ⚠️ SIEMPRE en finally
}
}
// permisos.availablePermits() → métrica excelente para exponer en /actuator/metrics
// new Semaphore(1) es un mutex que además se puede liberar desde OTRO hilo (a diferencia
// de un lock, que solo libera su propietario). Útil, y también peligroso.
// --- COUNTDOWNLATCH: esperar a que N cosas ocurran (una sola vez) ---
CountDownLatch listos = new CountDownLatch(3);
for (var componente : List.of(cache, colaMensajes, poolBd)) {
Thread.ofVirtual().start(() -> {
try { componente.iniciar(); } finally { listos.countDown(); } // ⚠️ en finally: si
}); // falla, no cuelgues a nadie
}
if (!listos.await(30, TimeUnit.SECONDS)) { // ✅ con timeout, no await() a secas
throw new IllegalStateException("Arranque incompleto en 30 s");
}
// Doble latch: el patrón estándar para tests de concurrencia (sección 12)
CountDownLatch salida = new CountDownLatch(1); // pistoletazo de salida
CountDownLatch fin = new CountDownLatch(nHilos); // meta
// cada hilo: salida.await(); trabajo(); fin.countDown();
// el test: salida.countDown(); fin.await(5, SECONDS); ← máxima colisión real
// --- CYCLICBARRIER: N hilos avanzan por fases sincronizadas ---
CyclicBarrier barrera = new CyclicBarrier(4, () -> log.info("Fase completada"));
// En cada hilo: calcularMiParte(); barrera.await(); // espera a los otros 3
// ⚠️ Si un hilo falla o se interrumpe, la barrera se ROMPE y los demás reciben
// BrokenBarrierException. Hay que reponerla con barrera.reset().
// --- PHASER: como CyclicBarrier pero con participantes dinámicos ---
Phaser fase = new Phaser(1); // 1 = el hilo coordinador
for (var tarea : tareas) {
fase.register(); // se apunta un participante MÁS
ejecutor.execute(() -> {
try { tarea.run(); } finally { fase.arriveAndDeregister(); }
});
}
fase.arriveAndAwaitAdvance(); // el coordinador espera la fase 0
log.info("Fase actual: {}", fase.getPhase());
// --- EXCHANGER: intercambio de búferes entre dos hilos ---
Exchanger<List<Registro>> intercambiador = new Exchanger<>();
// Productor: lleno = intercambiador.exchange(lleno); // entrega lleno, recibe vacío
// Consumidor: vacio = intercambiador.exchange(vacio); // entrega vacío, recibe lleno
5.6 Colecciones concurrentes
| Colección | Mecanismo | Úsala cuando… | Cuidado con… |
|---|---|---|---|
ConcurrentHashMap | Bloqueo por cubeta + CAS; lecturas sin lock | Es tu mapa por defecto en concurrencia | No admite null; el tamaño es aproximado; no llames a otras operaciones del mismo mapa dentro de computeIfAbsent |
CopyOnWriteArrayList | Cada escritura copia el array entero | Listas pequeñas con lecturas masivas y escrituras rarísimas: listeners, configuración | Escritura O(n): con 10.000 elementos y escrituras frecuentes es catastrófico. El iterador ve una instantánea |
ConcurrentLinkedQueue | Lock-free (Michael–Scott) | Cola FIFO no bloqueante e ilimitada | Sin take() bloqueante; size() es O(n) y aproximado; ilimitada = riesgo de OOM |
ArrayBlockingQueue | Array circular + 1 lock con 2 Condition | Cola acotada productor-consumidor: backpressure gratis | Capacidad fija; un solo lock (productores y consumidores compiten) |
LinkedBlockingQueue | Lista enlazada con dos locks (cabeza y cola) | Más throughput con muchos productores y consumidores | Ilimitada por defecto → pásale siempre capacidad; más basura por nodo |
SynchronousQueue | Capacidad 0: cada put espera un take | Handoff directo: es la cola de newCachedThreadPool | No almacena nada; sin consumidor, el productor se bloquea |
PriorityBlockingQueue | Heap + lock | Tareas con prioridad | Ilimitada; sin FIFO entre iguales; starvation de los de baja prioridad |
DelayQueue | Heap por tiempo de expiración | Reintentos programados, cachés con TTL, planificadores | Los elementos implementan Delayed; solo salen al expirar |
ConcurrentSkipListMap | Skip list | Mapa ordenado y concurrente (NavigableMap) | Más lento que ConcurrentHashMap si no necesitas orden |
LinkedTransferQueue | Lock-free con transfer() | Handoff donde el productor espera a que se consuma | Ilimitada |
// ConcurrentHashMap: las operaciones ATÓMICAS son su razón de existir
ConcurrentHashMap<String, Integer> conteo = new ConcurrentHashMap<>();
// ❌ NO es atómico aunque el mapa sea concurrente: dos operaciones = carrera
if (!conteo.containsKey(k)) conteo.put(k, 1); // check-then-act, otra vez
conteo.put(k, conteo.getOrDefault(k, 0) + 1); // lost update
// ✅ Atómico: una sola llamada
conteo.merge(k, 1, Integer::sum); // contar ocurrencias
conteo.putIfAbsent(k, 0);
conteo.computeIfAbsent(k, clave -> cargarCostoso(clave));// memoización perezosa
conteo.computeIfPresent(k, (clave, v) -> v > 1 ? v - 1 : null); // null = ELIMINA la entrada
conteo.compute(k, (clave, v) -> v == null ? 1 : v + 1);
conteo.replace(k, viejo, nuevo); // CAS sobre el valor
conteo.remove(k, valorEsperado); // borrado condicional
// ⚠️ REGLAS DE ORO de compute*/merge:
// 1. La función se ejecuta CON EL LOCK de la cubeta cogido → hazla CORTA. Nunca metas
// dentro una llamada HTTP, una consulta a la BD ni una espera.
// 2. NO toques el MISMO mapa dentro de la función: la clase lo detecta y lanza
// IllegalStateException("Recursive update"), o puede llegar a bloquearse.
// 3. Puede ejecutarse más de una vez si hay contención → debe ser pura.
// 4. Devolver null en compute/merge ELIMINA la entrada (nunca guarda null).
// Recorridos y agregaciones paralelas (usan el commonPool: cuidado)
conteo.forEach(4, (k, v) -> log.info("{}={}", k, v)); // 4 = umbral de paralelismo
long suma = conteo.reduceValues(1, Integer::sum);
Integer encontrado = conteo.search(1, (k, v) -> v > 100 ? v : null);
long n = conteo.mappingCount(); // ✅ preferible a size(): devuelve long y es aproximado
// Un Set concurrente: no existe ConcurrentHashSet, se construye así
Set<String> visitados = ConcurrentHashMap.newKeySet();
// ⚠️ POR QUÉ NO ADMITE null (pregunta clásica con buena respuesta): en un HashMap,
// get(k) == null es ambiguo ("no existe" o "existe con valor null") y se desambigua con
// containsKey(k). En un mapa CONCURRENTE esa desambiguación es imposible: entre el get y
// el containsKey otro hilo pudo cambiar la entrada. Doug Lea eliminó la ambigüedad
// prohibiendo null. Usa Optional o un objeto centinela.
// ⚠️ Y la trampa que sigue existiendo: la colección es thread-safe, el objeto guardado NO
Map<String, List<String>> porCliente = new ConcurrentHashMap<>();
porCliente.computeIfAbsent(k, x -> new ArrayList<>()).add(v); // ❌ ArrayList no es seguro
porCliente.computeIfAbsent(k, x -> new CopyOnWriteArrayList<>()).add(v); // ✅
// Alternativas históricas y por qué no usarlas
Map<String, String> viejo1 = new Hashtable<>(); // ❌ obsoleto, lock global
Map<String, String> viejo2 = Collections.synchronizedMap(new HashMap<>());
// ❌ envuelve TODO en un único lock (mata la concurrencia) y la iteración NO es segura sin
// sincronizar externamente sobre el propio mapa. ConcurrentHashMap gana siempre.
5.7 Productor-consumidor completo con BlockingQueue
import java.util.concurrent.*;
import java.util.List;
/**
* Pipeline productor-consumidor de producción: cola ACOTADA (backpressure), varios
* consumidores, apagado ordenado con "poison pill" y contabilidad de errores.
*/
public class Pipeline {
// Centinela para señalar el fin. Instancia única e identificable por referencia.
private static final String FIN = new String("##FIN##");
private final BlockingQueue<String> cola;
private final ExecutorService consumidores;
private final int nConsumidores;
public Pipeline(int capacidad, int nConsumidores) {
// ✅ ACOTADA: si los consumidores no siguen el ritmo, el productor se bloquea en
// put(). Ese bloqueo ES el backpressure. Con una cola ilimitada, en su lugar te
// comes la memoria hasta el OutOfMemoryError.
this.cola = new ArrayBlockingQueue<>(capacidad);
this.nConsumidores = nConsumidores;
this.consumidores = Executors.newFixedThreadPool(
nConsumidores, Thread.ofPlatform().name("consumidor-", 0).factory());
}
public void arrancar() {
for (int i = 0; i < nConsumidores; i++) consumidores.execute(this::bucleConsumidor);
}
private void bucleConsumidor() {
try {
while (true) {
String item = cola.take(); // bloquea si está vacía, sin quemar CPU
if (item == FIN) { // comparación por IDENTIDAD, a propósito
cola.put(FIN); // reenvía la píldora al siguiente consumidor
return;
}
try {
procesar(item);
} catch (RuntimeException e) {
// ✅ Un item malo NO puede matar al consumidor: aísla el error
log.error("Fallo procesando '{}'; se descarta", item, e);
metricas.incrementar("pipeline.errores");
}
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt(); // ✅ restaurar el flag
log.info("Consumidor interrumpido");
}
}
/** Productor: put() bloqueante = backpressure. */
public void publicar(String item) throws InterruptedException {
cola.put(item);
}
/** Variante con descarte controlado en lugar de bloqueo. */
public boolean publicarSiCabe(String item) throws InterruptedException {
boolean aceptado = cola.offer(item, 100, TimeUnit.MILLISECONDS);
if (!aceptado) metricas.incrementar("pipeline.descartados");
return aceptado;
}
public void cerrar() throws InterruptedException {
cola.put(FIN); // un consumidor la toma y la reenvía
consumidores.shutdown();
if (!consumidores.awaitTermination(30, TimeUnit.SECONDS)) {
log.warn("Consumidores sin terminar en 30 s; forzando");
consumidores.shutdownNow();
}
}
private void procesar(String item) { /* trabajo real */ }
}
| Método | Cola llena (insertar) | Cola vacía (extraer) | Cuándo usarlo |
|---|---|---|---|
add / remove | Lanza IllegalStateException | Lanza NoSuchElementException | Casi nunca en concurrencia |
offer / poll | Devuelve false | Devuelve null | Cuando quieres descartar o seguir sin esperar |
put / take | Bloquea | Bloquea | El caso normal: backpressure natural |
offer(t,u) / poll(t,u) | Bloquea con plazo, luego false | Bloquea con plazo, luego null | ✅ Producción: nunca esperar para siempre |
5.8 ThreadLocal: usos legítimos y fugas garantizadas
// ThreadLocal = una variable con un valor DISTINTO por hilo. Es "confinamiento por hilo"
// como técnica de thread-safety: si nadie comparte, nadie necesita sincronizar.
// ✅ USO LEGÍTIMO 1: reutilizar un objeto costoso y NO thread-safe
private static final ThreadLocal<SimpleDateFormat> FORMATO =
ThreadLocal.withInitial(() -> new SimpleDateFormat("yyyy-MM-dd"));
// (Con java.time no hace falta: DateTimeFormatter ES inmutable y thread-safe. Este era EL
// caso de uso histórico de ThreadLocal y hoy está mayormente obsoleto.)
private static final ThreadLocal<byte[]> BUFFER = ThreadLocal.withInitial(() -> new byte[8192]);
// ✅ USO LEGÍTIMO 2: contexto implícito de la petición (lo que hace medio Spring)
// SecurityContextHolder, RequestContextHolder, TransactionSynchronizationManager y el
// MDC de SLF4J son todos ThreadLocal por debajo.
public final class ContextoPeticion {
private static final ThreadLocal<String> TRACE_ID = new ThreadLocal<>();
public static void establecer(String traceId) { TRACE_ID.set(traceId); }
public static String actual() { return TRACE_ID.get(); }
public static void limpiar() { TRACE_ID.remove(); } // ⚠️ IMPRESCINDIBLE
}
// ✅ El único patrón seguro en un pool: try/finally con remove()
public void manejarPeticion(Peticion p) {
ContextoPeticion.establecer(p.traceId());
try { procesar(p); }
finally { ContextoPeticion.limpiar(); } // ⚠️ SIN esto: fuga de memoria Y de DATOS
}
// ❌ POR QUÉ FUGA EN UN POOL (y por qué es un problema de SEGURIDAD, no solo de memoria):
// 1. Los hilos de un pool VIVEN PARA SIEMPRE; cada uno tiene su ThreadLocalMap.
// 2. Sin remove(), el valor sigue ahí cuando ese hilo atiende la SIGUIENTE petición…
// de OTRO usuario. Resultado: el usuario B ve el traceId, el tenant o, peor, la
// identidad de seguridad del usuario A. Este bug existe en producción.
// 3. La clave del ThreadLocalMap es una WeakReference al ThreadLocal, pero el VALOR es
// una referencia fuerte: si el valor referencia al ClassLoader de la aplicación, un
// redespliegue fuga el ClassLoader entero (OutOfMemoryError: Metaspace).
// ❌ InheritableThreadLocal: el hijo hereda una COPIA del valor del padre al crearse.
private static final InheritableThreadLocal<String> TENANT = new InheritableThreadLocal<>();
// Problemas: (a) no funciona con pools (los hilos ya existían cuando se puso el valor),
// (b) con un hilo por tarea copia el valor en CADA creación: con un millón de
// virtual threads, un millón de copias en memoria.
// ✅ EL SUSTITUTO MODERNO: ScopedValue (sección 9.2)
// - inmutable: se enlaza para un ámbito dinámico y se desenlaza solo al salir
// - no se puede olvidar el remove(): no hay remove()
// - se hereda de forma eficiente en concurrencia estructurada, sin copiar
private static final ScopedValue<String> TENANT_SV = ScopedValue.newInstance();
Checklist — java.util.concurrent
6 · Executors y pools de hilos
6.1 ExecutorService y las factorías de Executors
// El framework Executor separa "qué tarea" de "en qué hilo se ejecuta".
Executor basico = tarea -> new Thread(tarea).start(); // la interfaz mínima: execute(Runnable)
ExecutorService pool = Executors.newFixedThreadPool(8);
pool.execute(() -> log.info("fire and forget")); // Runnable, sin resultado
Future<Integer> f = pool.submit(() -> 42); // Callable, con resultado
List<Future<String>> todos = pool.invokeAll(listaDeCallables); // espera a TODOS
String primero = pool.invokeAny(listaDeCallables); // el primero que acabe
List<Future<String>> conPlazo = pool.invokeAll(listaDeCallables, 5, TimeUnit.SECONDS);
| Factoría | Configuración real interna | Peligro |
|---|---|---|
newFixedThreadPool(n) | core = max = n, LinkedBlockingQueue ILIMITADA | La cola crece sin freno: la latencia se dispara y acabas en OutOfMemoryError. Sin backpressure |
newCachedThreadPool() | core = 0, max = Integer.MAX_VALUE, keepAlive 60 s, SynchronousQueue | Crea un hilo por cada tarea sin hueco: bajo una ráfaga, miles de hilos y OutOfMemoryError: unable to create native thread |
newSingleThreadExecutor() | 1 hilo, cola ilimitada | Serializa todo (a veces es lo que quieres), pero la cola sigue siendo ilimitada |
newScheduledThreadPool(n) | DelayedWorkQueue ilimitada | Ver 6.5: una excepción no capturada cancela silenciosamente las ejecuciones futuras |
newWorkStealingPool() | ForkJoinPool con paralelismo = núcleos | Sin orden FIFO; pensado solo para cómputo divisible |
newVirtualThreadPerTaskExecutor() | Un virtual thread nuevo por tarea; no es un pool | No limita la concurrencia: 100.000 llamadas simultáneas a la BD la hunden. Limítala con Semaphore |
newFixedThreadPool es peligroso, con números: pool de 10 hilos, cada tarea 100
ms → capacidad 100 tareas/s. Si llegan 200/s, se acumulan 100 por segundo en una cola ilimitada. A
los 5 minutos hay 30.000 tareas en cola: cada nueva petición espera 5 minutos antes de empezar, todas
expiran por timeout en el cliente… y tú sigues procesándolas. Es el colapso por cola
infinita. La solución no es una cola más grande: es una cola acotada con política
de rechazo. Fallar rápido siempre es mejor que fallar lento.
6.2 Construir un ThreadPoolExecutor a mano (lo correcto)
CÓMO DECIDE ThreadPoolExecutor qué hacer con una tarea nueva
execute(tarea)
│
▼
¿hilos activos < corePoolSize?
SÍ ──► crear un hilo nuevo y ejecutar ya
NO
▼
¿cabe en la cola?
SÍ ──► encolar (¡y NO se crean más hilos aunque haya CPU libre!)
NO
▼
¿hilos activos < maximumPoolSize?
SÍ ──► crear un hilo nuevo (por encima del core) y ejecutar
NO
▼
RejectedExecutionHandler ← aquí decides tú qué pasa
⚠️ CONSECUENCIA CRÍTICA Y CONTRAINTUITIVA: con una cola ILIMITADA nunca se llega al paso 3,
así que maximumPoolSize se IGNORA por completo. Un pool con core=2, max=100 y
LinkedBlockingQueue() sin capacidad NUNCA tendrá más de 2 hilos. Es el error de
configuración de pools número uno.
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
public final class Pools {
/** Pool bien configurado: acotado en hilos Y en cola, con nombres y rechazo explícito. */
public static ThreadPoolExecutor crear(String nombre, int core, int max, int capacidadCola) {
ThreadFactory fabrica = new ThreadFactory() {
private final AtomicInteger n = new AtomicInteger(1);
@Override public Thread newThread(Runnable r) {
Thread t = new Thread(r, nombre + "-" + n.getAndIncrement());
t.setDaemon(false); // trabajo de negocio: no daemon
t.setUncaughtExceptionHandler((h, e) ->
LoggerFactory.getLogger(Pools.class).error("Excepción en {}", h.getName(), e));
return t;
}
};
ThreadPoolExecutor pool = new ThreadPoolExecutor(
core, // hilos que se mantienen vivos
max, // techo absoluto
60L, TimeUnit.SECONDS, // los hilos extra mueren tras 60 s inactivos
new ArrayBlockingQueue<>(capacidadCola), // ✅ ACOTADA: backpressure real
fabrica, // ✅ nombres reconocibles en jstack
new ThreadPoolExecutor.CallerRunsPolicy() // ✅ rechazo: lo ejecuta quien lo envió
);
pool.allowCoreThreadTimeOut(false); // true si quieres bajar a 0 en reposo
// pool.prestartAllCoreThreads(); // evita el warm-up del primer pico
return pool;
}
}
| Política de rechazo | Qué hace | Cuándo usarla |
|---|---|---|
AbortPolicy (por defecto) | Lanza RejectedExecutionException | API síncrona: devuelve 429/503 al cliente. Fail fast, la opción honesta |
CallerRunsPolicy | La ejecuta el hilo que la envió | Backpressure automático: el productor se frena solo. Ojo: si el llamante es el hilo HTTP, bloqueas una petición |
DiscardPolicy | Tira la tarea en silencio | Solo para datos prescindibles. Nunca sin una métrica que lo cuente |
DiscardOldestPolicy | Descarta la más antigua y reintenta | Cuando lo reciente vale más que lo viejo (cotizaciones, posiciones GPS) |
| Propia | rejectedExecution(Runnable, ThreadPoolExecutor) | Lo habitual en producción: métrica + log + excepción de dominio, o reencolar en disco/Kafka |
// Handler propio: lo que de verdad quieres en un servicio observable
RejectedExecutionHandler handler = (tarea, ejecutor) -> {
metricas.contador("pool.rechazos", "pool", "pedidos").incrementar();
log.warn("Pool saturado: activos={} cola={} completadas={}",
ejecutor.getActiveCount(), ejecutor.getQueue().size(), ejecutor.getCompletedTaskCount());
if (ejecutor.isShutdown()) throw new RejectedExecutionException("Pool cerrado");
throw new ServicioSaturadoException("Inténtalo más tarde"); // → HTTP 503 + Retry-After
};
6.3 Cierre correcto: shutdown, shutdownNow, awaitTermination
// ⚠️ NINGUNO de los tres, por sí solo, espera a que las tareas terminen.
// shutdown() → no acepta tareas nuevas; las pendientes SE EJECUTAN. No bloquea.
// shutdownNow() → no acepta nuevas, VACÍA la cola (te la devuelve) e INTERRUMPE a los
// hilos activos. No bloquea. Solo funciona si tus tareas respetan la
// interrupción (¡sección 2.4!).
// awaitTermination() → BLOQUEA hasta que todo termine o expire el plazo. Devuelve boolean.
/** Patrón de cierre en dos fases: el recomendado en el javadoc de ExecutorService. */
public static void cerrarOrdenadamente(ExecutorService pool, Duration plazo) {
pool.shutdown(); // 1) deja de aceptar tareas
try {
if (!pool.awaitTermination(plazo.toMillis(), TimeUnit.MILLISECONDS)) {
log.warn("Tareas sin terminar en {}; interrumpiendo", plazo);
List<Runnable> noEjecutadas = pool.shutdownNow(); // 2) plan B: interrumpir
log.warn("{} tareas descartadas de la cola", noEjecutadas.size());
if (!pool.awaitTermination(5, TimeUnit.SECONDS)) {
log.error("El pool NO ha muerto: hay tareas que ignoran la interrupción");
}
}
} catch (InterruptedException e) {
pool.shutdownNow(); // 3) nos han interrumpido a nosotros
Thread.currentThread().interrupt(); // ✅ restaurar el flag
}
}
// Java 19+: ExecutorService es AutoCloseable → close() = shutdown() + awaitTermination
// INDEFINIDO (reinterrumpiendo si hace falta). Perfecto para tareas acotadas y tests.
try (ExecutorService pool = Executors.newVirtualThreadPerTaskExecutor()) {
for (var url : urls) pool.submit(() -> descargar(url));
} // aquí se espera a TODAS. Ojo: si una tarea no termina nunca, close() no vuelve.
// En Spring: el cierre lo gestiona el contenedor si declaras el bean correctamente
@Bean(destroyMethod = "shutdown")
public ThreadPoolTaskExecutor ejecutorPedidos() {
var e = new ThreadPoolTaskExecutor();
e.setCorePoolSize(8);
e.setMaxPoolSize(16);
e.setQueueCapacity(500); // ✅ acotada
e.setThreadNamePrefix("pedidos-");
e.setWaitForTasksToCompleteOnShutdown(true); // ✅ apagado ordenado (graceful)
e.setAwaitTerminationSeconds(30); // ✅ menor que terminationGracePeriodSeconds de K8s
e.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
e.initialize();
return e;
}
6.4 Future frente a CompletableFuture
// Future (Java 5): un contenedor de resultado con una API muy limitada.
Future<String> f = pool.submit(() -> llamarServicio());
String r1 = f.get(); // ❌ BLOQUEA para siempre: nunca hagas esto
String r2 = f.get(3, TimeUnit.SECONDS); // ✅ con timeout, siempre
boolean cancelado = f.cancel(true); // true = interrumpir el hilo si ya empezó
f.isDone(); f.isCancelled();
f.state(); // Java 19+: RUNNING, SUCCESS, FAILED, CANCELLED
f.resultNow(); f.exceptionNow(); // Java 19+: sin bloquear (si ya terminó)
// ❌ LO QUE NO PUEDE HACER Future: encadenar, combinar, reaccionar. Solo esperar bloqueando,
// lo cual anula media ventaja de haber ido en paralelo.
// ⚠️⚠️ LA TRAMPA MÁS PREGUNTADA DEL MÓDULO:
// ¿qué pasa si una tarea lanza una excepción en un ExecutorService?
pool.execute(() -> { throw new RuntimeException("A"); });
// → el hilo muere, se llama al UncaughtExceptionHandler y el pool crea otro hilo.
// Es RUIDOSO: lo verás en los logs.
Future<?> olvidado = pool.submit(() -> { throw new RuntimeException("B"); });
// → la excepción se GUARDA dentro del Future y NO se imprime en ningún sitio. Si nadie
// llama a get(), el error DESAPARECE. Silencio absoluto. Este es el bug "mi tarea
// programada dejó de funcionar y nadie se enteró".
// ✅ Al llamar a get() sí aparece, envuelta:
try {
olvidado.get();
} catch (ExecutionException e) {
Throwable causa = e.getCause(); // ✅ la excepción REAL está en getCause()
log.error("La tarea falló", causa);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
// ✅ Regla: envuelve SIEMPRE el cuerpo de una tarea "fire and forget" en try/catch, o usa
// CompletableFuture y añade un .exceptionally(...) al final de la cadena.
pool.execute(() -> {
try { trabajoQuePuedeFallar(); }
catch (Exception e) { log.error("Fallo en la tarea de fondo", e); }
});
6.5 ScheduledExecutorService
ScheduledExecutorService programador = Executors.newScheduledThreadPool(
2, Thread.ofPlatform().name("scheduler-", 0).factory());
// Una sola vez, con retardo
programador.schedule(() -> log.info("una vez"), 5, TimeUnit.SECONDS);
// PERIÓDICO A RITMO FIJO: cada 10 s desde el INICIO del anterior. Si una ejecución tarda
// 15 s, la siguiente arranca inmediatamente después (se acumulan).
programador.scheduleAtFixedRate(this::publicarMetricas, 0, 10, TimeUnit.SECONDS);
// PERIÓDICO CON RETARDO FIJO: 10 s desde el FIN del anterior. Nunca se solapan.
// ✅ Es el que quieres el 90% de las veces (limpiezas, sondeos, reintentos).
programador.scheduleWithFixedDelay(this::limpiarCaducados, 0, 10, TimeUnit.SECONDS);
FIXED RATE vs FIXED DELAY (periodo 10 s; la tarea tarda 4 s y luego 15 s)
scheduleAtFixedRate(…, 0, 10, SECONDS)
t=0 [tarea 4s]
t=10 [tarea ····························15s····························]
t=25 [tarea 4s] ← ¡empieza YA: la "cita" de t=20 se había perdido!
t=30 [tarea 4s] ← vuelve a la rejilla original
scheduleWithFixedDelay(…, 0, 10, SECONDS)
t=0 [tarea 4s]
t=14 [tarea ····························15s····························]
t=39 [tarea 4s] ← siempre 10 s de descanso: nunca se solapan ni se acumulan
ScheduledExecutorService: si una ejecución periódica
lanza una excepción no capturada, la tarea se cancela para siempre y no se registra
nada en ningún log. Tu job de limpieza nocturna dejó de ejecutarse en marzo y lo descubriste en
agosto. La defensa es obligatoria.
// ✅ Envoltura defensiva para CUALQUIER tarea periódica
static Runnable blindar(String nombre, Runnable tarea) {
return () -> {
long inicio = System.nanoTime();
try {
tarea.run();
metricas.contador("job.ok", "job", nombre).incrementar();
} catch (Throwable t) { // Throwable: también Error
metricas.contador("job.error", "job", nombre).incrementar();
log.error("La tarea periódica '{}' ha fallado; se mantiene la programación", nombre, t);
// ✅ NO relanzar: si relanzas, la programación se cancela para siempre
} finally {
metricas.temporizador("job.duracion", "job", nombre)
.registrar(Duration.ofNanos(System.nanoTime() - inicio));
}
};
}
programador.scheduleWithFixedDelay(blindar("limpieza", this::limpiar), 0, 10, TimeUnit.MINUTES);
// ⚠️ Otras limitaciones: no hay expresiones cron, ni persistencia, ni coordinación entre
// instancias. Para eso: @Scheduled de Spring (cron sí) y, para varios nodos, un bloqueo
// distribuido con ShedLock (sección 11.3) o un planificador como Quartz.
6.6 ForkJoinPool, work-stealing y el commonPool
POOL CLÁSICO (una cola compartida) FORK/JOIN (una deque por hilo)
┌──────── cola única ────────┐ H1: [t1][t2][t3] ← push/pop por la CABEZA
│ t1 t2 t3 t4 t5 t6 │ H2: [t4] (LIFO: mejor caché)
└──┬────┬────┬────┬──────────┘ H3: [] ← vacío: ROBA de la COLA de H1
▼ ▼ ▼ ▼ H4: [t5][t6] (FIFO al robar: se
H1 H2 H3 H4 lleva la tarea más
Contención en la cola compartida "gorda", sin dividir)
WORK-STEALING: un hilo sin trabajo roba del EXTREMO OPUESTO de la deque de otro.
· Casi sin contención: cada hilo trabaja normalmente en su propia deque.
· Buena localidad de caché: LIFO en local (los datos están calientes).
· Reparto automático de carga sin planificador central.
· Y la propiedad clave: cuando una tarea hace join(), su hilo NO se duerme; ejecuta otras
tareas mientras espera. Por eso ForkJoinPool tolera la recursión que mataría a un pool
fijo (el thread starvation deadlock de la sección 3.6).
DIVIDE Y VENCE con umbral (la parte que la gente olvida)
ordenar(0..1.000.000)
├─ ordenar(0..500.000) ──── fork ──► otro hilo
└─ ordenar(500.000..1M) ─── compute (este hilo)
…hasta que el trozo baje del UMBRAL (típicamente 1.000-10.000 elementos) y entonces
se resuelve secuencialmente: por debajo del umbral, dividir cuesta más que lo que
se gana paralelizando.
// RecursiveTask<V> devuelve valor; RecursiveAction no devuelve nada.
public class SumaParalela extends RecursiveTask<Long> {
private static final int UMBRAL = 10_000; // ⚠️ medir: ni muy bajo ni muy alto
private final long[] datos;
private final int desde, hasta;
SumaParalela(long[] datos, int desde, int hasta) {
this.datos = datos; this.desde = desde; this.hasta = hasta;
}
@Override protected Long compute() {
int longitud = hasta - desde;
if (longitud <= UMBRAL) { // caso base: secuencial
long suma = 0;
for (int i = desde; i < hasta; i++) suma += datos[i];
return suma;
}
int medio = desde + longitud / 2;
SumaParalela izquierda = new SumaParalela(datos, desde, medio);
SumaParalela derecha = new SumaParalela(datos, medio, hasta);
izquierda.fork(); // ✅ una mitad al pool…
long resultadoDerecha = derecha.compute(); // ✅ …y la otra en ESTE hilo (evita un fork inútil)
long resultadoIzquierda = izquierda.join();
return resultadoIzquierda + resultadoDerecha;
}
public static long sumar(long[] datos) {
// ✅ Pool propio y dimensionado, no el commonPool
try (ForkJoinPool pool = new ForkJoinPool(Runtime.getRuntime().availableProcessors())) {
return pool.invoke(new SumaParalela(datos, 0, datos.length));
}
}
}
El commonPool: el recurso compartido que hunde aplicaciones
// TODO esto comparte UN SOLO pool global de la JVM:
// · parallelStream() y Arrays.parallelSort()
// · CompletableFuture.supplyAsync/thenApplyAsync SIN executor explícito
// · los métodos forEach/reduce/search paralelos de ConcurrentHashMap
// · muchas librerías de terceros que no lo documentan
int paralelismo = ForkJoinPool.getCommonPoolParallelism();
// = availableProcessors() - 1 ← ¡uno menos! En un contenedor con 1 CPU asignada, el
// paralelismo es 0 y todo se ejecuta en el hilo llamante:
// parallelStream() se vuelve secuencial sin avisar.
// ❌ EL DESASTRE: E/S bloqueante en el commonPool
List<Respuesta> respuestas = urls.parallelStream()
.map(url -> clienteHttp.get(url)) // ⛔ 500 ms bloqueado × N tareas
.toList();
// Con 8 núcleos hay 7 hilos en el commonPool. Si los ocupas con llamadas HTTP lentas, TODO
// lo demás que use el commonPool en la JVM entera se para: otros parallelStream, los
// CompletableFuture sin executor, los ordenamientos paralelos… Es un fallo GLOBAL provocado
// por un trozo de código LOCAL, y es dificilísimo de atribuir.
// ✅ SOLUCIÓN 1: executor propio para E/S
try (var exec = Executors.newFixedThreadPool(32, Thread.ofPlatform().name("http-", 0).factory())) {
List<CompletableFuture<Respuesta>> futuros = urls.stream()
.map(url -> CompletableFuture.supplyAsync(() -> clienteHttp.get(url), exec))
.toList();
respuestas = futuros.stream().map(CompletableFuture::join).toList();
}
// ✅ SOLUCIÓN 2 (2026, la buena): un virtual thread por llamada
try (var exec = Executors.newVirtualThreadPerTaskExecutor()) {
var tareas = urls.stream().map(u -> (Callable<Respuesta>) () -> clienteHttp.get(u)).toList();
for (Future<Respuesta> futuro : exec.invokeAll(tareas)) respuestas.add(futuro.get());
}
// ✅ SOLUCIÓN 3: si de verdad necesitas bloquear DENTRO de un ForkJoinPool, usa
// ManagedBlocker, que avisa al pool para que compense creando un hilo extra temporal.
ForkJoinPool.managedBlock(new ForkJoinPool.ManagedBlocker() {
private Respuesta resultado;
@Override public boolean block() { resultado = clienteHttp.get(url); return true; }
@Override public boolean isReleasable() { return resultado != null; }
});
// ⚠️ Cambiar el tamaño del commonPool es GLOBAL: último recurso.
// -Djava.util.concurrent.ForkJoinPool.common.parallelism=16
parallelStream(): úsalo solo si se cumplen todas estas
condiciones: (1) hay decenas de miles de elementos, (2) el trabajo por elemento es cómputo puro
sin E/S ni bloqueos, (3) la fuente se divide bien (array, ArrayList,
IntStream.range; no LinkedList ni Files.lines), (4) las operaciones
no tienen estado ni efectos laterales y (5) lo has medido con JMH. Si falla cualquiera, un
bucle secuencial o un executor propio es mejor.
6.7 Monitorización de pools: qué exponer sí o sí
| Métrica | De dónde sale | Qué te dice / cuándo alertar |
|---|---|---|
| Hilos activos | getActiveCount() | Pegado al máximo de forma sostenida: pool pequeño o tareas lentas |
| Tamaño del pool | getPoolSize() / getLargestPoolSize() | Si largest == max, has tocado techo alguna vez |
| Longitud de la cola | getQueue().size() | La métrica más importante. Cola creciendo = saturación. Alerta al 70 % de capacidad |
| Tareas rechazadas | Contador en tu RejectedExecutionHandler | Cualquier valor > 0 sostenido es un incidente |
| Tareas completadas | getCompletedTaskCount() | Su derivada es el throughput real del pool |
| Tiempo en cola | Instrumentar la tarea (marca de tiempo al enviar) | Latencia oculta que no aparece en el tiempo de ejecución |
| Duración de la tarea | Histograma propio | p99 alto = candidato a timeout o a bulkhead |
// Micrometer lo instrumenta por ti (Spring Boot Actuator): una línea y ya tienes
// executor.active, executor.queued, executor.pool.size, executor.completed, executor.seconds
@Bean
public ExecutorService ejecutorPedidos(MeterRegistry registro) {
ThreadPoolExecutor pool = Pools.crear("pedidos", 8, 16, 500);
return ExecutorServiceMetrics.monitor(registro, pool, "pedidos");
}
// Alerta útil en Prometheus (cola por encima del 70% durante 5 minutos)
// executor_queued_tasks{name="pedidos"} / 500 > 0.7
# Spring Boot: pools declarativos (no crees los tuyos si estos te sirven)
spring:
task:
execution: # el TaskExecutor de @Async
pool:
core-size: 8
max-size: 16
queue-capacity: 500 # ✅ acotada; por defecto es Integer.MAX_VALUE (¡ilimitada!)
keep-alive: 60s
thread-name-prefix: app-async-
shutdown:
await-termination: true
await-termination-period: 30s
scheduling: # el TaskScheduler de @Scheduled
pool:
size: 4 # ⚠️ por defecto es 1: dos @Scheduled lentos se estorban
thread-name-prefix: app-sched-
server:
tomcat:
threads:
max: 200 # hilos de plataforma que atienden peticiones
min-spare: 10
accept-count: 100 # cola de conexiones del SO cuando no hay hilos libres
max-connections: 8192
Checklist — executors y pools
7 · CompletableFuture en profundidad
CompletableFuture (Java 8) es un Future que además implementa
CompletionStage: se puede encadenar, combinar y completar desde fuera, con
manejo de errores integrado. Es el modelo asíncrono estándar del JDK y, aunque los virtual threads han
reducido su necesidad, sigue siendo la mejor herramienta para orquestar varias llamadas en paralelo
dentro de una misma petición.
7.1 Crear
// ⚠️ REGLA NÚMERO UNO: pasa SIEMPRE tu propio executor. Sin él se usa el commonPool.
private static final ExecutorService IO = Executors.newVirtualThreadPerTaskExecutor();
// (o un ThreadPoolExecutor acotado si sigues en Java 11/17)
CompletableFuture<Usuario> fu = CompletableFuture.supplyAsync(() -> repo.buscar(id), IO);
CompletableFuture<Void> fv = CompletableFuture.runAsync(() -> auditoria.registrar(id), IO);
// Ya completados (útiles en tests y para cortocircuitar cachés)
CompletableFuture<String> ok = CompletableFuture.completedFuture("cacheado");
CompletableFuture<String> fallo = CompletableFuture.failedFuture(new TimeoutException()); // Java 9+
// Completado manualmente: el puente para integrar APIs de callbacks
CompletableFuture<String> manual = new CompletableFuture<>();
clienteLegado.enviar(peticion, new Callback() {
@Override public void onSuccess(String r) { manual.complete(r); }
@Override public void onError(Throwable t) { manual.completeExceptionally(t); }
});
// Executor con retardo (Java 9+): reintentos y backoff sin ocupar un hilo esperando
Executor tras200ms = CompletableFuture.delayedExecutor(200, TimeUnit.MILLISECONDS, IO);
CompletableFuture.supplyAsync(() -> reintentar(), tras200ms);
7.2 Transformar y componer: thenApply vs thenCompose
LA DIFERENCIA QUE PREGUNTAN SIEMPRE (es map vs flatMap, igual que en Stream/Optional)
thenApply(f) : f devuelve un VALOR normal T -> U
resultado: CompletableFuture<U>
thenCompose(f) : f devuelve OTRO CompletableFuture T -> CompletionStage<U>
resultado: CompletableFuture<U> ← APLANA el anidamiento
Si usas thenApply con una función que devuelve un future, obtienes
CompletableFuture<CompletableFuture<U>> y tendrás que hacer join() dentro de la cadena:
exactamente lo que querías evitar.
SUFIJOS: los tres sabores de cada operador
thenApply(f) → se ejecuta en el hilo que completó el future anterior (o en
el llamante si ya estaba completo). Sin cambio de contexto:
ideal para transformaciones baratas.
thenApplyAsync(f) → se ejecuta en el commonPool ❌ evítalo
thenApplyAsync(f, executor) → se ejecuta en TU executor ✅ el correcto
// --- Transformar el valor ---
CompletableFuture<String> nombre = CompletableFuture
.supplyAsync(() -> repo.buscar(id), IO)
.thenApply(Usuario::nombre) // Usuario -> String (barato: sin Async)
.thenApply(String::toUpperCase);
// --- Encadenar otra llamada asíncrona (dependiente de la anterior) ---
CompletableFuture<List<Pedido>> pedidos = CompletableFuture
.supplyAsync(() -> repo.buscar(id), IO)
.thenCompose(u -> CompletableFuture.supplyAsync(() -> pedidoRepo.deCliente(u.nif()), IO));
// ✅ thenCompose aplana. Con thenApply tendrías CF<CF<List<Pedido>>>
// --- Combinar dos futuros INDEPENDIENTES (se ejecutan en paralelo) ---
CompletableFuture<Usuario> fUsuario = CompletableFuture.supplyAsync(() -> repo.buscar(id), IO);
CompletableFuture<Saldo> fSaldo = CompletableFuture.supplyAsync(() -> banco.saldo(id), IO);
CompletableFuture<Ficha> ficha = fUsuario.thenCombine(fSaldo, Ficha::new);
// Tiempo total ≈ max(t1, t2) en lugar de t1 + t2. Ese es TODO el objetivo del ejercicio.
// --- Efectos laterales sin cambiar el valor ---
ficha.thenAccept(f -> log.info("Ficha lista para {}", f.nif())); // consume, devuelve CF<Void>
ficha.thenRun(() -> metricas.incrementar("ficha.ok")); // ignora el valor
// --- El primero que llegue de dos alternativas equivalentes ---
CompletableFuture<Precio> rapido = proveedorA().applyToEither(proveedorB(), p -> p);
// (acceptEither / runAfterEither / runAfterBoth completan la familia)
7.3 allOf y anyOf con recolección de resultados
// ⚠️ allOf devuelve CompletableFuture<Void>: NO te da los resultados. Hay que recogerlos.
List<CompletableFuture<Producto>> futuros = ids.stream()
.map(id -> CompletableFuture.supplyAsync(() -> catalogo.buscar(id), IO))
.toList(); // ⚠️ .toList() aquí es IMPORTANTE: si dejas el stream perezoso,
// los futuros se crearían de uno en uno al consumirlo y NO habría paralelismo.
// Materializa primero, espera después.
// ✅ Patrón estándar de recolección
CompletableFuture<List<Producto>> todos = CompletableFuture
.allOf(futuros.toArray(CompletableFuture[]::new))
.thenApply(v -> futuros.stream()
.map(CompletableFuture::join) // ✅ seguro: allOf garantiza que ya acabaron
.toList());
// ✅ Variante TOLERANTE A FALLOS: los que fallen aportan un valor por defecto
CompletableFuture<List<Producto>> losQuePuedan = CompletableFuture
.allOf(futuros.toArray(CompletableFuture[]::new))
.handle((v, error) -> futuros.stream()
.map(f -> f.handle((p, e) -> e == null ? p : Producto.NO_DISPONIBLE).join())
.toList());
// ⚠️ Comportamiento clave: si UN future falla, el CF<Void> resultante falla, pero los DEMÁS
// siguen ejecutándose (no se cancelan). Si quieres cancelarlos, hazlo tú, o usa
// StructuredTaskScope (sección 9), que sí lo hace por diseño.
// anyOf: el primero que termine (¡incluso si termina con EXCEPCIÓN!)
CompletableFuture<Object> primero = CompletableFuture.anyOf(
replica("eu-west-1"), replica("eu-central-1"), replica("us-east-1"));
// ⚠️ Devuelve CompletableFuture<Object> (no tipado) y NO ignora los fallos: si la réplica
// más rápida devuelve un error, ese error es tu resultado. Para "el primero que tenga
// ÉXITO" hay que envolver cada rama con .exceptionally(...) y filtrar, o usar
// StructuredTaskScope con la política de éxito.
7.4 Manejo de errores: exceptionally, handle, whenComplete
| Operador | Se ejecuta… | ¿Puede cambiar el resultado? | Uso típico |
|---|---|---|---|
exceptionally(fn) | Solo si hubo error | Sí: devuelve el valor de recuperación | Fallback: valor por defecto, caché, degradación |
handle(bi) | Siempre (valor o error) | Sí: transforma ambos casos | Convertir a un Resultado/Either; normalizar |
whenComplete(bi) | Siempre | No: propaga tal cual | Efectos laterales: cerrar recursos, métricas, logs |
exceptionallyCompose(fn) | Solo si hubo error | Sí, con otro future | Reintentar o llamar a un servicio alternativo (Java 12+) |
// ⚠️ Las excepciones vienen ENVUELTAS en CompletionException (o ExecutionException en get()).
// Hay que desenvolverlas o los instanceof no funcionarán nunca.
static Throwable raiz(Throwable t) {
return (t instanceof CompletionException || t instanceof ExecutionException)
&& t.getCause() != null ? t.getCause() : t;
}
CompletableFuture<Precio> precio = CompletableFuture
.supplyAsync(() -> servicioPrecios.consultar(sku), IO)
// 1) Recuperación tipada
.exceptionally(e -> {
Throwable causa = raiz(e);
if (causa instanceof TimeoutException) {
metricas.incrementar("precios.timeout");
return cache.ultimoConocido(sku); // degradación aceptable
}
if (causa instanceof PrecioNoEncontradoException) return Precio.CERO;
throw new CompletionException(causa); // ✅ lo que no sé tratar, lo relanzo
})
// 2) Efectos laterales que NO deben alterar el resultado
.whenComplete((valor, error) -> {
if (error != null) log.error("Fallo consultando el precio de {}", sku, raiz(error));
else log.debug("Precio de {}: {}", sku, valor);
cronometro.parar(); // se ejecuta en ambos casos
});
// 3) handle: colapsar éxito y error en un tipo de dominio (muy limpio para agregaciones)
record Resultado<T>(T valor, Throwable error) {
boolean ok() { return error == null; }
}
CompletableFuture<Resultado<Precio>> nuncaFalla = CompletableFuture
.supplyAsync(() -> servicioPrecios.consultar(sku), IO)
.handle((v, e) -> new Resultado<>(v, e == null ? null : raiz(e)));
// 4) exceptionallyCompose (Java 12+): reintentar con OTRO future
CompletableFuture<Precio> conAlternativa = CompletableFuture
.supplyAsync(() -> proveedorPrincipal.precio(sku), IO)
.exceptionallyCompose(e -> CompletableFuture.supplyAsync(() -> proveedorSecundario.precio(sku), IO));
// ❌ EL ERROR MÁS COMÚN: no terminar la cadena. Sin exceptionally/handle/whenComplete, una
// excepción se queda DENTRO del future y desaparece sin dejar rastro en los logs.
CompletableFuture.runAsync(() -> enviarEmail(pedido), IO); // ⛔ fallo silencioso garantizado
CompletableFuture.runAsync(() -> enviarEmail(pedido), IO)
.exceptionally(e -> { log.error("No se envió el email del pedido {}", pedido.id(), e); return null; });
7.5 Timeouts
// Java 9+: timeouts declarativos, sin hilos extra (usan un temporizador interno)
CompletableFuture<Precio> conFallo = CompletableFuture
.supplyAsync(() -> servicio.lento(), IO)
.orTimeout(2, TimeUnit.SECONDS); // falla con TimeoutException
CompletableFuture<Precio> conDefecto = CompletableFuture
.supplyAsync(() -> servicio.lento(), IO)
.completeOnTimeout(Precio.CERO, 2, TimeUnit.SECONDS); // se completa con un valor
// ⚠️ MUY IMPORTANTE: orTimeout NO cancela el trabajo subyacente. La llamada HTTP sigue en
// marcha ocupando su hilo y su conexión; simplemente ya nadie espera el resultado. Los
// timeouts de la cadena son la ÚLTIMA red de seguridad, no la primera:
// 1.º timeouts en el cliente (connect + read) ← imprescindible
// 2.º timeouts en la consulta SQL (queryTimeout)
// 3.º orTimeout en la composición ← este
// Ver módulo 08, sección de resiliencia.
// Java 8 (sin orTimeout): se emula con un scheduler
static <T> CompletableFuture<T> conTimeout(CompletableFuture<T> f, long ms, ScheduledExecutorService s) {
CompletableFuture<T> timeout = new CompletableFuture<>();
ScheduledFuture<?> aviso = s.schedule(
() -> timeout.completeExceptionally(new TimeoutException("> " + ms + " ms")),
ms, TimeUnit.MILLISECONDS);
return f.applyToEither(timeout, x -> x)
.whenComplete((v, e) -> aviso.cancel(false)); // libera el temporizador
}
7.6 Ejemplo realista: agregar cuatro servicios en paralelo
/**
* Caso típico de BFF ("backend for frontend"): construir la pantalla de detalle de un
* pedido, que necesita datos de cuatro servicios distintos.
*
* SECUENCIAL : 120 + 200 + 90 + 150 = 560 ms
* PARALELO : max(120, 200, 90, 150) ≈ 200 ms ← 2,8× mejor latencia
*/
@Service
public class VistaPedidoService {
private static final Logger log = LoggerFactory.getLogger(VistaPedidoService.class);
private final PedidoClient pedidos;
private final ClienteClient clientes;
private final EnvioClient envios;
private final RecomendacionClient recomendaciones;
private final ExecutorService io; // ✅ executor propio inyectado
public VistaPedidoService(PedidoClient pedidos, ClienteClient clientes, EnvioClient envios,
RecomendacionClient recomendaciones,
@Qualifier("ioExecutor") ExecutorService io) {
this.pedidos = pedidos; this.clientes = clientes; this.envios = envios;
this.recomendaciones = recomendaciones; this.io = io;
}
public VistaPedido cargar(String idPedido, String traceId) {
// 1) El pedido es OBLIGATORIO: si falla, la respuesta entera falla.
CompletableFuture<Pedido> fPedido = CompletableFuture
.supplyAsync(() -> pedidos.porId(idPedido), io)
.orTimeout(1, TimeUnit.SECONDS);
// 2) El cliente DEPENDE del pedido → thenCompose (no thenApply)
CompletableFuture<Cliente> fCliente = fPedido
.thenCompose(p -> CompletableFuture.supplyAsync(() -> clientes.porNif(p.nif()), io))
.orTimeout(800, TimeUnit.MILLISECONDS);
// 3) El envío es OPCIONAL: si falla, degradamos con un valor "desconocido"
CompletableFuture<Envio> fEnvio = fPedido
.thenCompose(p -> CompletableFuture.supplyAsync(() -> envios.delPedido(p.id()), io))
.completeOnTimeout(Envio.desconocido(), 500, TimeUnit.MILLISECONDS)
.exceptionally(e -> {
log.warn("[{}] Envío no disponible para {}", traceId, idPedido, e);
metricas.incrementar("vista.envio.degradado");
return Envio.desconocido();
});
// 4) Las recomendaciones son PRESCINDIBLES y NO dependen del pedido: arrancan en
// paralelo desde el primer instante.
CompletableFuture<List<Producto>> fReco = CompletableFuture
.supplyAsync(() -> recomendaciones.para(idPedido), io)
.completeOnTimeout(List.of(), 300, TimeUnit.MILLISECONDS)
.exceptionally(e -> List.of());
// 5) Ensamblar. join() aquí es aceptable porque estamos en el hilo de la petición
// (Tomcat o un virtual thread), NO en un hilo del pool de E/S.
try {
return CompletableFuture
.allOf(fPedido, fCliente, fEnvio, fReco)
.thenApply(v -> new VistaPedido(
fPedido.join(), fCliente.join(), fEnvio.join(), fReco.join()))
.get(2, TimeUnit.SECONDS); // ✅ presupuesto TOTAL de la petición
} catch (TimeoutException e) {
// Cancelamos lo que quede para no seguir gastando recursos en balde
fPedido.cancel(true); fCliente.cancel(true); fEnvio.cancel(true); fReco.cancel(true);
throw new PasarelaTiempoAgotadoException(idPedido, e);
} catch (ExecutionException e) {
throw new VistaNoDisponibleException(idPedido, e.getCause());
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new PeticionCanceladaException(idPedido, e);
}
}
}
7.7 Los cinco errores típicos con CompletableFuture
// ❌ 1. Usar el commonPool para E/S (el pecado original). Ver 6.6.
CompletableFuture.supplyAsync(() -> http.get(url)); // ⛔ commonPool
CompletableFuture.supplyAsync(() -> http.get(url), miIoPool); // ✅
// ❌ 2. Olvidar el manejo de errores → el fallo se pierde para siempre
CompletableFuture.runAsync(this::sincronizar, io); // ⛔
CompletableFuture.runAsync(this::sincronizar, io)
.exceptionally(e -> { log.error("Fallo sincronizando", e); return null; }); // ✅
// ❌ 3. Bloquear con join()/get() DENTRO de una etapa del mismo pool
cf.thenApplyAsync(v -> otroFuture.join(), pool); // ⛔ ocupa un hilo del pool esperando
cf.thenCompose(v -> otroFuture); // ✅ composición, sin bloquear nada
// ❌ 4. Perder el paralelismo por no materializar el stream
var resultados = ids.stream()
.map(id -> CompletableFuture.supplyAsync(() -> buscar(id), io))
.map(CompletableFuture::join) // ⛔ SECUENCIAL: crea, espera, crea, espera…
.toList();
var futuros = ids.stream().map(id -> CompletableFuture.supplyAsync(() -> buscar(id), io)).toList();
var ok = futuros.stream().map(CompletableFuture::join).toList(); // ✅ crea TODOS, luego espera
// ❌ 5. Perder el contexto (traceId, seguridad, tenant): los ThreadLocal NO viajan de un
// hilo a otro. En Spring, SecurityContextHolder y el MDC de logs se vacían.
CompletableFuture.supplyAsync(() -> { log.info("sin traceId"); return 1; }, io); // ⛔
// ✅ Propagación explícita del contexto
String traceId = MDC.get("traceId");
CompletableFuture.supplyAsync(() -> {
MDC.put("traceId", traceId);
try { return trabajo(); } finally { MDC.clear(); } // ⚠️ limpiar: es un pool
}, io);
// ✅ O mejor: un TaskDecorator (Spring) / ContextSnapshot (Micrometer Context Propagation)
// que envuelva TODAS las tareas del executor, sin ensuciar cada llamada.
// ❌ BONUS: cancel(true) NO interrumpe el hilo que ejecuta un supplyAsync. A diferencia de
// Future.cancel(true) de un ExecutorService, aquí solo se marca el future como cancelado;
// el trabajo sigue. Si necesitas cancelación real, usa StructuredTaskScope (sección 9) o
// comprueba un flag dentro de la tarea.
8 · Virtual threads (Project Loom)
8.1 Qué son exactamente
Un virtual thread (JEP 444, estable en Java 21) es un
java.lang.Thread gestionado por la JVM, no por el sistema operativo. No tiene
un hilo del kernel asignado en exclusiva: se monta sobre un carrier thread (un
hilo de plataforma de un ForkJoinPool interno) solo mientras ejecuta código, y se
desmonta en cuanto se bloquea. Su pila vive en el heap y crece y se encoge según
haga falta.
MONTAJE Y DESMONTAJE — el mecanismo completo
Virtual threads (millones) Carrier threads (= nº de núcleos)
┌────┬────┬────┬────┬────┐ ┌──────────┬──────────┐
│ V1 │ V2 │ V3 │ … │ Vn │ ──────► │ Carrier1 │ Carrier2 │ ──► hilos del SO
└────┴────┴────┴────┴────┘ └──────────┴──────────┘
pila en el HEAP (cientos de bytes pila del SO (≈1 MB cada uno)
a pocos KB; crece y decrece) ForkJoinPool interno en modo FIFO
1. V1 se MONTA en Carrier1 y ejecuta código Java normal
2. V1 llama a socket.read() → no hay datos todavía
3. La JVM (con java.base reescrito para Loom) DESMONTA V1:
· copia su pila (continuation) al heap
· Carrier1 queda LIBRE inmediatamente
4. Carrier1 monta V2 y sigue trabajando ← ¡el hilo del SO nunca se bloquea!
5. Llegan los datos → el planificador vuelve a montar V1 (posiblemente en Carrier2:
NO hay afinidad de carrier)
6. V1 continúa justo después del read(), como si nada hubiera pasado
Para TU CÓDIGO nada de esto es visible: escribes código bloqueante, secuencial y
depurable. Ese es el punto: recuperar el estilo síncrono con el coste del asíncrono.
COMPARATIVA
Hilo de plataforma Virtual thread
Gestionado por Sistema operativo JVM
Pila ≈1 MB, fija, en el SO Heap, elástica (~KB)
Coste de creación ≈0,5-1 ms ≈1 µs (casi como un objeto normal)
Cambio de contexto 1-10 µs (syscall) ~100 ns (sin syscall)
Cuántos caben Miles Millones
¿Pooling? Sí, obligatorio NO, nunca (son desechables)
¿daemon? Configurable Siempre daemon
¿prioridad? Configurable (ignorada) Fija, no configurable
ThreadLocal Sí Sí, pero desaconsejado (n× memoria)
¿Aparecen en jstack? Sí NO por defecto → usa Thread.dump_to_file
// Cinco formas de usarlos (y una que NO debes usar)
Thread v1 = Thread.startVirtualThread(() -> log.info("hola"));
Thread v2 = Thread.ofVirtual().name("tarea-", 0).start(tarea);
Thread v3 = Thread.ofVirtual().unstarted(tarea); v3.start();
ThreadFactory fabricaV = Thread.ofVirtual().name("worker-", 0).factory();
// ✅ La forma canónica: un ejecutor que crea un virtual thread NUEVO por tarea
try (var exec = Executors.newVirtualThreadPerTaskExecutor()) {
for (var peticion : peticiones) exec.submit(() -> atender(peticion));
} // close() espera a que terminen todas
// ❌ NUNCA hagas esto: un pool de virtual threads es lo peor de los dos mundos
// ExecutorService malo = Executors.newFixedThreadPool(200, Thread.ofVirtual().factory());
// Comprobar y ajustar el planificador (propiedades del sistema, no API pública)
Thread.currentThread().isVirtual();
// -Djdk.virtualThreadScheduler.parallelism=8 (por defecto: availableProcessors)
// -Djdk.virtualThreadScheduler.maxPoolSize=256 (carriers máximos, para compensar bloqueos)
8.2 Qué problema resuelven y cuándo NO ayudan
EL PROBLEMA HISTÓRICO: "thread per request" tenía un techo bajísimo
1 petición = 1 hilo de plataforma = 1 MB + 1 hilo del SO
→ el techo NO era la CPU (al 5%), era el NÚMERO DE HILOS.
La respuesta de la industria (2010-2020) fue la programación asíncrona/reactiva: pocos
hilos + callbacks. Funciona, y tiene un precio brutal:
· el código se fragmenta en cadenas de operadores
· las trazas de pila dejan de significar nada (40 frames de Netty)
· los depuradores no sirven; los perfiladores tampoco
· los ThreadLocal (seguridad, MDC, transacciones) se rompen
· una llamada bloqueante colada por error hunde el event loop entero
Loom devuelve el estilo SÍNCRONO con el coste del ASÍNCRONO:
"el hilo vuelve a ser una unidad de concurrencia baratísima, no un recurso escaso"
DÓNDE AYUDA (mucho) — todo lo que ESPERA
· Llamadas HTTP a otros servicios (el caso estrella de los microservicios)
· Consultas a base de datos con JDBC bloqueante
· Lectura/escritura de ficheros, colas, cachés remotas
· Servidores con miles de conexiones concurrentes mayormente inactivas
DÓNDE NO AYUDA (y hay que decirlo claro)
· CPU-BOUND: si el trabajo es cálculo, el techo son los núcleos. 10.000 virtual threads
calculando primos van IGUAL o PEOR (más contención) que 8 hilos de plataforma. Para
esto: ForkJoinPool con paralelismo = núcleos.
· Cuando el cuello de botella es un recurso EXTERNO limitado (20 conexiones de BD, una
API con 10 req/s). Los virtual threads solo mueven la cola: hay que limitar con Semaphore.
· Cuando el código usa mucho ThreadLocal con objetos grandes: 100.000 hilos × un buffer
de 8 KB = 800 MB. Ahí ThreadLocal deja de ser una optimización.
· Cuando lo que sobra son hilos ociosos, no peticiones: si tu app va bien con 50 hilos,
Loom no te dará ni un 1% (aunque tampoco te quitará nada).
8.3 Pinning: cuando un virtual thread secuestra su carrier
Hay situaciones en las que la JVM no puede desmontar un virtual thread bloqueado, así que el carrier se queda atrapado con él. Si eso ocurre a la vez en tantos virtual threads como carriers hay, la aplicación se para: es el equivalente moderno de bloquear el event loop.
| Causa de pinning | Estado en 2026 | Qué hacer |
|---|---|---|
Bloquearse dentro de un synchronized (o esperando entrar en uno) |
Fijaba el carrier en Java 21–23. Resuelto en Java 24 (JEP 491, “Synchronize Virtual Threads without Pinning”) | En Java 21–23: sustituir por ReentrantLock en las rutas que bloquean. En Java 24+ ya no hace falta |
| Método nativo o upcall desde código nativo (JNI, FFM) | Sigue fijando: la JVM no puede capturar una pila nativa | Aislar esas llamadas en un pool de hilos de plataforma dedicado |
Object.wait() |
Fijaba antes; mejorado junto con synchronized en Java 24 |
Migrar a Condition/BlockingQueue de todos modos |
| Cargador de clases o inicializador estático que bloquea | Casos residuales | No hacer E/S en bloques static { } |
// ❌ En Java 21-23 este método FIJA el carrier durante toda la llamada HTTP
public synchronized Respuesta consultar(String sku) { // ⛔ synchronized + bloqueo de red
return http.get("/precios/" + sku);
}
// Con 8 carriers por defecto, 8 llamadas simultáneas paralizan la aplicación entera.
// ✅ Arreglo válido en TODAS las versiones: ReentrantLock nunca fija
private final ReentrantLock lock = new ReentrantLock();
public Respuesta consultarOk(String sku) {
lock.lock();
try { return http.get("/precios/" + sku); }
finally { lock.unlock(); }
}
// ✅✅ Mejor aún: ¿por qué hay un lock alrededor de una llamada de red? Casi nunca hace
// falta. Si es para limitar concurrencia, la herramienta correcta es un Semaphore.
private final Semaphore limite = new Semaphore(20);
public Respuesta consultarMejor(String sku) throws InterruptedException {
if (!limite.tryAcquire(1, TimeUnit.SECONDS)) throw new ServicioSaturadoException();
try { return http.get("/precios/" + sku); } finally { limite.release(); }
}
// ✅ Y para código nativo inevitable (una librería JNI de cifrado, por ejemplo):
private static final ExecutorService NATIVO =
Executors.newFixedThreadPool(4, Thread.ofPlatform().name("nativo-", 0).factory());
public byte[] cifrar(byte[] datos) throws Exception {
return NATIVO.submit(() -> libreriaNativa.cifrar(datos)).get(); // aislado en plataforma
}
# DETECTAR PINNING — el método depende de tu versión de Java
# Java 21-23: propiedad del sistema (imprime una traza cada vez que se fija un carrier)
java -Djdk.tracePinnedThreads=full -jar app.jar # full = traza completa
java -Djdk.tracePinnedThreads=short -jar app.jar # short = solo los frames problemáticos
# Salida: "Thread[#25,ForkJoinPool-1-worker-1,5,CarrierThreads] ... <== monitors:1"
# ⚠️ Esta propiedad se retiró en Java 24, cuando JEP 491 eliminó el pinning por synchronized.
# Java 21+ y método recomendado hoy: eventos JFR
java -XX:StartFlightRecording=duration=60s,filename=app.jfr,settings=profile -jar app.jar
jfr summary app.jfr
jfr print --events jdk.VirtualThreadPinned app.jfr # cada fijación, con su pila
jfr print --events jdk.VirtualThreadSubmitFailed app.jfr # el planificador no pudo montar
# Volcado de hilos que SÍ incluye los virtual threads (jstack no los muestra)
jcmd <pid> Thread.dump_to_file -format=json /tmp/hilos.json
jcmd <pid> Thread.dump_to_file -format=text /tmp/hilos.txt
# El formato JSON agrupa por StructuredTaskScope: es la única forma de leer un volcado con
# 200.000 hilos sin volverse loco.
8.4 No hacer pooling, y qué pasa con ThreadLocal
// ❌❌ EL ANTIPATRÓN NÚMERO UNO DE LOOM
ExecutorService malo = Executors.newFixedThreadPool(200, Thread.ofVirtual().factory());
/* Por qué está mal:
· Un pool existe para AMORTIZAR el coste de crear hilos. Crear un virtual thread cuesta
~1 µs: no hay nada que amortizar.
· El pool vuelve a imponer un LÍMITE artificial de 200 tareas simultáneas, que es justo
lo que Loom viene a eliminar.
· Los hilos se reutilizan → vuelven los problemas de ThreadLocal sucio.
Mentalidad correcta: un virtual thread es DESECHABLE, como un objeto. Crea uno por tarea
y déjalo morir. */
// ✅ Lo correcto
try (var exec = Executors.newVirtualThreadPerTaskExecutor()) { /* … */ }
// ✅ Y si necesitas LIMITAR la concurrencia (que casi siempre la necesitas), la herramienta
// NO es el tamaño del pool: es un semáforo. La diferencia es sutil y clave: el pool
// limita CUÁNTOS HILOS existen; el semáforo limita CUÁNTAS OPERACIONES concurrentes se
// hacen contra un recurso concreto. Cada recurso, su semáforo.
public class ClienteConLimite {
private final Semaphore limiteBd = new Semaphore(20); // = maximumPoolSize de Hikari
private final Semaphore limitePago = new Semaphore(5); // lo que aguanta la pasarela
private final ExecutorService exec = Executors.newVirtualThreadPerTaskExecutor();
public void procesar(List<Pedido> pedidos) {
for (Pedido p : pedidos) {
exec.submit(() -> { // 100.000 virtual threads: sin problema
limiteBd.acquire(); // pero solo 20 hablan con la BD a la vez
try { repo.guardar(p); } finally { limiteBd.release(); }
limitePago.acquire(); // y solo 5 con la pasarela
try { pasarela.cobrar(p); } finally { limitePago.release(); }
return null;
});
}
}
}
// THREADLOCAL Y MEMORIA con virtual threads
// Funciona, pero cambia la aritmética: antes tenías 200 hilos × 1 valor = 200 copias.
// Ahora puedes tener 500.000 hilos × 1 valor = 500.000 copias.
private static final ThreadLocal<byte[]> BUFFER = ThreadLocal.withInitial(() -> new byte[65536]);
// 500.000 × 64 KB = 32 GB 💥 Con virtual threads, un ThreadLocal "de reutilización" deja de
// ser una optimización y se convierte en una bomba de memoria.
// ✅ Para datos de contexto pequeños: ScopedValue (sección 9.2).
// ✅ Para buffers reutilizables: un pool explícito de buffers, no ThreadLocal.
8.5 Comparativa: 10.000 peticiones HTTP
import java.net.URI;
import java.net.http.*;
import java.time.Duration;
import java.util.concurrent.*;
import java.util.function.Supplier;
import java.util.stream.IntStream;
/**
* Compara el mismo trabajo (10.000 llamadas HTTP de ~100 ms) con un pool de plataforma y
* con un virtual thread por tarea. Apunta a un endpoint local que duerme 100 ms, para medir
* el efecto del modelo de hilos y no la red.
*
* Ejecuta con: java -Xmx512m ComparativaLoom.java
*/
public class ComparativaLoom {
static final int PETICIONES = 10_000;
static final HttpClient CLIENTE = HttpClient.newBuilder()
.connectTimeout(Duration.ofSeconds(2))
.executor(Executors.newVirtualThreadPerTaskExecutor())
.build();
public static void main(String[] args) throws Exception {
var uri = URI.create("http://localhost:8080/lento"); // endpoint que duerme 100 ms
medir("Pool de plataforma (200 hilos)", () -> Executors.newFixedThreadPool(200), uri);
medir("Pool de plataforma (1.000 hilos)", () -> Executors.newFixedThreadPool(1_000), uri);
medir("Virtual thread por tarea", Executors::newVirtualThreadPerTaskExecutor, uri);
}
static void medir(String etiqueta, Supplier<ExecutorService> fabrica, URI uri) throws Exception {
System.gc();
long memInicial = usada();
long t0 = System.nanoTime();
try (ExecutorService exec = fabrica.get()) {
var tareas = IntStream.range(0, PETICIONES)
.mapToObj(i -> (Callable<Integer>) () -> CLIENTE
.send(HttpRequest.newBuilder(uri).build(),
HttpResponse.BodyHandlers.discarding())
.statusCode())
.toList();
exec.invokeAll(tareas); // espera a todas
}
long ms = (System.nanoTime() - t0) / 1_000_000;
System.out.printf("%-38s %6d ms memoria: %+5d MB%n",
etiqueta, ms, (usada() - memInicial) / (1024 * 1024));
}
static long usada() {
Runtime r = Runtime.getRuntime();
return r.totalMemory() - r.freeMemory();
}
}
| Modelo | Tiempo total (orden de magnitud) | Memoria de los hilos | Por qué |
|---|---|---|---|
| Secuencial (1 hilo) | ≈ 1.000 s (17 min) | ≈ 1 MB | 10.000 × 100 ms en serie |
| Pool de 200 hilos de plataforma | ≈ 5 s | ≈ 200 MB reservados | 50 tandas × 100 ms; el techo es el número de hilos |
| Pool de 1.000 hilos de plataforma | ≈ 1,2 s | ≈ 1 GB reservados | 10 tandas; empieza a doler el cambio de contexto y la memoria |
| Pool de 10.000 hilos de plataforma | Suele fallar | ≈ 10 GB | OutOfMemoryError: unable to create native thread |
| Virtual thread por tarea | ≈ 0,3–0,7 s | ≈ 20–40 MB | Las 10.000 peticiones están en vuelo a la vez; el límite pasa a ser el servidor |
8.6 Integración con Spring Boot y el nuevo cuello de botella
# Spring Boot 3.2+ : una línea y Tomcat/Jetty atiende cada petición en un virtual thread.
# También cambia @Async, el TaskScheduler, los listeners de Kafka/RabbitMQ y el WebClient.
spring.threads.virtual.enabled=true
# Requiere Java 21+. Comprueba además:
# · que ninguna librería crítica use synchronized alrededor de E/S (si estás en 21-23)
# · que tus ThreadLocal se limpien (ahora hay muchísimos más "hilos")
# · que los pools de recursos externos estén dimensionados (ver abajo)
// Configuración explícita si necesitas más control
@Configuration
public class ConfiguracionVirtual {
@Bean
public TomcatProtocolHandlerCustomizer<?> virtualThreadHandler() {
return handler -> handler.setExecutor(Executors.newVirtualThreadPerTaskExecutor());
}
@Bean(name = "ioExecutor", destroyMethod = "close")
public ExecutorService ioExecutor() {
return Executors.newVirtualThreadPerTaskExecutor();
}
@Bean // @Async sobre virtual threads
public AsyncTaskExecutor applicationTaskExecutor() {
return new TaskExecutorAdapter(Executors.newVirtualThreadPerTaskExecutor());
}
}
EL NUEVO CUELLO DE BOTELLA: has movido el problema, no lo has eliminado
ANTES (200 hilos de Tomcat, 20 conexiones de Hikari)
200 peticiones máximo en el sistema → como mucho 200 esperan por 20 conexiones.
La cola está limitada por el número de hilos: incómodo, pero acotado.
DESPUÉS (virtual threads, 20 conexiones de Hikari)
50.000 peticiones simultáneas → 49.980 virtual threads haciendo cola en
Hikari.getConnection(). Todas expiran por connection-timeout (30 s por defecto) y el
cliente ve 500 en masa. Antes te protegía el límite de hilos; ahora no hay límite.
┌──────────────────────────────────────────────────────────────────────┐
│ 50.000 virtual threads ──► Semaphore(20) ──► Hikari(20) ──► BD │
│ ▲ │
│ AQUÍ pones el límite explícito y FALLAS RÁPIDO │
└──────────────────────────────────────────────────────────────────────┘
QUÉ HAY QUE HACER AL ACTIVAR VIRTUAL THREADS
1. Limitar la concurrencia POR RECURSO con Semaphore (BD, cada API externa, disco).
2. Bajar connection-timeout de Hikari (2-3 s) y devolver 503 con Retry-After en lugar de
acumular esperas de 30 s.
3. Poner timeouts en TODAS las llamadas remotas (ahora hay muchas más en vuelo).
4. Revisar los ThreadLocal: memoria y fugas de contexto entre peticiones.
5. Instrumentar: virtual threads en vuelo, permisos libres de cada semáforo y tiempo de
espera en la cola del semáforo.
6. Revisar el rate limiting DE ENTRADA: sin el límite de hilos, tu app aceptará toda la
carga que le echen. El backpressure ya no es implícito: hay que diseñarlo.
Checklist — virtual threads
9 · Concurrencia estructurada y ScopedValue
StructuredTaskScope y
ScopedValue forman parte de Project Loom pero, a diferencia de los virtual threads,
no son estables: han pasado por varias rondas de preview (JEP 428, 437, 453, 462,
480, 499, 505…) y la API ha cambiado entre versiones. En Java 21–23 se usa el estilo
new StructuredTaskScope.ShutdownOnFailure(); en Java 25 se rediseñó hacia
StructuredTaskScope.open(Joiner…). Requieren --enable-preview mientras sigan en
esa fase. Antes de usarlo en producción, comprueba el javadoc de tu versión exacta y no
afirmes en una entrevista que ya es estable: di que “es preview, con la API en evolución, y explico el
modelo”. Eso es lo que se valora.
9.1 StructuredTaskScope: subtareas con ciclo de vida acotado
EL PROBLEMA QUE RESUELVE — "concurrencia no estructurada"
Con un ExecutorService, las subtareas se van "de casa" y no vuelven:
· si el padre falla, las hijas siguen ejecutándose (fugas de trabajo)
· si una hija falla, el padre no se entera hasta que hace get()
· cancelar todo requiere código manual y propenso a olvidos
· la relación padre-hijo NO existe en el volcado de hilos: pierdes el contexto
La CONCURRENCIA ESTRUCTURADA aplica al hilo la misma idea que los bloques {} aplican al
flujo de control: si una tarea se divide en subtareas concurrentes, TODAS terminan (o se
cancelan) antes de salir del bloque. El árbol de llamadas vuelve a ser un árbol.
try (var scope = …) { ─┐
st1 = scope.fork(tarea1); │ las subtareas NO pueden sobrevivir al bloque
st2 = scope.fork(tarea2); │
scope.join(); │ ← punto de reunión obligatorio
} ─┘ close() garantiza que nada queda vivo
// ESTILO Java 21-23 (preview): política "todas o ninguna"
// Compila y ejecuta con: java --enable-preview --release 21 Agregador.java
import java.util.concurrent.StructuredTaskScope;
record Ficha(Usuario usuario, List<Pedido> pedidos, Saldo saldo) { }
Ficha cargarFicha(String nif) throws InterruptedException, ExecutionException {
// ShutdownOnFailure: si UNA subtarea falla, se CANCELAN automáticamente las demás
try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
StructuredTaskScope.Subtask<Usuario> sUsuario = scope.fork(() -> api.usuario(nif));
StructuredTaskScope.Subtask<List<Pedido>> sPedidos = scope.fork(() -> api.pedidos(nif));
StructuredTaskScope.Subtask<Saldo> sSaldo = scope.fork(() -> api.saldo(nif));
// Cada fork crea un VIRTUAL THREAD: 3 llamadas HTTP realmente en paralelo
scope.join(); // espera a las tres (o al primer fallo)
scope.throwIfFailed(); // si alguna falló, relanza su excepción aquí
return new Ficha(sUsuario.get(), sPedidos.get(), sSaldo.get()); // .get() ya no bloquea
} // close(): garantiza que NINGÚN hilo hijo sigue vivo al salir del try
}
// Con plazo global para toda la operación
Ficha cargarConPlazo(String nif) throws Exception {
try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
var sUsuario = scope.fork(() -> api.usuario(nif));
var sPedidos = scope.fork(() -> api.pedidos(nif));
var sSaldo = scope.fork(() -> api.saldo(nif));
scope.joinUntil(Instant.now().plusSeconds(2)); // presupuesto TOTAL: 2 s
scope.throwIfFailed(TiempoAgotadoException::new);
return new Ficha(sUsuario.get(), sPedidos.get(), sSaldo.get());
}
}
// Política "el primero que tenga ÉXITO gana": consultar réplicas y quedarse con la más rápida
String consultarRapido(String clave) throws Exception {
try (var scope = new StructuredTaskScope.ShutdownOnSuccess<String>()) {
scope.fork(() -> replica("eu-west-1").leer(clave));
scope.fork(() -> replica("eu-central-1").leer(clave));
scope.fork(() -> replica("us-east-1").leer(clave));
scope.join(); // vuelve en cuanto UNA termina bien; cancela las otras dos
return scope.result(); // lanza si TODAS fallaron
}
}
// Manejo parcial: quiero los resultados que se hayan podido obtener
List<Producto> loQueSePueda(List<String> skus) throws InterruptedException {
try (var scope = new StructuredTaskScope<Producto>()) { // sin política: no cancela nada
var subtareas = skus.stream().map(sku -> scope.fork(() -> catalogo.buscar(sku))).toList();
scope.join();
return subtareas.stream()
.filter(s -> s.state() == StructuredTaskScope.Subtask.State.SUCCESS)
.map(StructuredTaskScope.Subtask::get)
.toList();
}
}
// ⚠️ Java 25 (JEP 505, quinto preview) reorganizó la API hacia un método de fábrica y
// "joiners", en lugar de subclases. El modelo mental es el mismo; los nombres cambian:
// try (var scope = StructuredTaskScope.open(Joiner.<Ficha>allSuccessfulOrThrow())) { … }
// CONSULTA EL JAVADOC DE TU VERSIÓN antes de escribir código definitivo.
| Aspecto | ExecutorService + Future | StructuredTaskScope |
|---|---|---|
| Ciclo de vida de las subtareas | Independiente del llamante: pueden sobrevivirle | Acotado al bloque try: imposible que sobrevivan |
| Si una subtarea falla | Las demás siguen; lo descubres al hacer get() | Se cancelan automáticamente (con ShutdownOnFailure) |
| Si el llamante se cancela | Las hijas siguen ejecutándose (fuga) | Se propaga la cancelación hacia abajo |
| Legibilidad | Cadenas de futuros y join() dispersos | Bloque con entrada y salida claras |
| Diagnóstico | Hilos sueltos sin relación entre ellos | Árbol padre-hijo visible en Thread.dump_to_file |
| Madurez | Estable desde Java 5 | Preview, API en evolución |
9.2 ScopedValue frente a ThreadLocal
// ScopedValue (preview): contexto implícito INMUTABLE, con ámbito dinámico y acotado.
public final class Contexto {
public static final ScopedValue<String> TENANT = ScopedValue.newInstance();
public static final ScopedValue<String> TRACE_ID = ScopedValue.newInstance();
}
// Enlazar: el valor solo existe DENTRO de la lambda; al salir se desenlaza solo.
void manejarPeticion(Peticion p) {
ScopedValue.where(Contexto.TENANT, p.tenant())
.where(Contexto.TRACE_ID, p.traceId())
.run(() -> procesar(p)); // .call(...) si devuelve valor
} // aquí ya NO hay valor enlazado: no hay remove() que olvidar
// Leer, en cualquier punto de la pila de llamadas por debajo
void procesar(Peticion p) {
if (Contexto.TENANT.isBound()) {
log.info("[{}] procesando para tenant {}", Contexto.TRACE_ID.get(), Contexto.TENANT.get());
}
// Contexto.TENANT.get() sin enlace lanza NoSuchElementException: es intencionado,
// te obliga a razonar sobre el ámbito en lugar de recibir un null silencioso.
}
// Reenlace anidado para una parte del flujo (no muta: crea un ámbito nuevo)
ScopedValue.where(Contexto.TENANT, "sistema").run(this::tareaDeMantenimiento);
// ✅ Y la propiedad que lo hace brillar: HERENCIA en concurrencia estructurada.
// Las subtareas de un StructuredTaskScope ven los ScopedValue del padre SIN copiarlos
// (comparten la estructura inmutable), así que funciona con un millón de virtual threads.
try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
var a = scope.fork(() -> api.usuario(Contexto.TENANT.get())); // ✅ ve el tenant del padre
var b = scope.fork(() -> api.saldo(Contexto.TENANT.get())); // ✅ también
scope.join();
scope.throwIfFailed();
}
ThreadLocal | ScopedValue | |
|---|---|---|
| Mutabilidad | Mutable en cualquier momento (set) | Inmutable: se enlaza para un ámbito |
| Ciclo de vida | Vive mientras viva el hilo → hay que llamar a remove() | Acotado al bloque: se desenlaza automáticamente |
| Fugas en pools | Sí, si olvidas remove() (memoria y datos) | Imposible: no hay nada que limpiar |
| Herencia a hijos | InheritableThreadLocal: copia el valor | Compartición estructural, sin copias |
| Coste con 1 M de hilos | 1 M de copias en memoria | Prácticamente cero |
| Madurez | Desde Java 1.2, estable y omnipresente | Preview (JEP 429 y siguientes) |
10 · Programación reactiva: la comparativa honesta
10.1 De callbacks a reactive streams
// GENERACIÓN 1 — CALLBACKS (2005): funciona y es infernal de leer y de manejar errores
buscarUsuario(id, usuario ->
buscarPedidos(usuario, pedidos ->
calcularTotal(pedidos, total ->
enviarEmail(usuario, total, ok -> log.info("hecho"),
err -> log.error("email", err)),
err -> log.error("total", err)),
err -> log.error("pedidos", err)),
err -> log.error("usuario", err));
// "callback hell": la lógica se lee hacia dentro, el manejo de errores se duplica en cada
// nivel y no hay forma limpia de cancelar ni de aplicar un timeout global.
// GENERACIÓN 2 — FUTUROS COMPONIBLES (Java 8): plano, con errores centralizados
buscarUsuarioAsync(id)
.thenCompose(this::buscarPedidosAsync)
.thenApply(this::calcularTotal)
.thenCompose(t -> enviarEmailAsync(t))
.exceptionally(e -> { log.error("Fallo en el flujo", e); return null; });
// Mejor. Pero sigue siendo un valor ÚNICO: no sirve para flujos de N elementos ni tiene
// backpressure.
// GENERACIÓN 3 — REACTIVE STREAMS (2015, java.util.concurrent.Flow en Java 9):
// flujos de 0..N elementos CON contrapresión negociada entre productor y consumidor.
// Publisher → produce
// Subscriber → consume y PIDE (request(n)): "solo puedo con 10 más"
// Subscription → el canal de negociación
// Processor → ambas cosas
// El JDK define las interfaces (Flow) pero NO una implementación de operadores: para eso
// están Project Reactor (Spring), RxJava, Akka Streams o Mutiny (Quarkus).
// GENERACIÓN 4 — VIRTUAL THREADS (2023): vuelve el código secuencial…
Usuario u = api.usuario(id); // bloqueante, legible, depurable
List<Pedido> p = api.pedidos(u.nif()); // …y baratísimo, porque el hilo es virtual
enviarEmail(u, calcularTotal(p));
10.2 Project Reactor en cinco minutos
import reactor.core.publisher.*;
import reactor.core.scheduler.Schedulers;
// Mono<T> = 0 o 1 elemento | Flux<T> = 0..N elementos
Mono<Usuario> mono = Mono.just(new Usuario("12345678Z"));
Flux<Integer> flux = Flux.range(1, 10);
// ⚠️ Nada se ejecuta hasta que alguien se SUSCRIBE (lazy). Un pipeline sin subscribe() es
// código muerto: es el error número uno de quien empieza con Reactor.
flux.map(i -> i * 2).subscribe(System.out::println);
// Operadores esenciales
Flux<Pedido> pedidos = Flux.fromIterable(ids)
.flatMap(id -> repo.buscarReactivo(id), 8) // 8 = concurrencia máxima (¡importante!)
.filter(Pedido::estaActivo)
.map(this::enriquecer)
.timeout(Duration.ofSeconds(2))
.retryWhen(reactor.util.retry.Retry.backoff(3, Duration.ofMillis(200)))
.onErrorResume(e -> { log.error("fallo", e); return Flux.empty(); })
.doOnNext(p -> metricas.incrementar("pedido.procesado"))
.subscribeOn(Schedulers.boundedElastic()); // dónde se SUSCRIBE (fuente)
// publishOn(...) cambia el hilo de los operadores SIGUIENTES; subscribeOn afecta al origen.
// ⚠️ NUNCA hagas E/S bloqueante en Schedulers.parallel() (el event loop): usa boundedElastic
// para lo bloqueante inevitable, o el reactivo pierde todo el sentido.
// Combinar
Mono<Ficha> ficha = Mono.zip(api.usuarioMono(nif), api.saldoMono(nif))
.map(t -> new Ficha(t.getT1(), t.getT2()));
// BACKPRESSURE: el consumidor controla el ritmo
Flux.range(1, 1_000_000)
.onBackpressureBuffer(1000) // o onBackpressureDrop / onBackpressureLatest
.subscribe(new BaseSubscriber<Integer>() {
@Override protected void hookOnSubscribe(Subscription s) { request(10); } // pide 10
@Override protected void hookOnNext(Integer v) {
procesar(v);
request(1); // pide uno más cuando ha terminado
}
});
// Esto es lo que un pool de hilos NO te da gratis y una cola acotada sí (a su manera):
// el consumidor comunica su capacidad al productor en lugar de acumular en memoria.
10.3 ¿Reactivo o virtual threads en 2026?
| Criterio | Reactivo (Reactor/WebFlux) | Virtual threads (Loom) |
|---|---|---|
| Legibilidad | Cadenas de operadores; curva de aprendizaje alta | Código secuencial normal |
| Depuración | Trazas inútiles sin Hooks.onOperatorDebug() (que cuesta rendimiento) | Trazas normales, puntos de ruptura normales |
| Perfilado | El tiempo aparece atribuido a hilos del event loop | Cada hilo es una tarea: los flame graphs tienen sentido |
| Concurrencia de E/S masiva | Excelente | Excelente |
| Backpressure | Integrado en el modelo (request(n)) | No existe: lo diseñas tú con semáforos y colas acotadas |
| Streaming y eventos (SSE, WebSocket, Kafka) | Modelo natural: operadores de ventana, agrupación, tiempo | Se puede, pero reimplementas mucho |
| Composición de flujos complejos | Muy potente (window, buffer, groupBy, merge) | Manual |
| Riesgo típico | Una llamada bloqueante colada hunde el event loop | Saturar un recurso externo por falta de límites |
| Contexto (seguridad, MDC) | Context reactivo: distinto y contagioso | ScopedValue / ThreadLocal normales |
- Proyecto nuevo, API REST clásica con BD: Spring MVC + virtual threads. Más simple, más fácil de contratar gente que lo mantenga, y el rendimiento es equivalente para el 95 % de las cargas.
- Streaming real, SSE, WebSockets, pasarelas con miles de conexiones abiertas o pipelines de eventos con backpressure: reactivo sigue ganando, y sigue siendo la base de Spring Cloud Gateway.
- Ya tienes WebFlux en producción y funciona: no reescribas nada por moda. El coste de una migración es real y la ganancia puede ser cero.
- Lo que no debes hacer: elegir reactivo “porque es más moderno” sin necesitar backpressure ni streaming. Pagarás la complejidad todos los días y solo cobrarás el beneficio si de verdad tienes ese problema.
WebClient y el detalle de Spring Cloud Gateway se desarrollan en
módulo 04 · Spring Boot.
11 · Patrones de concurrencia en producción
11.1 Catálogo con diagramas
PRODUCTOR-CONSUMIDOR WORKER POOL
P1 ─┐ ┌──► W1 ──┐
P2 ─┼──► [cola acotada] ──► C tareas ──┼──► W2 ──┼──► resultados
P3 ─┘ │ └──► W3 ──┘
└ backpressure Reparto de carga; cada worker toma lo siguiente
Desacopla ritmos de producción que haya disponible. Es lo que hace un
y consumo. ThreadPoolExecutor.
PIPELINE (etapas especializadas) FAN-OUT / FAN-IN (scatter-gather)
[leer] → q1 → [validar] → q2 ┌──► servicio A ──┐
→ [enriquecer] → q3 │ │
→ [escribir] petición ┼──► servicio B ──┼──► agregar → respuesta
Cada etapa con su propio pool │ │
dimensionado a su coste. └──► servicio C ──┘
Ojo: la etapa más lenta manda. CompletableFuture.allOf o StructuredTaskScope.
THROTTLING (Semaphore) BULKHEAD (compartimentos)
1000 tareas ──► [20 permisos] pool "pagos" (8 hilos) ← si se satura,
──► recurso escaso pool "catálogo" (16 hilos) los demás siguen
Limita la concurrencia, no la pool "informes" (4 hilos) funcionando
tasa de llegada. Aísla fallos por dominio.
CIRCUIT BREAKER IDEMPOTENCIA + REINTENTOS
cerrado ──fallos>umbral──► abierto cliente envía Idempotency-Key
▲ │ servidor: ¿ya la vi? → devuelve la misma respuesta
└──éxito── semiabierto ◄───┘ Sin esto, reintentar = cobrar dos veces.
Deja de llamar a lo que está roto.
| Patrón | Problema que resuelve | Implementación en Java |
|---|---|---|
| Productor-consumidor | Ritmos distintos entre quien genera y quien procesa | ArrayBlockingQueue + hilos consumidores (sección 5.7) |
| Worker pool | Repartir tareas homogéneas entre N trabajadores | ThreadPoolExecutor acotado |
| Pipeline | Etapas con costes muy distintos que conviene dimensionar aparte | Varias colas y pools; o Flux con publishOn |
| Fan-out / fan-in | Agregar respuestas de varios servicios | allOf + join, o StructuredTaskScope |
| Throttling | Proteger un recurso con capacidad limitada | Semaphore con tryAcquire(timeout) |
| Rate limiting | Limitar operaciones por unidad de tiempo | Token bucket propio, Guava RateLimiter, Resilience4j, Redis |
| Circuit breaker | No insistir contra un servicio caído | Resilience4j (módulo 08) |
| Bulkhead | Que un dominio lento no consuma todos los hilos | Pools separados o @Bulkhead de Resilience4j |
| Idempotencia + reintentos | Reintentar sin duplicar efectos | Clave de idempotencia + tabla de deduplicación |
| Procesamiento por lotes | Reducir el coste por elemento | Acumular en BlockingQueue y volcar por tamaño o por tiempo |
| Tareas programadas idempotentes | Varias instancias ejecutando el mismo job | ShedLock, o SELECT … FOR UPDATE SKIP LOCKED |
11.2 Throttling con Semaphore y rate limiting con token bucket
Son cosas distintas y se confunden: throttling limita cuántas operaciones hay simultáneamente (concurrencia); rate limiting limita cuántas se hacen por segundo (tasa). Un semáforo de 10 permisos con operaciones de 1 ms permite 10.000 op/s; un rate limiter de 10 op/s permite como máximo 10 aunque duren un nanosegundo.
import java.time.Duration;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicLong;
/**
* Rate limiter TOKEN BUCKET sin hilos de fondo: los tokens se "rellenan" calculando el
* tiempo transcurrido en cada intento. Permite ráfagas de hasta 'capacidad' y una tasa
* media de 'tokensPorSegundo'. Thread-safe mediante CAS sobre un único long empaquetado.
*/
public final class TokenBucket {
private final long capacidad; // tamaño del cubo (ráfaga máxima)
private final double tokensPorNano; // tasa de relleno
private final AtomicLong estado; // nanos del último relleno
private final Semaphore tokens; // tokens disponibles
public TokenBucket(long capacidad, double tokensPorSegundo) {
if (capacidad <= 0 || tokensPorSegundo <= 0) throw new IllegalArgumentException();
this.capacidad = capacidad;
this.tokensPorNano = tokensPorSegundo / 1_000_000_000d;
this.tokens = new Semaphore((int) capacidad);
this.estado = new AtomicLong(System.nanoTime());
}
/** Rellena el cubo según el tiempo transcurrido. Idempotente y sin bloqueos. */
private void rellenar() {
long ahora = System.nanoTime();
long anterior = estado.get();
long transcurrido = ahora - anterior;
long nuevos = (long) (transcurrido * tokensPorNano);
if (nuevos <= 0) return;
// CAS: solo un hilo consigue "consumir" el intervalo y añadir los tokens
if (estado.compareAndSet(anterior, anterior + (long) (nuevos / tokensPorNano))) {
int hueco = (int) Math.min(nuevos, capacidad - tokens.availablePermits());
if (hueco > 0) tokens.release(hueco);
}
}
/** No bloquea: dice si la operación está permitida ahora mismo. */
public boolean intentar() {
rellenar();
return tokens.tryAcquire();
}
/** Espera hasta obtener un token o hasta agotar el plazo. */
public boolean intentar(Duration plazo) throws InterruptedException {
long finNanos = System.nanoTime() + plazo.toNanos();
while (System.nanoTime() < finNanos) {
rellenar();
if (tokens.tryAcquire()) return true;
Thread.sleep(1); // con virtual threads, esta espera es baratísima
}
return false;
}
public int disponibles() { return tokens.availablePermits(); } // métrica a exponer
}
// Uso en un cliente que respeta la cuota de una API externa (100 req/s, ráfagas de 200)
private final TokenBucket cuota = new TokenBucket(200, 100);
public Respuesta llamar(Peticion p) throws InterruptedException {
if (!cuota.intentar(Duration.ofMillis(500))) {
metricas.incrementar("api.rate_limited");
throw new DemasiadasPeticionesException("Cuota agotada"); // → HTTP 429
}
return http.enviar(p);
}
// Alternativas que NO deberías reimplementar si puedes usarlas:
// Guava: rate limiter de tasa suave (smooth), bloqueante
com.google.common.util.concurrent.RateLimiter limitador =
com.google.common.util.concurrent.RateLimiter.create(100.0); // 100 permisos/s
limitador.acquire(); // bloquea lo necesario
if (limitador.tryAcquire(50, TimeUnit.MILLISECONDS)) { /* … */ }
// Resilience4j: declarativo, con métricas y configuración externalizada
@RateLimiter(name = "apiPagos", fallbackMethod = "cobrarDegradado")
@Bulkhead(name = "apiPagos", type = Bulkhead.Type.SEMAPHORE) // throttling de concurrencia
@CircuitBreaker(name = "apiPagos", fallbackMethod = "cobrarDegradado")
@Retry(name = "apiPagos")
public Recibo cobrar(Pedido p) { return pasarela.cobrar(p); }
private Recibo cobrarDegradado(Pedido p, Throwable t) {
colaPendientes.publicar(p); // se procesará más tarde
return Recibo.pendiente(p.id());
}
// ⚠️ Para varias instancias del servicio, el rate limiting LOCAL no basta: 10 pods × 100
// req/s = 1.000 req/s contra el proveedor. Necesitas un limitador DISTRIBUIDO (Redis con
// un script Lua, o el rate limiting del API Gateway). Ver módulo 08.
11.3 Tareas programadas idempotentes con bloqueo distribuido
EL PROBLEMA: 3 réplicas del mismo servicio, un @Scheduled cada minuto
pod-1 ──┐
pod-2 ──┼──► "enviar recordatorios de pago" a las 09:00
pod-3 ──┘
Resultado: el cliente recibe TRES correos. Y en el caso de un cobro, tres cargos.
REGLA: toda tarea programada en un sistema con más de una instancia necesita
(a) EXCLUSIÓN (que solo una la ejecute) y (b) IDEMPOTENCIA (que si se ejecuta dos
veces, el efecto sea el mismo). Las dos: (a) puede fallar en un fallo de red.
// ShedLock: bloqueo distribuido sobre la BD (o Redis, ZooKeeper, MongoDB…)
@Configuration
@EnableScheduling
@EnableSchedulerLock(defaultLockAtMostFor = "PT10M") // seguro contra pods muertos
public class ConfiguracionTareas {
@Bean
public LockProvider lockProvider(DataSource dataSource) {
return new JdbcTemplateLockProvider(
JdbcTemplateLockProvider.Configuration.builder()
.withJdbcTemplate(new JdbcTemplate(dataSource))
.usingDbTime() // ✅ usa la hora de la BD: evita relojes desfasados
.build());
}
}
@Component
public class TareasNocturnas {
@Scheduled(cron = "0 0 3 * * *", zone = "Europe/Madrid")
@SchedulerLock(name = "conciliacionDiaria",
lockAtMostFor = "PT30M", // si el pod muere, el lock caduca en 30 min
lockAtLeastFor = "PT1M") // evita doble ejecución por relojes distintos
public void conciliar() {
// ✅ Además de exclusión, IDEMPOTENCIA: procesa por marca de agua, no "lo de ayer"
LocalDate ultima = repo.ultimaFechaConciliada();
for (LocalDate dia = ultima.plusDays(1); dia.isBefore(LocalDate.now()); dia = dia.plusDays(1)) {
conciliarDia(dia); // cada día se marca como hecho al terminar
}
}
}
// ⚠️ lockAtMostFor es un seguro, no una promesa: si la tarea tarda MÁS, el lock caduca y
// otra instancia puede empezar. Por eso la idempotencia sigue siendo obligatoria.
11.4 Colas en base de datos con SKIP LOCKED
Cuando no quieres añadir Kafka ni RabbitMQ para un volumen moderado, una tabla puede ser una cola perfecta,
siempre que varios workers no se peleen por las mismas filas. La clave es
FOR UPDATE SKIP LOCKED: cada worker bloquea las filas que va a procesar y
salta las que ya están bloqueadas por otro, en lugar de esperar.
-- PostgreSQL (también en MySQL 8+, Oracle y SQL Server con READPAST)
CREATE TABLE tarea_pendiente (
id BIGSERIAL PRIMARY KEY,
tipo TEXT NOT NULL,
carga JSONB NOT NULL,
estado TEXT NOT NULL DEFAULT 'PENDIENTE',
intentos INT NOT NULL DEFAULT 0,
procesar_en TIMESTAMPTZ NOT NULL DEFAULT now(), -- para reintentos con backoff
creado_en TIMESTAMPTZ NOT NULL DEFAULT now()
);
-- Índice que hace eficiente la consulta de la cola
CREATE INDEX idx_cola ON tarea_pendiente (estado, procesar_en) WHERE estado = 'PENDIENTE';
-- El corazón del patrón: cada worker se lleva un lote SIN esperar a los demás
BEGIN;
SELECT id, tipo, carga
FROM tarea_pendiente
WHERE estado = 'PENDIENTE'
AND procesar_en <= now()
ORDER BY procesar_en
LIMIT 20
FOR UPDATE SKIP LOCKED; -- ← sin SKIP LOCKED, 10 workers harían cola por la misma fila
-- …procesar en la aplicación y marcar el resultado…
UPDATE tarea_pendiente SET estado = 'HECHA' WHERE id = ANY($1);
-- Reintento con backoff exponencial para las que fallen
UPDATE tarea_pendiente
SET intentos = intentos + 1,
procesar_en = now() + (interval '10 seconds' * power(2, intentos)),
estado = CASE WHEN intentos + 1 >= 5 THEN 'FALLIDA' ELSE 'PENDIENTE' END
WHERE id = $1;
COMMIT;
// El worker en Java: transacción corta, lote acotado y virtual threads para el trabajo de E/S
@Component
public class WorkerCola {
private final ColaRepository repo;
private final Semaphore limiteExterno = new Semaphore(10); // protege la API de destino
@Scheduled(fixedDelay = 1000) // fixedDelay, no fixedRate: nunca se solapa
public void procesarLote() {
List<Tarea> lote = repo.reclamarLote(20); // SELECT … FOR UPDATE SKIP LOCKED
if (lote.isEmpty()) return;
try (var exec = Executors.newVirtualThreadPerTaskExecutor()) {
for (Tarea t : lote) {
exec.submit(() -> {
limiteExterno.acquire();
try { ejecutar(t); repo.marcarHecha(t.id()); }
catch (Exception e) { repo.marcarFallo(t.id(), e.getMessage()); }
finally { limiteExterno.release(); }
return null;
});
}
} // close() espera a todo el lote
}
}
// Ventajas: cero infraestructura extra, transaccionalidad con tus datos (patrón outbox) y
// visibilidad total con SQL. Límites: no escala a millones de mensajes por minuto ni ofrece
// fan-out a varios consumidores. Para eso, Kafka (módulo 08).
12 · Testing de código concurrente
12.1 Por qué estos tests son flakies
Un test concurrente explora un entrelazado de los millones posibles, y el que explora depende de la máquina, la carga y el JIT. Consecuencias que hay que aceptar:
- Un test que pasa no demuestra que el código sea correcto. Solo que ese entrelazado concreto funcionó.
- Un test que falla una vez entre cien está detectando un bug real. No lo desactives: es la única pista que vas a tener.
- Nunca uses
Thread.sleeppara sincronizar un test. O es demasiado corto (falla en CI, que es más lenta) o demasiado largo (la suite tarda una eternidad). Usa latches o Awaitility. - Repite. Un bucle de 1.000 iteraciones con muchos hilos encuentra lo que una sola ejecución no encuentra jamás.
12.2 El patrón del doble latch: máxima colisión
import org.junit.jupiter.api.RepeatedTest;
import org.junit.jupiter.api.Test;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
import static org.assertj.core.api.Assertions.assertThat;
class ContadorConcurrenteTest {
/**
* Patrón estándar: TODOS los hilos esperan en el mismo latch y salen a la vez. Así se
* maximiza la probabilidad de colisión real, en lugar de que se ejecuten en fila.
*/
@RepeatedTest(20) // ✅ repetir: un solo intento no prueba nada
void elContadorNoPierdeIncrementos() throws Exception {
int hilos = 32, porHilo = 1_000;
var contador = new AtomicInteger(); // la clase bajo prueba
var salida = new CountDownLatch(1); // pistoletazo de salida
var fin = new CountDownLatch(hilos); // meta
var errores = new CopyOnWriteArrayList<Throwable>(); // los fallos de un hilo hijo
// NO hacen fallar el test solos
try (var exec = Executors.newVirtualThreadPerTaskExecutor()) {
for (int i = 0; i < hilos; i++) {
exec.submit(() -> {
try {
salida.await(); // todos esperan aquí
for (int j = 0; j < porHilo; j++) contador.incrementAndGet();
} catch (Throwable t) {
errores.add(t); // ✅ recoger, no ignorar
} finally {
fin.countDown(); // ✅ en finally: nunca cuelgues el test
}
return null;
});
}
salida.countDown(); // ¡ya!
assertThat(fin.await(10, TimeUnit.SECONDS))
.as("los hilos deben terminar en 10 s")
.isTrue(); // ✅ con timeout: un test no se cuelga
}
assertThat(errores).isEmpty();
assertThat(contador.get()).isEqualTo(hilos * porHilo);
}
/** Comprobar que una operación es idempotente bajo concurrencia. */
@Test
void soloUnHiloCreaElRecurso() throws Exception {
var servicio = new ServicioIdempotente();
int hilos = 50;
var salida = new CountDownLatch(1);
var creaciones = new AtomicInteger();
try (var exec = Executors.newVirtualThreadPerTaskExecutor()) {
var futuros = IntStream.range(0, hilos)
.mapToObj(i -> exec.submit(() -> {
salida.await();
if (servicio.crearSiNoExiste("clave-unica")) creaciones.incrementAndGet();
return null;
}))
.toList();
salida.countDown();
for (var f : futuros) f.get(5, TimeUnit.SECONDS); // propaga excepciones al test
}
assertThat(creaciones.get()).isEqualTo(1); // exactamente una creación
}
}
12.3 Awaitility: esperar por una condición, no por un reloj
import static org.awaitility.Awaitility.await;
import static java.time.Duration.ofSeconds;
import static java.time.Duration.ofMillis;
// ❌ Frágil y lento: adivinar cuánto tarda
@Test void malTest() throws Exception {
servicio.procesarAsincrono(pedido);
Thread.sleep(2000); // ⛔ ¿y si CI tarda 2,1 s?
assertThat(repo.buscar(pedido.id())).isPresent();
}
// ✅ Espera activa con plazo: rápido cuando va bien, tolerante cuando va lento
@Test void buenTest() {
servicio.procesarAsincrono(pedido);
await().atMost(ofSeconds(5)) // plazo máximo
.pollInterval(ofMillis(50)) // cada cuánto comprueba
.pollDelay(ofMillis(0)) // sin espera inicial
.untilAsserted(() -> assertThat(repo.buscar(pedido.id())).isPresent());
}
// Variantes útiles
await().atMost(ofSeconds(3)).until(() -> cola.isEmpty());
await().atMost(ofSeconds(3)).until(contador::get, valor -> valor >= 100);
await().during(ofSeconds(1)).atMost(ofSeconds(3)) // se mantiene estable 1 s
.until(() -> estado.get() == Estado.LISTO);
12.4 Tests deterministas: inyecta el ejecutor
// ✅ La mejor técnica: hacer que la concurrencia sea una DEPENDENCIA, no un detalle interno.
public class ServicioPedidos {
private final Executor executor; // ← inyectado, no creado dentro
public ServicioPedidos(Executor executor) { this.executor = executor; }
public void confirmar(Pedido p) {
repo.guardar(p);
executor.execute(() -> notificaciones.enviar(p)); // asíncrono en producción
}
}
// En el test: un ejecutor SÍNCRONO. El test es 100% determinista y sin esperas.
Executor sincrono = Runnable::run; // ejecuta en el hilo del test
var servicio = new ServicioPedidos(sincrono);
servicio.confirmar(pedido);
verify(notificaciones).enviar(pedido); // ya ha ocurrido: sin sleeps ni latches
// En Spring: SyncTaskExecutor hace lo mismo
@TestConfiguration
static class ConfigSincrona {
@Bean @Primary
public TaskExecutor taskExecutor() { return new SyncTaskExecutor(); }
}
// ✅ Y para la lógica que depende del tiempo, inyecta un Clock: nada de LocalDateTime.now()
private final Clock clock;
Instant ahora = clock.instant();
// En el test: Clock.fixed(Instant.parse("2026-01-15T10:00:00Z"), ZoneOffset.UTC)
12.5 Herramientas serias: jcstress y aserciones
// jcstress (Java Concurrency Stress tests): la herramienta del equipo de OpenJDK para
// comprobar propiedades del MODELO DE MEMORIA. Ejecuta millones de veces el mismo par de
// acciones y clasifica TODOS los resultados observados. Es lo único que demuestra de verdad
// si un resultado es posible.
@JCStressTest
@Outcome(id = "1, 1", expect = Expect.ACCEPTABLE, desc = "ambos ven la escritura")
@Outcome(id = "0, 0", expect = Expect.ACCEPTABLE_INTERESTING, desc = "¡reordenamiento observado!")
@State
public class TestReordenamiento {
int x, y;
@Actor public void hilo1(II_Result r) { x = 1; r.r1 = y; }
@Actor public void hilo2(II_Result r) { y = 1; r.r2 = x; }
}
// mvn clean verify && java -jar target/jcstress.jar -t TestReordenamiento
// Verás que "0, 0" ocurre de verdad, y con volatile deja de ocurrir. Es la mejor forma de
// convencerte a ti mismo (y a un compañero escéptico) de que el JMM importa.
# Aserciones activadas: muchas clases del JDK y de librerías comprueban invariantes con assert
java -ea -jar app.jar # -ea = -enableassertions (jamás en producción: coste y semántica)
# Detectar carreras de datos en tiempo de ejecución: Java NO tiene un ThreadSanitizer como
# C++ o Go. Las alternativas reales son:
# · jcstress → para propiedades del JMM en unidades pequeñas
# · SpotBugs + fb-contrib (detectores de concurrencia estáticos: IS2_INCONSISTENT_SYNC,
# UG_SYNC_SET_UNSYNC_GET, LI_LAZY_INIT_UPDATE_STATIC…)
# · Error Prone con @GuardedBy / @ThreadSafe (comprobación estática de anotaciones JCIP)
# · Pruebas de carga con aserciones de invariantes al final (p. ej. "la suma cuadra")
# · JFR + análisis de eventos de bloqueo (sección 13)
# Anotaciones de documentación ejecutable (JCIP / jsr305): merecen la pena
# @ThreadSafe @NotThreadSafe @Immutable @GuardedBy("this")
13 · Depuración y diagnóstico
13.1 Volcados de hilos: cómo se leen
# 1) Capturar. Toma SIEMPRE 3 volcados separados 5-10 s: lo importante es qué CAMBIA.
jps -l # localizar el PID
for i in 1 2 3; do jcmd $PID Thread.print -l > hilos-$i.txt; sleep 7; done
# Alternativas equivalentes
jstack -l $PID > hilos.txt
kill -3 $PID # el volcado va a stdout del proceso
# 2) Con virtual threads, jstack NO los muestra. Usa:
jcmd $PID Thread.dump_to_file -format=json /tmp/hilos.json
jcmd $PID Thread.dump_to_file -format=text /tmp/hilos.txt
# 3) Triaje rápido en 30 segundos
grep -c '^"' hilos-1.txt # ¿cuántos hilos hay? (¿han crecido?)
grep 'java.lang.Thread.State' hilos-1.txt | sort | uniq -c | sort -rn # reparto por estado
grep -A 3 'BLOCKED' hilos-1.txt | grep 'waiting to lock' | sort | uniq -c | sort -rn
grep 'Found one Java-level deadlock' hilos-1.txt
grep -c 'pedidos-worker' hilos-1.txt # ¿se ha desbocado un pool concreto?
# 4) ¿Qué hilo consume la CPU? (Linux) — el truco que hay que saber
top -H -p $PID # localiza el TID que quema CPU
printf '%x\n' <TID> # conviértelo a hexadecimal
grep -A 20 'nid=0x<hex>' hilos-1.txt # ese "nid" es tu hilo en el volcado
| Lo que ves en los 3 volcados | Diagnóstico probable | Siguiente paso |
|---|---|---|
Los mismos hilos en BLOCKED sobre el mismo monitor | Contención grave o deadlock | grep "Found one Java-level deadlock"; reducir la sección crítica |
Muchos RUNNABLE en SocketRead/SocketInputStream | Espera de red, no CPU | Timeouts, circuit breaker, revisar el servicio remoto |
Todos los hilos del pool en getConnection de Hikari | Pool de conexiones agotado | Buscar conexiones sin cerrar; revisar maximumPoolSize y transacciones largas |
| Número de hilos creciendo volcado a volcado | Fuga de hilos (executor creado por petición) | Buscar new Thread / Executors.new… dentro de métodos |
| Un hilo siempre en el mismo método con CPU alta | Bucle infinito, regex catastrófica, HashMap corrupto | async-profiler sobre ese hilo |
Casi todo en TIMED_WAITING en pools | Normal: pools en reposo | Nada |
Muchos hilos en parkNanos sobre el mismo Semaphore | Throttling funcionando (o mal dimensionado) | Medir el tiempo de espera; ajustar permisos |
13.2 JFR: eventos de bloqueo con impacto casi nulo
# Arrancar con grabación (impacto < 2%: se puede dejar en producción)
java -XX:StartFlightRecording=duration=120s,filename=/dumps/app.jfr,settings=profile -jar app.jar
# O activarla en caliente sobre un proceso vivo
jcmd $PID JFR.start name=diag settings=profile duration=120s filename=/dumps/app.jfr
jcmd $PID JFR.dump name=diag filename=/dumps/parcial.jfr
jcmd $PID JFR.stop name=diag
# EVENTOS CLAVE PARA CONCURRENCIA
jfr summary /dumps/app.jfr # visión general
jfr print --events jdk.JavaMonitorEnter /dumps/app.jfr # espera en synchronized (¡el oro!)
jfr print --events jdk.JavaMonitorWait /dumps/app.jfr # Object.wait()
jfr print --events jdk.ThreadPark /dumps/app.jfr # LockSupport.park: locks del j.u.c.
jfr print --events jdk.ThreadStart,jdk.ThreadEnd /dumps/app.jfr
jfr print --events jdk.VirtualThreadPinned /dumps/app.jfr # pinning (sección 8.3)
jfr print --events jdk.VirtualThreadStart /dumps/app.jfr
jfr print --events jdk.ExecutionSample /dumps/app.jfr # muestreo de CPU
# Y ábrelo en JDK Mission Control (JMC): la pestaña "Lock Instances" te da, ordenado,
# QUÉ objeto se bloquea más, CUÁNTO tiempo total y QUÉ pila lo pide. Es la forma más rápida
# que existe de encontrar el lock que te está costando dinero.
13.3 async-profiler en modo wall-clock
# async-profiler tiene DOS modos, y elegir mal es el error habitual:
# -e cpu → dónde se GASTA CPU (para problemas de cómputo)
# -e wall → dónde se PASA EL TIEMPO (para problemas de LATENCIA: esperas y bloqueos)
#
# Si tu servicio está lento con la CPU al 10%, el modo cpu no te dirá nada útil: todo el
# tiempo está en esperas, que el modo cpu no muestra. Usa wall-clock.
./profiler.sh -d 60 -e wall -t -f /tmp/wall.html $PID # -t = separar por hilo
./profiler.sh -d 60 -e cpu -f /tmp/cpu.html $PID
./profiler.sh -d 60 -e lock -f /tmp/lock.html $PID # contención de locks
./profiler.sh -d 60 -e alloc -f /tmp/alloc.html $PID # presión de asignación
# El flame graph de wall-clock te muestra, literalmente, en qué línea espera tu petición.
# En un backend típico verás: getConnection, socketRead, y ahí está el 90% de tu latencia.
13.4 Métricas mínimas de concurrencia
| Métrica | Fuente | Umbral de alerta razonable |
|---|---|---|
jvm.threads.live / peak | Micrometer / JMX | Crecimiento sostenido = fuga de hilos |
executor.queued | Micrometer | > 70 % de la capacidad durante 5 min |
| Tareas rechazadas | Contador propio en el handler | Cualquier valor sostenido |
hikaricp.connections.pending | Micrometer | > 0 sostenido = pool insuficiente o fuga |
hikaricp.connections.usage | Micrometer | p99 alto = transacciones largas |
| Permisos libres del semáforo | availablePermits() como gauge | 0 sostenido = recurso saturado |
| Hilos en deadlock | ThreadMXBean.findDeadlockedThreads() | > 0 = incidente inmediato |
| Latencia p99 por endpoint | http.server.requests | Según SLO |
14 · Buenas prácticas y malas prácticas
14.1 Buenas prácticas (en orden de rentabilidad)
Diseño: elimina el problema antes de resolverlo
- Prefiere la inmutabilidad. Un objeto inmutable (campos
final, sin setters, colecciones copiadas y envueltas,record) es thread-safe gratis y para siempre. Es la única forma de concurrencia que no se puede escribir mal. - No compartas estado mutable. Si dos hilos no ven el mismo objeto, no hay nada que sincronizar. Pasa copias, valores o mensajes.
- Confinamiento por hilo. Si el estado tiene que ser mutable, que lo toque un solo hilo (variables locales, un único consumidor por cola). Es la segunda mejor opción.
- Estado compartido en un único sitio. Concentra la mutación en una clase
pequeña, bien documentada (
@ThreadSafe,@GuardedBy) y probada. No la repartas por diez servicios. - Prefiere
java.util.concurrenta construirlo tú. Está escrito por Doug Lea, probado con jcstress y usado por millones de aplicaciones. Tuwait/notifyno.
Ejecución: pools, tiempos y visibilidad
- Timeouts en todo. Conexión, lectura,
tryLock,tryAcquire,get,awaitTermination,orTimeout. Una espera sin plazo es una caída futura. - Colas acotadas y rechazo explícito. Es mejor devolver 503 rápido que
acumular trabajo hasta el
OutOfMemoryError. - Nombra los hilos.
pedidos-worker-3se diagnostica en diez segundos;pool-4-thread-7no se diagnostica. - Dimensiona con datos. CPU-bound ≈ núcleos; IO-bound, con la fórmula de Goetz y validando con la ley de Little (sección 1.5). Y mide.
- Separa pools por dominio (bulkhead): que los informes lentos no se coman los hilos de los pagos.
- Cierra ordenadamente:
shutdown→awaitTermination→shutdownNow, y contry (var exec = …)cuando el ámbito sea local. - Instrumenta. Hilos vivos, cola, rechazos, conexiones pendientes. Sin métricas, la concurrencia es una caja negra.
- Propaga la interrupción o restaura el flag: es el contrato de cancelación de toda la plataforma.
- Documenta el porqué de cada
synchronizedy cadavolatile: “protege el invariante X junto con Y”. Sin eso, el siguiente lo borrará.
14.2 Malas prácticas (y por qué son un error, no una preferencia)
- Sincronizar sobre
Stringliterales o envoltorios.synchronized ("cuenta-" + id)es un bug garantizado: los literales se internan y losIntegerpequeños se cachean, así que compartes monitor con código ajeno de toda la JVM (y con dos ids distintos que produzcan la misma cadena no compartes nada). Usa un objeto dedicado oConcurrentHashMap<String, Object>concomputeIfAbsent. synchronizedsobrethisen una clase pública. Cualquiera puede hacersynchronized (miObjeto)desde fuera y bloquear tu clase. Usa unprivate final Object candado = new Object();.- Tragarse
InterruptedException.catch (InterruptedException e) {}rompe la cancelación de toda la aplicación y hace queshutdownNowno funcione. Thread.sleeppara sincronizar (ni en tests ni en producción). Es una carrera con el reloj: siempre pierdes en la máquina de otro.wait/notifyen código nuevo. Requieresynchronized, buclewhilecontra spurious wakeups ynotifyAll.BlockingQueue,CountDownLatch,SemaphoreoConditionhacen lo mismo sin trampas.Executors.newFixedThreadPool/newCachedThreadPoolen producción. Cola ilimitada (OOM silencioso) o hilos ilimitados (OOM ruidoso). Construye elThreadPoolExecutor.- E/S bloqueante en el
commonPool(o sea, enparallelStreamy ensupplyAsyncsin executor). Bloqueas el pool que comparte toda la JVM. - Pooling de virtual threads. Es como reciclar folios: gratis de crear, y al agruparlos vuelves a poner el límite que querías eliminar.
- Doble comprobación sin
volatile. El clásico DCL roto: puedes devolver una referencia a un objeto a medio construir. - Locks anidados sin orden global. Es la receta exacta del deadlock. Si no
puedes ordenar, usa
tryLockcon timeout. - Llamar a código ajeno con el lock tomado (alien method): un
listener, un
equalssobrescrito o una llamada HTTP dentro de unsynchronizedpueden bloquear tu monitor indefinidamente. ThreadLocalsinremove()en pools. Fuga de memoria y, peor, datos de un usuario visibles para el siguiente.SimpleDateFormat,RandomoCalendarcompartidos. No son thread-safe. UsaDateTimeFormatteryThreadLocalRandom.@Transactional+@Asyncen el mismo método. La transacción vive en unThreadLocal: el hilo nuevo no la hereda. Y en general, no ejecutes una tarea asíncrona que dependa de la transacción del llamante.- Usar el estado de un singleton de Spring como si fuera de la petición.
Un
@Servicees único y compartido: un campo mutable ahí es una fuga de datos entre usuarios. - “Arreglar” un fallo intermitente con un
sleepo con mássynchronized. Si no sabes qué invariante se rompía, no lo has arreglado: lo has escondido.
15 · Errores comunes y cómo solucionarlos
| Síntoma / error | Causa real | Solución |
|---|---|---|
| El contador da menos de lo esperado | i++ no es atómico (leer-sumar-escribir) |
AtomicInteger, LongAdder o synchronized |
Un bucle no termina aunque cambies el boolean |
Falta de visibilidad: el JIT lo ha izado a un registro | volatile (o AtomicBoolean, o interrupción) |
| Funciona en el portátil y falla en producción | Más núcleos, otro JIT y más carga exponen la carrera latente | Buscar el estado mutable compartido; no “funciona en local” como prueba |
| La aplicación se cuelga sin CPU | Deadlock, o pool agotado esperando un recurso | 3 volcados con jcmd Thread.print; buscar Found one Java-level deadlock |
OutOfMemoryError: unable to create native thread |
Fuga de hilos: se crea un executor por petición o newCachedThreadPool |
Un pool por dominio, creado una vez; comprobar jvm.threads.live |
OutOfMemoryError: Java heap space con la app “lenta” |
Cola ilimitada de un newFixedThreadPool acumulando tareas |
Cola acotada + RejectedExecutionHandler que devuelva 503 |
RejectedExecutionException |
Pool saturado, o se envían tareas después de shutdown() |
Dimensionar/degradar; y no enviar tras el cierre (comprobar el ciclo de vida) |
| Una excepción desaparece sin traza | submit() guarda la excepción en el Future; nadie llama a get() |
execute() para fire-and-forget, o get()/whenComplete siempre |
Una tarea @Scheduled deja de ejecutarse |
Lanzó una excepción: scheduleAtFixedRate cancela la tarea para siempre |
try/catch (Throwable) alrededor de todo el cuerpo |
| Datos de otro usuario en la respuesta | ThreadLocal sin remove(), o campo mutable en un singleton |
remove() en finally; estado por petición en parámetros |
| Fuga de memoria que crece con el tráfico | ThreadLocal retenido por hilos de pool que nunca mueren |
remove(); revisar filtros y interceptors; volcado de heap |
ConcurrentModificationException |
Modificar una colección mientras se itera (aunque sea en un solo hilo) | Iterator.remove(), removeIf, copia, o colección concurrente |
NullPointerException en ConcurrentHashMap |
La clase prohíbe claves y valores null por diseño |
getOrDefault, un valor centinela, o Optional como valor |
HikariPool-1 - Connection is not available |
Conexiones sin cerrar, transacciones largas o pool pequeño frente a la concurrencia | try-with-resources; acortar transacciones; alinear pool y concurrencia real |
| Los virtual threads no mejoran nada | El trabajo es CPU-bound, o hay un pool/semáforo estrangulando antes | Medir con wall-clock; buscar el cuello real (BD, API externa) |
| Con virtual threads, la latencia empeora a ratos | Pinning por synchronized bloqueante o JNI |
JFR jdk.VirtualThreadPinned; cambiar a ReentrantLock; JDK 24+ |
parallelStream() hunde toda la aplicación |
E/S bloqueante en el commonPool, compartido por la JVM |
No hacer E/S ahí; usar un executor propio o virtual threads |
La transacción no se propaga a @Async |
El contexto transaccional vive en un ThreadLocal del hilo llamante |
Rediseñar: transacción completa dentro del método asíncrono (patrón outbox) |
| El test pasa en local y falla en CI (o 1 de cada 50) | Thread.sleep como sincronización; entrelazado distinto |
Latches, Awaitility, ejecutor síncrono inyectado, @RepeatedTest |
| Un job programado se ejecuta N veces | N réplicas del servicio con el mismo @Scheduled |
ShedLock o SKIP LOCKED + idempotencia (sección 11.3) |
| Escalabilidad plana al añadir núcleos | Sección crítica dominante (ley de Amdahl) o false sharing | Reducir el lock, particionar el estado, LongAdder, @Contended |
IllegalMonitorStateException |
wait/notify sin poseer el monitor, o unlock en otro hilo |
Llamar dentro de synchronized; unlock() en finally del mismo hilo |
| Reintentos que duplican cobros o correos | Operación no idempotente reintentada tras un timeout | Clave de idempotencia + deduplicación en BD |
| Todo el sistema cae porque un servicio remoto está lento | Sin timeout, sin bulkhead y sin circuit breaker: los hilos se agotan | Timeouts agresivos, pools separados, Resilience4j (módulo 08) |
16 · FAQ de entrevista
volatile vs synchronized vs atómicos: ¿cuándo cada uno?
volatile da visibilidad y orden, pero no atomicidad:
sirve para banderas y para publicar una referencia ya construida, nunca para
contador++. Los atómicos (AtomicInteger,
AtomicReference) añaden operaciones compuestas atómicas sobre una sola
variable mediante CAS, sin bloquear. synchronized (y ReentrantLock) es lo
único que te permite hacer atómica una operación sobre varias variables a la vez,
es decir, mantener un invariante que las relacione.
Regla práctica: bandera → volatile; contador o referencia única →
atómico; invariante entre dos o más campos → lock.
Diferencia entre wait() y sleep().
Object.wait() es un método de instancia que exige poseer el monitor de ese objeto,
libera el monitor mientras espera y se despierta con
notify/notifyAll (o espuriamente, de ahí el while).
Thread.sleep() es estático, no suelta ningún lock y se despierta solo
por tiempo o por interrupción. Dormir con un lock tomado es una de las mejores formas de hundir el
rendimiento. En código nuevo, ninguno de los dos: usa Condition.await(),
BlockingQueue o CountDownLatch.
sleep() vs yield() vs onSpinWait().
sleep(ms) pasa a TIMED_WAITING durante un tiempo garantizado (como
mínimo). Thread.yield() es una sugerencia al planificador para ceder el turno:
el hilo sigue RUNNABLE y la JVM puede ignorarla, así que no sirve para coordinar nada,
solo como ayuda en pruebas o en spin largos. Thread.onSpinWait() (Java 9+) es
una pista al procesador dentro de un bucle de espera activa muy corto (emite PAUSE en
x86), útil solo en código de altísimo rendimiento. Ninguno es un mecanismo de sincronización.
¿Cómo se evitan los deadlocks?
Un deadlock necesita cuatro condiciones simultáneas (exclusión mutua, retención y espera, no preemción, espera circular); basta romper una:
- Orden global de adquisición. Si todo el código toma los locks siempre en el
mismo orden (por ejemplo, por
ido porSystem.identityHashCode), no puede haber ciclo. Es la solución definitiva. - Un solo lock de granularidad mayor, cuando la contención lo permita.
tryLockcon timeout y reintento con backoff: rompe la “no preemción”.- No llamar a código ajeno con el lock tomado (ni HTTP, ni listeners).
- Sin locks: inmutabilidad, colas y confinamiento por hilo.
Y para detectarlo: jcmd <pid> Thread.print lo dice literalmente, o
ThreadMXBean.findDeadlockedThreads() como métrica de salud.
¿Qué es CAS y qué relación tiene con los locks?
Compare-And-Swap es una instrucción del procesador (LOCK CMPXCHG en x86,
LDXR/STXR en ARM) que hace atómicamente: “si esta dirección vale
esperado, escribe nuevo; si no, dime que has fallado”. Sobre ella se
construyen los atómicos: incrementAndGet es un bucle
leo → calculo → intento CAS → si falla, repito. Es no bloqueante (nadie
duerme; siempre hay un hilo que progresa), así que con poca contención es mucho más rápido que un
lock. Con mucha contención, los reintentos se desperdician y un lock —o
LongAdder, que reparte el contador— puede rendir mejor. Sus limitaciones: solo abarca
una variable y sufre el problema ABA.
¿Por qué ConcurrentHashMap no admite null?
Porque haría ambiguo el resultado de get(): si devuelve
null, ¿la clave no está o está asociada a null? En un
HashMap puedes resolverlo con containsKey(), pero en un mapa concurrente
esa comprobación no vale: entre el get y el containsKey otro hilo puede
haber cambiado el mapa, así que no existe respuesta correcta. Doug Lea decidió eliminar la
ambigüedad de raíz. Además, permitir null rompería
computeIfAbsent/merge, que usan null como “ausente”. Si
necesitas representar ausencia, usa un valor centinela o un Optional como valor.
¿Qué es el false sharing y cómo se evita?
La caché de la CPU no trabaja con bytes sino con líneas de 64 bytes. Si dos hilos en núcleos distintos escriben en dos variables diferentes que caen en la misma línea, el protocolo de coherencia invalida la línea entera en el otro núcleo a cada escritura: los hilos no comparten datos, pero sí el coste. Se puede perder un orden de magnitud.
Se evita separando las variables calientes (padding), con
@jdk.internal.vm.annotation.Contended (interno del JDK) o, en la práctica, usando las
clases que ya lo hacen: LongAdder reparte el contador en celdas alineadas justo por
esto. Es la explicación de la mayoría de “no escala aunque no haya locks”.
¿Cuántos hilos necesito?
Depende de qué hacen. Si es CPU-bound:
availableProcessors() (±1); más hilos solo añaden cambios de contexto. Si es
IO-bound: núcleos × utilizaciónObjetivo × (1 + espera/cómputo), y
valida con la ley de Little (concurrencia = tasa × latencia). Con
virtual threads, la pregunta cambia: no dimensionas hilos, dimensionas el
recurso escaso (conexiones a la BD, cuota de la API) con un Semaphore.
Y la respuesta honesta en una entrevista: “calculo un punto de partida con esa fórmula y luego lo ajusto midiendo latencia p99 y saturación del recurso, porque el número real depende del cuello de botella, no del cálculo”.
¿Los virtual threads sustituyen a la programación reactiva?
Para el caso mayoritario —un servicio que atiende peticiones y llama a base de datos y a otros
servicios— sí: dan la misma escalabilidad con código secuencial, depurable y con pilas legibles, que
era el 90 % del motivo para usar reactivo. Lo que no sustituyen es el modelo de
flujos: backpressure de extremo a extremo, streaming, operadores de composición
temporal (window, buffer, debounce, retryWhen) y
la programación orientada a eventos. Si tu problema es “procesar un río de datos con control de
flujo”, Reactor sigue siendo la herramienta. Si es “atender muchas peticiones bloqueantes”, virtual
threads.
¿Es thread-safe un singleton de Spring?
Spring garantiza que se crea una sola vez de forma segura, pero
no que su código sea thread-safe. Un @Service es una única
instancia atendiendo todas las peticiones a la vez, así que cualquier campo mutable de
instancia es estado compartido. Si guardas ahí el usuario actual o un acumulador, tienes
una fuga de datos entre usuarios.
Lo correcto: beans sin estado, con dependencias final inyectadas por
constructor; el estado de la petición viaja en parámetros y variables locales. Si de verdad
necesitas estado por petición, usa @Scope("request") o pásalo explícitamente (y en
Java 21+, ScopedValue antes que ThreadLocal).
Diferencia entre Runnable y Callable.
Runnable.run() no devuelve nada y no puede lanzar excepciones comprobadas.
Callable<V>.call() devuelve V y sí puede lanzar
Exception. En un ExecutorService, execute(Runnable) no da
resultado y las excepciones van al UncaughtExceptionHandler;
submit(…) devuelve un Future y captura la excepción dentro de él. Nota
útil: Executors.callable(runnable) adapta uno al otro.
¿Qué pasa si una tarea lanza una excepción en un ExecutorService?
Depende de cómo la enviaste, y es una de las trampas más habituales:
execute(runnable): la excepción sale del hilo, elUncaughtExceptionHandlerla registra y el pool reemplaza el hilo muerto. Se ve en el log.submit(…): la excepción se guarda en elFuturey no aparece en ningún log. Si nadie llama aget(), el fallo es invisible: es la causa clásica de “la tarea no hizo nada y no hay error”.scheduleAtFixedRate: la tarea se cancela para siempre. Envuélvela siempre entry/catch (Throwable).CompletableFuture: la excepción viaja por la cadena; sinexceptionally/handle/whenCompletese pierde igual.
¿Por qué ThreadLocal es peligroso en un pool?
Porque los hilos del pool no mueren: se reutilizan durante toda la vida de la
aplicación. El valor que dejaste ahí sigue asociado al hilo cuando llega la siguiente petición, con
dos consecuencias: (1) fuga de memoria —la entrada la retiene el
Thread, y la clave es débil pero el valor no—; y (2)
fuga de datos, mucho peor: la petición del usuario B ve el contexto del usuario A.
La única cura es remove() en un finally, en el filtro o
interceptor que lo puso. Con virtual threads el problema de reutilización desaparece
(cada tarea es un hilo nuevo), pero el coste de memoria por millones de hilos sigue importando, y
ScopedValue es la alternativa correcta.
¿Qué es el pinning y por qué debería preocuparme?
Un virtual thread se ejecuta “montado” sobre un hilo de plataforma (carrier). Cuando se
bloquea, normalmente se desmonta y libera el carrier. Pero si en ese momento hay una operación que
no se puede desmontar —clásicamente un synchronized bloqueante, y siempre en un
frame nativo (JNI)— el virtual thread queda fijado al carrier, que se
queda bloqueado. Con el pool de carriers por defecto del tamaño del número de núcleos, unos pocos
hilos fijados pueden estrangular la aplicación.
Se detecta con -Djdk.tracePinnedThreads=full (útil, aunque marcado como
obsoleto en versiones recientes) y, mejor, con el evento JFR
jdk.VirtualThreadPinned. La solución era sustituir synchronized por
ReentrantLock; desde JDK 24, con JEP 491, los bloques synchronized ya no
fijan el hilo, así que en JDK 24+ el problema queda prácticamente reducido al código nativo.
¿Qué garantiza exactamente happens-before?
No es “ocurrir antes en el tiempo”, sino una relación de visibilidad y orden: si A happens-before B, todo lo que A escribió es visible para B, y B no puede observar A a medias. Es la única garantía que da el modelo de memoria; sin una relación así entre dos accesos a la misma variable (uno de ellos escritura), tienes una carrera de datos y el resultado no está definido, por mucho que “funcione”.
Las fuentes principales: el orden del programa dentro de un hilo; liberar un monitor antes de
adquirirlo; escribir un volatile antes de leerlo; Thread.start() respecto
a todo lo anterior; el final de un hilo respecto a join(); enviar una tarea a un
executor respecto a su ejecución; y la transitividad, que es lo que hace útil todo lo demás.
¿Por qué está eliminado Thread.stop()?
Porque lanzaba un ThreadDeath en un punto arbitrario del hilo:
podía interrumpir a mitad de una actualización de dos campos relacionados, dejando el objeto en un
estado inconsistente, y podía soltar los monitores sin restaurar invariantes. No hay forma de
escribir código robusto frente a eso, así que el método se deprecó muy pronto y en Java 20 se
convirtió en un lanzador de UnsupportedOperationException. La única cancelación
correcta es cooperativa: interrupt() más código que compruebe el flag
y propague InterruptedException.
Diferencia entre thenApply y thenCompose.
thenApply(f) transforma el valor con una función normal:
T → U, y devuelve CompletableFuture<U>.
thenCompose(f) encadena otra operación asíncrona:
T → CompletableFuture<U>, y aplana el resultado. Si usas
thenApply con una función que devuelve un futuro, obtienes
CompletableFuture<CompletableFuture<U>>, que es el equivalente de
Optional<Optional<T>> o del map vs flatMap de los
streams: la señal de que te has equivocado de operador.
¿Es SimpleDateFormat thread-safe? ¿Y StringBuilder?
Ninguno de los dos. SimpleDateFormat guarda estado interno de parsing y,
compartido entre hilos, produce fechas silenciosamente incorrectas (un bug perfecto: sin excepción y
difícil de reproducir). La solución moderna es DateTimeFormatter, que es inmutable y
thread-safe. StringBuilder tampoco lo es —su hermano sincronizado es
StringBuffer—, pero en la práctica casi siempre es una variable local, y entonces no
hay problema: el confinamiento por hilo es una forma válida de seguridad. Otros
sospechosos habituales: Random (usa ThreadLocalRandom),
Calendar, HashMap y los Collections.unmodifiable*, que solo
protegen de la modificación, no de la concurrencia.
17 · Ejercicios y retos
17.1 Ocho ejercicios guiados
17.2 Cinco retos
18 · Resumen y recursos
18.1 Quince ideas para llevarte
- Concurrencia ≠ paralelismo. Concurrencia es estructurar el programa para gestionar muchas cosas a la vez; paralelismo es ejecutarlas simultáneamente. Casi siempre necesitas lo primero.
- El problema no son los hilos, es la memoria compartida mutable. Quítala y la concurrencia se vuelve fácil.
- Sin una relación happens-before no hay garantía ninguna. “A mí me funciona” no es evidencia: los bugs de concurrencia son probabilísticos.
volatileda visibilidad, no atomicidad. Los atómicos dan atomicidad sobre una variable. Los locks son lo único que protege un invariante entre varias.i++son tres operaciones. Ahí empieza todo.- Nunca te tragues
InterruptedException. Propágala o restaura el flag: es el contrato de cancelación de la plataforma. - Los deadlocks se evitan por diseño con un orden global de adquisición, no con
suerte ni con
sleep. - Cola acotada, nombres y política de rechazo en todos los pools. Las factorías
de
Executorsson cómodas y peligrosas. - Nunca hagas E/S bloqueante en el
commonPoolni, por tanto, en unparallelStream. - Timeouts en todo. Toda espera sin plazo es una caída pendiente de fecha.
- Los virtual threads hacen viable el estilo bloqueante a gran escala, pero no aceleran la CPU ni multiplican tu base de datos: mueven el cuello de botella, no lo eliminan.
- Con virtual threads no se hace pooling: se limita el recurso escaso con
un
Semaphore. - La concurrencia estructurada convierte las tareas concurrentes en un bloque con alcance, con cancelación y propagación de errores automáticas. Es el futuro del código concurrente en Java.
- Reactivo sigue teniendo su sitio —flujos y backpressure—, pero ya no es el precio obligatorio de la escalabilidad.
- Sin métricas y volcados de hilos, la concurrencia es una caja negra. Instrumenta antes de que arda: hilos vivos, cola, rechazos y conexiones pendientes.
18.2 Recursos para profundizar
Libros
- Java Concurrency in Practice — Brian Goetz y otros (2006). Sigue siendo
el libro: envejeció en las APIs (no hay
CompletableFutureni Loom), pero el modelo mental, el modelo de memoria y las fórmulas de dimensionado no han envejecido nada. Los capítulos 2-5 y 10-11 son obligatorios. - Java Concurrency Stress (jcstress) y los artículos de Aleksey Shipilëv sobre el JMM: la referencia técnica cuando quieras saber qué está realmente garantizado.
- Effective Java — Joshua Bloch, capítulo 11 (“Concurrency”): reglas cortas y demoledoras sobre sincronización y ejecutores.
- The Art of Multiprocessor Programming — Herlihy y Shavit: la teoría (CAS, estructuras no bloqueantes, linealizabilidad) si quieres el nivel académico.
Especificaciones, JEPs y podcasts
- JEP 444 · Virtual Threads (estable en Java 21) y las anteriores 425 y 436: explican el diseño, no solo la API.
- JEP 453 / 462 / 480 / 499 · Structured Concurrency y JEP 429 / 446 / 464 / 481 · Scoped Values: el estado real (y cambiante) de estas dos APIs. Consulta siempre la JEP de tu versión antes de usarlas.
- JEP 491 · Synchronize Virtual Threads without Pinning (Java 24): la que elimina la principal fuente de pinning.
- Javadoc de
java.util.concurrent: la documentación del paquete es material didáctico de primera, no una lista de métodos. Lee las notas deConcurrentHashMapy deThreadPoolExecutorenteras al menos una vez. - JLS, capítulo 17 (“Threads and Locks”): la definición formal de happens-before. Denso, pero es la fuente de la verdad.
- Inside Java Podcast y Inside Java Newscast (José Paumard, Ron Pressler, Alan Bateman): los episodios sobre Loom explican decisiones de diseño que no están en ninguna documentación.
- JDK Mission Control y async-profiler: aprende a usarlos antes del incidente.
@Async, @Scheduled, el
TaskExecutor y la configuración de virtual threads. Las transacciones y los pools de
conexiones están en módulo 05, y la resiliencia distribuida
(circuit breakers, bulkheads, reintentos) en
módulo 08. Si necesitas repasar
final, la inmutabilidad o los record, vuelve a
módulo 01; para CompletableFuture en el contexto de las
novedades del lenguaje, a módulo 02.