Guía completa de concurrencia en Java: de hilos tradicionales a CompletableFuture y ForkJoin
Objetivo: entender patrones y APIs de concurrencia en Java, cuándo usar cada una, ejemplos completos y un proyecto mínimo para probar buenas prácticas. Va dirigido a desarrolladores Java con experiencia intermedia que quieren elevar su manejo de concurrencia y paralelismo.
1. Conceptos rápidos y cuándo elegir qué
- Hilo (Thread): unidad básica, evita crear muchos hilos manualmente.
- ExecutorService / ThreadPool: usar para gestionar recursos y evitar overhead.
- CompletableFuture: programación asíncrona con composición fácil; ideal para IO-bound y pipelines.
- ForkJoinPool / RecursiveTask: para tareas CPU-bound divididas (divide & conquer).
- Síncronización: synchronized, ReentrantLock, volatile, atomics según granularidad y rendimiento.
2. Buenas prácticas resumidas (antes de ver código)
- No crear hilos sin control; usa pools.
- Preferir estructuras no bloqueantes (CompletableFuture, atomics) cuando sea posible.
- Minimizar alcance de locks y evitar locks anidados para prevenir deadlocks.
- Medir y perfilar: los cambios de concurrencia afectan latencia y throughput.
- Usar mecanismos de coordinación (CountDownLatch, Semaphore, Phaser) cuando necesites sincronización entre tareas.
3. Proyecto mínimo: procesador paralelo de archivos (Resumen de la estructura)
Objetivo del proyecto: leer múltiples archivos de texto y contar la frecuencia de palabras en paralelo, combinando resultados de forma segura y eficiente.
/wordcount-parallel ├─ pom.xml └─ src/main/java/com/example/wordcount ├─ App.java ├─ FileScanner.java ├─ WordCounter.java └─ MergeUtils.java
Por qué esta estructura
Separamos responsabilidades: escaneo de archivos, conteo por archivo y fusión de resultados. Esto facilita paralelizar por archivo y usar distintas estrategias (CompletableFuture o ForkJoin).
4. Código completo esencial
pom.xml (solo dependencias básicas)
<project xmlns="http://maven.apache.org/POM/4.0.0" ...>
<modelVersion>4.0.0</modelVersion>
<groupId>com.example</groupId>
<artifactId>wordcount-parallel</artifactId>
<version>1.0-SNAPSHOT</version>
<properties>
<maven.compiler.source>11</maven.compiler.source>
<maven.compiler.target>11</maven.compiler.target>
</properties>
</project>
App.java (orquestador)
package com.example.wordcount;
import java.nio.file.Path;
import java.nio.file.Paths;
import java.util.Map;
import java.util.concurrent.*;
public class App {
public static void main(String[] args) throws Exception {
if (args.length == 0) {
System.out.println("Usar: java -jar wordcount-parallel ");
return;
}
Path dir = Paths.get(args[0]);
// Pool tamaño razonable: number of cores para CPU-bound, más para IO-bound
int cores = Runtime.getRuntime().availableProcessors();
ExecutorService pool = Executors.newFixedThreadPool(Math.max(2, cores));
try {
// Escanear archivos
var files = FileScanner.listTextFiles(dir);
// Estrategia 1: CompletableFuture por archivo
var futures = files.stream()
.map(path -> CompletableFuture.supplyAsync(() -> WordCounter.countWords(path), pool))
.toArray(CompletableFuture[]::new);
CompletableFuture.allOf(futures).join();
// Combinar resultados
ConcurrentHashMap<String, Long> merged = new ConcurrentHashMap<>();
for (var f : futures) {
Map<String, Long> map = (Map<String, Long>) f.join();
MergeUtils.mergeInto(merged, map);
}
// Mostrar top 20
merged.entrySet().stream()
.sorted((a,b) -> Long.compare(b.getValue(), a.getValue()))
.limit(20)
.forEach(e -> System.out.println(e.getKey() + " -> " + e.getValue()));
} finally {
pool.shutdown();
pool.awaitTermination(10, TimeUnit.SECONDS);
}
}
}
FileScanner.java
package com.example.wordcount;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import java.util.stream.Collectors;
public class FileScanner {
public static List<Path> listTextFiles(Path dir) throws IOException {
try (var stream = Files.walk(dir)) {
return stream
.filter(Files::isRegularFile)
.filter(p -> p.toString().endsWith(".txt"))
.collect(Collectors.toList());
}
}
}
WordCounter.java
package com.example.wordcount;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.HashMap;
import java.util.Map;
import java.util.regex.Pattern;
public class WordCounter {
private static final Pattern SPLIT = Pattern.compile("\\W+");
public static Map<String, Long> countWords(Path path) {
try {
var map = new HashMap<String, Long>();
Files.lines(path)
.flatMap(line -> java.util.Arrays.stream(SPLIT.split(line.toLowerCase())))
.filter(s -> !s.isBlank())
.forEach(w -> map.merge(w, 1L, Long::sum));
return map;
} catch (IOException e) {
throw new RuntimeException(e);
}
}
}
MergeUtils.java
package com.example.wordcount;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
public class MergeUtils {
public static void mergeInto(ConcurrentHashMap<String, Long> target, Map<String, Long> src) {
src.forEach((k,v) -> target.merge(k, v, Long::sum));
}
}
5. Alternativa: usar ForkJoin para archivos muy grandes (divide & conquer)
Si cada archivo es enorme y se puede dividir eficientemente, usa ForkJoinPool y RecursiveTask. Ejemplo: dividir un archivo en rangos de líneas y procesar recursivamente. ForkJoin es más apropiado para tareas CPU-bound con trabajo dividible.
6. Errores comunes y cómo evitarlos
- Crear pools sin límites: usar Executors.newCachedThreadPool sin entender puede agotar recursos.
- Bloqueos largos dentro de un pool compartido: si un task bloquea IO y saturas el pool, todo se detiene — separar pools para IO y CPU.
- Uso incorrecto de mutable shared state: preferir ConcurrentHashMap, AtomicLong o reducir la necesidad de locks (partitioning, map-reduce).
- Olvidar shutdown del ExecutorService: siempre cerrar pools y manejar timeouts.
7. Debugging y profiling
- Usa jstack para capturar hilos en bloqueo y detectar deadlocks.
- Flight Recorder y Java Mission Control para hotspots y bloqueo de locks.
- Logs con thread id/context para reproducir condiciones concurrentes.
8. Optimización rápida (3 pasos)
- Medir: identifica CPU-bound vs IO-bound.
- Elegir pool correcto: fixed pool para CPU-bound (size ~ cores), cached o larger pool para IO-bound.
- Reduce contención: sharding de datos, usar ConcurrentHashMap con compute/merge, evitar locks globales.
9. Patrón recomendado para pipelines asíncronos
Combina CompletableFuture con un pool dedicado para transformaciones pesadas. Ejemplo:
CompletableFuture.supplyAsync(() -> readIO(path), ioPool)
.thenApplyAsync(data -> parseAndTransform(data), cpuPool)
.thenAcceptAsync(result -> persist(result), ioPool);
Así separas recursos y evitas que operaciones CPU saturen los hilos que deben atender IO.
10. Seguridad y consistencia
Para datos críticos: utiliza transacciones en la capa persistencia; evita lógica crítica en callbacks no controlados. Si necesitas consistencia fuerte en memoria, opta por locks bien diseñados o estructuras transaccionales externas (databases, queues).
¿Qué seguir aprendiendo?
Explora Project Loom (fibers/virtual threads) en versiones recientes de Java para simplificar concurrencia; reevalúa pools y backpressure cuando uses virtual threads. Empieza probando la versión preview de Loom y comparando latencia y uso de memoria con tu implementación basada en ThreadPool.
Consejo avanzado: si tu aplicación tiene mezcla de IO y CPU intensivo, segmenta executors por tipo de trabajo, instrumenta cada pool con métricas (task latency, queue size) y adapta dinámicamente el sizing en función de la carga.
Advertencia: cambiar a soluciones no bloqueantes sin pruebas puede introducir condiciones de carrera difíciles de reproducir — añade tests de concurrencia y escenarios de carga antes de desplegar en producción.
Siguiente paso: implementa la versión con ForkJoin y una con CompletableFuture + pools separados, profila ambas con JFR bajo carga representativa y compara throughput y latencia.
¿Quieres comentar?
Inicia sesión con Telegram para participar en la conversación