Saltar al contenido principal

Proyecto: procesamiento concurrente controlado

Construiremos una aplicación Java 17 que procesa los .txt de una carpeta con un pool acotado. Cada tarea devuelve un resultado independiente; el coordinador agrega resultados, registra fallos y continúa el lote.

Requisitos​

  • contar líneas y palabras por archivo;
  • devolver nombre y estado OK o ERROR;
  • usar ExecutorService, Callable y Future, no un hilo por fichero;
  • aplicar timeout y cancelación cooperativa;
  • tolerar fallos parciales;
  • cerrar siempre el executor;
  • usar logging y tests deterministas.

No se garantiza el orden de finalización. El listado final se ordenará por nombre para ofrecer una salida estable, no porque las tareas terminen en ese orden.

Estructura​

procesador-textos/
├── pom.xml
├── README.md
└── src
├── main/java/es/skilly/java5/
│ ├── Main.java
│ ├── ProcesadorLote.java
│ ├── TareaArchivo.java
│ └── ResultadoArchivo.java
└── test/java/es/skilly/java5/
└── ProcesadorLoteTest.java

pom.xml​

<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>es.skilly</groupId><artifactId>procesador-textos</artifactId><version>1.0.0</version>
<properties><maven.compiler.release>17</maven.compiler.release><project.build.sourceEncoding>UTF-8</project.build.sourceEncoding></properties>
<dependencies>
<dependency><groupId>org.slf4j</groupId><artifactId>slf4j-api</artifactId><version>2.0.17</version></dependency>
<dependency><groupId>org.slf4j</groupId><artifactId>slf4j-simple</artifactId><version>2.0.17</version><scope>runtime</scope></dependency>
<dependency><groupId>org.junit.jupiter</groupId><artifactId>junit-jupiter</artifactId><version>5.14.2</version><scope>test</scope></dependency>
</dependencies>
<build><plugins>
<plugin><groupId>org.apache.maven.plugins</groupId><artifactId>maven-compiler-plugin</artifactId><version>3.13.0</version></plugin>
<plugin><groupId>org.apache.maven.plugins</groupId><artifactId>maven-surefire-plugin</artifactId><version>3.5.4</version></plugin>
<plugin><groupId>org.codehaus.mojo</groupId><artifactId>exec-maven-plugin</artifactId><version>3.5.0</version><configuration><mainClass>es.skilly.java5.Main</mainClass></configuration></plugin>
</plugins></build>
</project>

Modelo de resultado​

package es.skilly.java5;

import java.nio.file.Path;

public class ResultadoArchivo {
private final Path archivo;
private final int lineas;
private final int palabras;
private final boolean correcto;
private final String error;

private ResultadoArchivo(Path archivo, int lineas, int palabras, boolean correcto, String error) {
this.archivo = archivo; this.lineas = lineas; this.palabras = palabras;
this.correcto = correcto; this.error = error;
}
public static ResultadoArchivo ok(Path a, int l, int p) { return new ResultadoArchivo(a,l,p,true,null); }
public static ResultadoArchivo error(Path a, String e) { return new ResultadoArchivo(a,0,0,false,e); }
public Path getArchivo() { return archivo; }
public int getLineas() { return lineas; }
public int getPalabras() { return palabras; }
public boolean isCorrecto() { return correcto; }
public String getError() { return error; }
}

Tarea por archivo​

package es.skilly.java5;

import java.nio.file.*;
import java.util.List;
import java.util.concurrent.Callable;

public class TareaArchivo implements Callable<ResultadoArchivo> {
private final Path archivo;
public TareaArchivo(Path archivo) { this.archivo = archivo; }

@Override public ResultadoArchivo call() throws InterruptedException {
try {
if (Thread.currentThread().isInterrupted()) {
throw new InterruptedException("Tarea interrumpida");
}
List<String> lineas = Files.readAllLines(archivo);
int palabras = 0;
for (String linea : lineas) {
if (Thread.currentThread().isInterrupted()) throw new InterruptedException("Tarea interrumpida");
String limpia = linea.trim();
if (!limpia.isEmpty()) palabras += limpia.split("\\s+").length;
}
return ResultadoArchivo.ok(archivo, lineas.size(), palabras);
} catch (InterruptedException e) {
throw e;
} catch (Exception e) {
return ResultadoArchivo.error(archivo, e.getMessage());
}
}
}

Cada Callable posee sus datos y devuelve un resultado. Evitamos que workers muten una lista compartida. La interrupción es cooperativa: la tarea consulta la señal antes y durante el procesamiento, pero no prometemos que Files.readAllLines pueda cancelarse inmediatamente mientras está leyendo.

Coordinador​

package es.skilly.java5;

import java.nio.file.*;
import java.time.Duration;
import java.util.*;
import java.util.concurrent.*;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

