Guía completa de concurrencia en Java: de hilos tradicionales a CompletableFuture y ForkJoin

java Guía completa de concurrencia en Java: de hilos tradicionales a CompletableFuture y ForkJoin

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)

  1. Medir: identifica CPU-bound vs IO-bound.
  2. Elegir pool correcto: fixed pool para CPU-bound (size ~ cores), cached o larger pool para IO-bound.
  3. 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.

Comentarios
¿Quieres comentar?

Inicia sesión con Telegram para participar en la conversación


Comentarios (0)

Aún no hay comentarios. ¡Sé el primero en comentar!

Iniciar Sesión