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
OKoERROR; - usar
ExecutorService,CallableyFuture, 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.listo el executor.