public class ProcesadorLote implements AutoCloseable {
private static final Logger LOG = LoggerFactory.getLogger(ProcesadorLote.class);
private final ExecutorService executor;

public ProcesadorLote(int hilos) {
if (hilos < 1) throw new IllegalArgumentException("hilos debe ser positivo");
executor = Executors.newFixedThreadPool(hilos);
}

public List<ResultadoArchivo> procesar(Path carpeta, Duration timeout) throws InterruptedException {
List<Path> archivos;
try (var stream = Files.list(carpeta)) {
archivos = stream.filter(p -> p.toString().endsWith(".txt")).sorted().toList();
} catch (Exception e) {
throw new IllegalArgumentException("No se puede listar la carpeta", e);
}
return procesarArchivos(archivos, timeout);
}

List<ResultadoArchivo> procesarArchivos(List<Path> archivos, Duration timeout) throws InterruptedException {
LOG.info("Lote iniciado: {} archivos", archivos.size());
Map<Path, Callable<ResultadoArchivo>> tareas = new LinkedHashMap<>();
for (Path archivo : archivos) tareas.put(archivo, new TareaArchivo(archivo));
return procesarTareas(tareas, timeout);
}

List<ResultadoArchivo> procesarTareas(Map<Path, Callable<ResultadoArchivo>> tareas, Duration timeout) throws InterruptedException {
Map<Path, Future<ResultadoArchivo>> futures = new LinkedHashMap<>();
tareas.forEach((archivo, tarea) -> futures.put(archivo, executor.submit(tarea)));

List<ResultadoArchivo> resultados = new ArrayList<>();
for (var entrada : futures.entrySet()) {
try {
ResultadoArchivo r = entrada.getValue().get(timeout.toMillis(), TimeUnit.MILLISECONDS);
resultados.add(r);
if (r.isCorrecto()) LOG.info("Procesado: {}", entrada.getKey());
else LOG.warn("Fallido: {} - {}", entrada.getKey(), r.getError());
} catch (TimeoutException e) {
entrada.getValue().cancel(true);
LOG.warn("Timeout y cancelación solicitada: {}", entrada.getKey());
resultados.add(ResultadoArchivo.error(entrada.getKey(), "timeout"));
} catch (ExecutionException e) {
Throwable causa = e.getCause();
LOG.warn("Fallo procesando {}: {}", entrada.getKey(), causa.toString());
resultados.add(ResultadoArchivo.error(entrada.getKey(), causa.toString()));
} catch (CancellationException e) {
LOG.warn("Tarea cancelada: {}", entrada.getKey());
resultados.add(ResultadoArchivo.error(entrada.getKey(), "cancelado"));
} catch (InterruptedException e) {
futures.values().forEach(f -> f.cancel(true));
Thread.currentThread().interrupt();
throw e;
}
}
resultados.sort(Comparator.comparing(r -> r.getArchivo().getFileName().toString()));
return resultados;
}

@Override public void close() {
executor.shutdown();
try {
if (!executor.awaitTermination(2, TimeUnit.SECONDS)) {
executor.shutdownNow();
if (!executor.awaitTermination(2, TimeUnit.SECONDS)) {
LOG.warn("El executor no terminó después de solicitar cancelación");
return;
}
}
LOG.info("Executor terminado");
} catch (InterruptedException e) {
executor.shutdownNow();
Thread.currentThread().interrupt();
LOG.warn("Cierre interrumpido; se solicitó cancelación");
}
}

boolean estaTerminado() { return executor.isTerminated(); }
}

La clase implementa AutoCloseable, no ExecutorService: así podemos usar el coordinador con try-with-resources en Java 17 y mantener el cierre explícito del executor. shutdown() solicita el cierre; solo isTerminated() confirma que todas las tareas terminaron. Los helpers con visibilidad de paquete permiten probar de forma determinista fallos y timeouts sin cambiar la API pública.

En esta implementación, timeout limita cuánto espera el coordinador por cada Future cuando llega a él. No representa un tiempo máximo absoluto desde el submit ni desde el inicio real de la tarea.

Entrada​

package es.skilly.java5;
import java.nio.file.*; import java.time.Duration;
public class Main {
public static void main(String[] args) throws Exception {
Path carpeta = Path.of(args.length == 0 ? "datos" : args[0]);
try (ProcesadorLote procesador = new ProcesadorLote(2)) {
for (ResultadoArchivo r : procesador.procesar(carpeta, Duration.ofSeconds(2))) {
System.out.printf("%s -> %s (líneas=%d, palabras=%d)%n", r.getArchivo().getFileName(), r.isCorrecto()?"OK":"ERROR: "+r.getError(), r.getLineas(), r.getPalabras());
}
}
}
}

Tests​

package es.skilly.java5;
import static org.junit.jupiter.api.Assertions.*;
import java.nio.file.*; import java.time.Duration; import java.util.*;
import java.util.concurrent.*;
import org.junit.jupiter.api.*; import org.junit.jupiter.api.io.TempDir;

class ProcesadorLoteTest {
@TempDir Path dir;

@Test void cuentaUnArchivo() throws Exception {
Files.writeString(dir.resolve("a.txt"), "uno dos\ntres\n");
try (ProcesadorLote p = new ProcesadorLote(2)) {
ResultadoArchivo r = p.procesar(dir, Duration.ofSeconds(2)).get(0);
assertAll(() -> assertTrue(r.isCorrecto()), () -> assertEquals(2,r.getLineas()), () -> assertEquals(3,r.getPalabras()));
}
}

@Test void procesaVariosSinDependerDelScheduling() throws Exception {
Files.writeString(dir.resolve("b.txt"), "b"); Files.writeString(dir.resolve("a.txt"), "a");
try (ProcesadorLote p = new ProcesadorLote(2)) {
List<ResultadoArchivo> rs = p.procesar(dir, Duration.ofSeconds(2));
assertEquals(List.of("a.txt","b.txt"), rs.stream().map(r->r.getArchivo().getFileName().toString()).toList());
}
}

@Test void loteVacio() throws Exception {
try (ProcesadorLote p = new ProcesadorLote(1)) { assertTrue(p.procesar(dir, Duration.ofSeconds(2)).isEmpty()); }
}

@Test void conservaFalloParcial() throws Exception {
Path valido=dir.resolve("valido.txt"), roto=dir.resolve("roto.txt"); Files.writeString(valido,"uno dos");
try (ProcesadorLote p = new ProcesadorLote(2)) {
List<ResultadoArchivo> rs = p.procesarArchivos(List.of(valido, roto), Duration.ofSeconds(2));
assertAll(
() -> assertEquals(1, rs.stream().filter(ResultadoArchivo::isCorrecto).count()),
() -> assertEquals(1, rs.stream().filter(r -> !r.isCorrecto()).count())
);
}
}

@Test void cierraElExecutor() {
ProcesadorLote p = new ProcesadorLote(1); p.close(); assertTrue(p.estaTerminado());
}

@Test void timeoutCancelaDeFormaCooperativa() throws Exception {
Path lento = dir.resolve("lento.txt");
CountDownLatch iniciada = new CountDownLatch(1), interrumpida = new CountDownLatch(1);
Callable<ResultadoArchivo> tarea = () -> {
iniciada.countDown();
try {
new CountDownLatch(1).await();
return ResultadoArchivo.ok(lento, 0, 0);
} catch (InterruptedException e) {
interrumpida.countDown();
throw e;
}
};
ExecutorService coordinador = Executors.newSingleThreadExecutor();
ProcesadorLote p = new ProcesadorLote(1);
try {
Future<List<ResultadoArchivo>> ejecucion = coordinador.submit(
() -> p.procesarTareas(Map.of(lento, tarea), Duration.ofMillis(100))
);
assertTrue(iniciada.await(1, TimeUnit.SECONDS));
ResultadoArchivo r = ejecucion.get(2, TimeUnit.SECONDS).get(0);
assertAll(() -> assertFalse(r.isCorrecto()), () -> assertEquals("timeout", r.getError()));
assertTrue(interrumpida.await(1, TimeUnit.SECONDS));
} finally {
p.close();
coordinador.shutdownNow();
assertTrue(coordinador.awaitTermination(2, TimeUnit.SECONDS));
}
assertTrue(p.estaTerminado());
}
}

Los seis tests no dependen de sleeps, orden del scheduler, CPU ni archivos personales. También comprueban un fallo parcial real, timeout, cancelación cooperativa y terminación explícita del executor.

README.md​

# Procesador de textos

Requisitos: JDK 17 y Maven.

1. Crea una carpeta `datos` en la raíz del proyecto.
2. Dentro, crea uno o dos archivos `.txt`, por ejemplo `ejemplo.txt`, con cualquier texto.
3. Ejecuta las pruebas: `mvn clean verify`.
4. Ejecuta la aplicación: `mvn compile exec:java -Dexec.args="datos"`.

Cada `.txt` produce nombre, estado, número de líneas y palabras. El orden mostrado es alfabético; no representa el orden de finalización concurrente.

Validación razonada​

Garantizado: pool acotado, resultado por archivo, fallos parciales conservados, timeout explícito y cierre en close().

Posible: cualquier orden de ejecución o finalización compatible con el pool.

Observado: debe registrarse al ejecutar tests y aplicación, sin convertir tiempos u orden en requisitos.

Errores habituales​

  • Compartir una lista mutable entre workers sin necesidad.
  • Crear un hilo por archivo.
  • Confundir timeout con terminación.
  • Tragar interrupciones.
  • Cancelar todo el lote por un fallo parcial no requerido.
  • Olvidar cerrar el stream de Files.list o el executor.