La creación y destrucción de hilos son operaciones costosas en términos de recursos del sistema. Además, el cambio de contexto entre hilos requiere una frecuente intervención de la CPU, lo que subraya la ineficiencia de crear un nuevo hilo para cada tarea. Para optimizar el uso de recursos y mejorar el rendimiento de aplicaciones concurrentes, los pools de hilos (o "piscinas de hilos") emergen como una solución fundamental. Estos administran un conjunto de hilos de trabajo que pueden reutilizarse para ejecutra múltiples tareas, eliminando la sobrecarga de crear hilos repetidamente.
Jerarquía de Clases Clave
El ecosistema de hilos en Java se organiza a través de una jerarquía de interfaces y clases, facilitando la gestión de la concurrencia:
Executor: La interfaz más básica, que define un único métodoexecute(Runnable command)para ejecutar tareas. Es la abstracción fundamental para enviar tareas.ExecutorService: ExtiendeExecutory añade funcionalidades para la gestión de ciclos de vida de hilos, incluyendo métodos para apagar el pool, y métodos para enviar tareas que devuelven resultados (submit) o colecciones de tareas (invokeAll,invokeAny).AbstractExecutorService: Una clase abstracta que proporciona implementaciones predeterminadas para varios métodos deExecutorService, facilitando la creación de implementaciones concretas de pools de hilos.ScheduledExecutorService: ExtiendeExecutorServicey permite la ejecución de tareas de forma programada o periódica.
Tipos Comunes de Thread Pools
La clase Executors ofrece métodos de fábrica para crear fácilmente diferentes tipos de pools de hilos con configuraciones predefinidas:
FixedThreadPool
Este pool mantiene un número fijo de hilos en todo momento. Su tamaño de núcleo (corePoolSize) es igual a su tamaño máximo (maximumPoolSize), lo que significa que los hilos no se destruyen por inactividad. Si se envían tareas cuando todos los hilos están ocupados, las nuevas tareas se colocan en una cola de espera. Los hilos inactivos de este pool recogerán tareas de la cola.
public static ExecutorService newFixedThreadPool(int cantidadHilos) {
return new ThreadPoolExecutor(cantidadHilos, cantidadHilos,
0L, TimeUnit.MILLISECONDS,
new LinkedBlockingQueue<Runnable>());
}
CachedThreadPool
Un pool flexible que crea nuevos hilos según sea necesario para ejecutar tareas, pero reutiliza hilos existentes cuando están disponibles. No tiene hilos de núcleo fijos (corePoolSize = 0) y puede crecer hasta un número ilimitado de hilos (maximumPoolSize = Integer.MAX_VALUE). Los hilos inactivos que no han procesado una tarea durante 60 segundos son terminados.
public static ExecutorService newCachedThreadPool() {
return new ThreadPoolExecutor(0, Integer.MAX_VALUE, 60L, TimeUnit.SECONDS,
new SynchronousQueue<Runnable>());
}
SingleThreadPool
Este pool garantiza que todas las tareas se ejecuten de forma secuencial por un único hilo de trabajo. Las tareas se procesan en el orden en que se envían, siguiendo una lógica de cola FIFO.
public static ExecutorService newSingleThreadExecutor() {
return new FinalizableDelegatedExecutorService
(new ThreadPoolExecutor(1, 1, 0L, TimeUnit.MILLISECONDS,
new LinkedBlockingQueue<Runnable>()));
}
ThreadPoolExecutor
La clase ThreadPoolExecutor es la implementación más versátil y configurable de ExecutorService en Java. Ofrece varios constructores que permiten personalizar el comportamiento del pool de hilos mediante la especificación de sus parámetros internos.
Parámetros Esenciales del ThreadPoolExecutor
Para configurar un ThreadPoolExecutor, es fundamental comprender sus parámetros constructores:
corePoolSize
Representa el número mínimo de hilos que siempre se mantendrán activos en el pool, incluso si están inactivos. Cuando se envían nuevas tareas y el número actual de hilos es inferior a corePoolSize, el pool creará nuevos hilos de núcleo para ejecutarlas. Por defecto, los hilos de núcleo no se eliminan por inactividad, a menos que se configure allowCoreThreadTimeOut a true.
maximumPoolSize
Define el número máximo total de hilos (núcleo + adicionales) que el pool puede llegar a tener. Este límite permite al pool manejar picos de carga creando hilos adicionales cuando la cola de tareas está llena y el número de hilos es inferior a maximumPoolSize. Una vez que la carga disminuye, los hilos adicionales (no de núcleo) pueden ser terminados.
keepAliveTime
Es el tiempo que un hilo adicional (no de núcleo) puede permanecer inactivo sin tareas antes de ser terminado. Si allowCoreThreadTimeOut se establece en true, esta regla también se aplica a los hilos de núcleo.
TimeUnit unit
La unidad de tiempo para el parámetro keepAliveTime. Es una enumeración que incluye valores como NANOSECONDS, MICROSECONDS, MILLISECONDS, SECONDS, MINUTES, HOURS, DAYS.
BlockingQueue<Runnable> workQueue
La cola de tareas que almacena los objetos Runnable pendientes de ejecución. Cuando todos los hilos de núcleo están ocupados, las nuevas tareas se colocan en esta cola. Si la cola se llena, el pool intentará crear hilos adicionales (no de núcleo) hasta maximumPoolSize para manejar las tareas.
ThreadFactory threadFactory
Una interfaz funcional para crear nuevos hilos. Permite personalizar los hilos creados por el pool, por ejemplo, asignándoles nombres específicos, prioridad o configurando si son hilos demonio. Requiere la implementación del método Thread newThread(Runnable r).
RejectedExecutionHandler handler
La estrategia de rechazo que el pool aplicará cuando no pueda aceptar una nueva tarea. Esto ocurre si el pool está apagado, o si ha alcanzado su maximumPoolSize y su workQueue está llena. Existen varias políticas predefinidas:
AbortPolicy: (Política por defecto) Lanza unaRejectedExecutionException, deteniendo el procesamiento de la tarea. Útil cuando se desea una notificación inmediata de que una tarea no pudo ser aceptada.
int tamanoNucleo = 2;
int tamanoMaximo = 5;
long tiempoVida = 60L;
int capacidadCola = 10;
ThreadPoolExecutor executor = new ThreadPoolExecutor(
tamanoNucleo,
tamanoMaximo,
tiempoVida,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(capacidadCola),
new ThreadPoolExecutor.AbortPolicy() // Política por defecto
);
CallerRunsPolicy: La tarea rechazada es ejecutada por el hilo que intentó enviarla al pool. Esto transfiere la carga al invocador, útil para desacelerar la tasa de envío de tareas.
new ThreadPoolExecutor.CallerRunsPolicy();
DiscardPolicy: Descarta silenciosamente la tarea rechazada. No se lanza ninguna excepción ni se realiza ninguna acción. Adecuada cuando la pérdida de tareas es aceptable.
new ThreadPoolExecutor.DiscardPolicy();
DiscardOldestPolicy: Elimina la tarea más antigua de la cola de trabajo y luego intenta re-enviar la tarea actual. Útil para mantener las tareas más recientes.
new ThreadPoolExecutor.DiscardOldestPolicy();
- Política de Rechazo Personalizada: Implementando la interfaz
RejectedExecutionHandler, se puede definir un comportamiento específico para las tareas rechazadas, como registrar el evento o intentar reintentar la tarea.
public class ManejadorRechazoPersonalizado implements RejectedExecutionHandler {
@Override
public void rejectedExecution(Runnable r, ThreadPoolExecutor executor) {
// Lógica personalizada: log, reintento, etc.
System.out.println("Tarea denegada: " + r.toString());
}
}
// Uso de la política personalizada
new ThreadPoolExecutor(tamanoNucleo, tamanoMaximo, tiempoVida, TimeUnit.SECONDS,
new ArrayBlockingQueue<>(capacidadCola),
new ManejadorRechazoPersonalizado());
Mecanismo de Operación
Modelo Productor-Consumidor
El funcionamiento de un ThreadPoolExecutor se basa en el patrón Productor-Consumidor. Los "productores" son los hilos que envían tareas al pool, y los "consumidores" son los hilos de trabajo del pool que extraen y ejecutan esas tareas.
Inicialización del Pool
Al crear un ThreadPoolExecutor, inicialmente no hay hilos activos. Estos se crean a medida que se necesitan, o se pueden precargar:
ThreadPoolExecutor gestorTareas = new ThreadPoolExecutor(
5, 5, 0L, TimeUnit.SECONDS, new ArrayBlockingQueue<>(10));
Para crear los hilos de núcleo antes de que lleguen las tareas, se pueden usar los métodos prestartCoreThread() o prestartAllCoreThreads():
// Inicia un único hilo de núcleo
public boolean prestartCoreThread() {
return workerCountOf(control.get()) < corePoolSize &&
addWorker(null, true);
}
// Inicia todos los hilos de núcleo
public int prestartAllCoreThreads() {
int contador = 0;
while (addWorker(null, true))
++contador;
return contador;
}
Envío de Tareas al Pool
ThreadPoolExecutor ofrece dos métodos principales para enviar tareas:
void execute(Runnable command): Heredado de la interfazExecutor. Se usa para tareas que no devuelven un resultado y no se espera que lancen excepciones verificadas (checked exceptions).Future<T> submit(Callable<T> task)oFuture<?> submit(Runnable task)o<T> Future<T> submit(Runnable task, T result): Estos métodos devuelven un objetoFuture, que permite obtener el resultado de la tarea (si es unCallable), verificar su estado o cancelarla. También encapsulan y permiten manejar excepciones lanzadas durante la ejecución. Internamente, todos los métodossubmitdelegan la ejecución aexecute.
public Future<?> submit(Runnable tarea) {
if (tarea == null) throw new NullPointerException();
RunnableFuture<Void> tareaFutura = newTaskFor(tarea, null);
execute(tareaFutura);
return tareaFutura;
}
public <T> Future<T> submit(Callable<T> tarea) {
if (tarea == null) throw new NullPointerException();
RunnableFuture<T> tareaFutura = newTaskFor(tarea);
execute(tareaFutura);
return tareaFutura;
}
public <T> Future<T> submit(Runnable tarea, T resultado) {
if (tarea == null) throw new NullPointerException();
RunnableFuture<T> tareaFutura = newTaskFor(tarea, resultado);
execute(tareaFutura);
return tareaFutura;
}
Diferencias Clave entre submit y execute:
- Gestión de Excepciones: Las excepciones lanzadas por tareas enviadas con
execute()no son directamente capturables por el código invocador, a menos que se configure unUncaughtExceptionHandler. Consubmit(), las excepciones se encapsulan en unExecutionExceptiony se lanzan al llamar aFuture.get(). - Cancelación de Tareas:
submit()devuelve un objetoFuture, lo que permite intentar cancelar la tarea usandoFuture.cancel(). Conexecute(), no hay un objetoFuturedirectamente accesible para gestionar la cancelación.
Flujo de Envío de Tareas
Cuando una tarea se envía mediante execute(), el pool sigue esta lógica:
- Hilos de Núcleo Disponibles: Si el número actual de hilos activos es menor que
corePoolSize, se crea un nuevo hilo de núcleo para ejecutar la tarea inmediatamente. Este hilo, al terminar su tarea, buscará más tareas en la cola. - Cola de Tareas: Si el número de hilos ya es igual o mayor que
corePoolSize, la tarea se intenta añadir a laworkQueue. Si la adición es exitosa, un hilo de núcleo inactivo (si lo hay) la tomará de la cola para ejecutarla. - Hilos Adicionales: Si la
workQueueestá llena, el pool verifica si el número actual de hilos es menor quemaximumPoolSize. Si es así, se crea un nuevo hilo no de núcleo para ejecutar la tarea de forma directa. Este nuevo hilo prioriza la ejecución de la tarea que causó su creación antes de buscar en la cola. - Política de Rechazo: Si el pool ha alcanzado
maximumPoolSizey laworkQueuetambién está llena, la tarea se considera "rechazada" y se aplica laRejectedExecutionHandlerconfigurada.
Implementación de execute:
public void execute(Runnable comando) {
if (comando == null)
throw new NullPointerException();
int estadoControl = ctl.get();
// ¿El conteo de trabajadores es menor que el tamaño del núcleo?
if (workerCountOf(estadoControl) < corePoolSize) {
// Añadir un hilo de núcleo para la tarea
if (addWorker(comando, true))
return;
estadoControl = ctl.get();
}
// Si el pool está en ejecución y la cola de trabajo acepta la tarea
if (isRunning(estadoControl) && workQueue.offer(comando)) {
int recheck = ctl.get();
if (! isRunning(recheck) && remove(comando))
// Ejecutar la política de rechazo si el pool se ha detenido
reject(comando);
else if (workerCountOf(recheck) == 0)
// Si no hay trabajadores, añadir uno (no de núcleo si la tarea es null)
addWorker(null, false);
}
// Si no se pudo añadir a la cola y no se puede añadir un trabajador
else if (!addWorker(comando, false))
// Ejecutar la política de rechazo
reject(comando);
}
Implementación de addWorker:
El método addWorker intenta crear y arrancar un nuevo hilo de trabajo (Worker).
private boolean addWorker(Runnable primeraTarea, boolean esNucleo) {
// Bucle para manejar cambios en el estado del pool de hilos
retry:
for (;;) {
int estadoControl = ctl.get();
int estadoEjecucion = runStateOf(estadoControl);
// Comprobar si la cola está vacía solo si es necesario para apagar
if (estadoEjecucion >= SHUTDOWN &&
! (estadoEjecucion == SHUTDOWN &&
primeraTarea == null &&
! workQueue.isEmpty()))
return false;
for (;;) { // Bucle interno para ajustar el conteo de trabajadores
int contadorTrabajadores = workerCountOf(estadoControl);
if (contadorTrabajadores >= CAPACITY ||
contadorTrabajadores >= (esNucleo ? corePoolSize : maximumPoolSize))
return false; // Se ha alcanzado el límite
if (compareAndIncrementWorkerCount(estadoControl))
break retry; // Se incrementó el conteo, salir del bucle principal
estadoControl = ctl.get(); // Volver a leer ctl si CAS falló
if (runStateOf(estadoControl) != estadoEjecucion)
continue retry; // El estado de ejecución cambió, reintentar desde el principio
}
}
boolean trabajadorArrancado = false;
boolean trabajadorAnadido = false;
Worker nuevoTrabajador = null;
try {
nuevoTrabajador = new Worker(primeraTarea); // Crear el Worker
final Thread hiloAsociado = nuevoTrabajador.thread;
if (hiloAsociado != null) {
final ReentrantLock bloqueoPrincipal = this.mainLock;
bloqueoPrincipal.lock(); // Bloquear para modificar la colección de trabajadores
try {
// Recomprobar el estado mientras se mantiene el bloqueo
int estadoEjecucion = runStateOf(ctl.get());
if (estadoEjecucion < SHUTDOWN ||
(estadoEjecucion == SHUTDOWN && primeraTarea == null)) {
if (hiloAsociado.isAlive())
throw new IllegalThreadStateException();
workers.add(nuevoTrabajador);
int tamanoActual = workers.size();
if (tamanoActual > largestPoolSize)
largestPoolSize = tamanoActual;
trabajadorAnadido = true;
}
} finally {
bloqueoPrincipal.unlock();
}
if (trabajadorAnadido) {
// Iniciar el hilo del trabajador
hiloAsociado.start();
trabajadorArrancado = true;
}
}
} finally {
if (! trabajadorArrancado)
addWorkerFailed(nuevoTrabajador);
}
return trabajadorArrancado;
}
La Clase Worker
Dentro de ThreadPoolExecutor, cada hilo de trabajo está representado por una instancia de la clase interna Worker. Esta clase encapsula tanto el Thread subyacente como la primera Runnable que debe ejecutar. Además, Worker extiende AbstractQueuedSynchronizer (AQS) para gestionar de forma eficiente el estado de interrupción y exclusión mutua de los hilos de trabajo.
private final HashSet<Worker> hilosTrabajadores = new HashSet<Worker>();
private final class Worker
extends AbstractQueuedSynchronizer
implements Runnable {
/** El hilo en el que se ejecuta este Worker. Nulo si la fábrica falla. */
final Thread hiloInterno;
/** La primera tarea a ejecutar. Puede ser nula. */
Runnable primeraTarea;
/** Contador de tareas completadas por este hilo. */
volatile long tareasCompletadas;
Worker(Runnable tareaInicial) {
setState(-1); // Inhibir interrupciones hasta que se ejecute runWorker
this.primeraTarea = tareaInicial;
this.hiloInterno = getThreadFactory().newThread(this);
}
}
Cuando se crea un Worker, su constructor llama a la ThreadFactory para obtener el Thread real, pasándose a sí mismo (el Worker) como el Runnable para ese hilo. Así, cuando el hilo se inicia, ejecutará el método run() del Worker, que a su vez invoca runWorker().
Ejecución de Tareas: runWorker
El método runWorker contiene el bucle principal de vida de un hilo de trabajo. Se encarga de ejecutar la primera tarea asignada (si existe) y luego de forma continua, busca nuevas tareas en la cola de trabajo para ejecutarlas.
final void runWorker(Worker trabajador) {
Thread hiloActual = Thread.currentThread();
Runnable tarea = trabajador.primeraTarea;
trabajador.primeraTarea = null;
trabajador.unlock(); // Permitir interrupciones; el estado -1 se convierte en 0
boolean completadoAbruptamente = true;
try {
// Bucle para obtener y ejecutar tareas
while (tarea != null || (tarea = obtenerTarea()) != null) {
trabajador.lock(); // Bloquear para prevenir interrupciones externas durante la ejecución de la tarea
try {
// Comprobar el estado del pool y la interrupción del hilo
if ((runStateAtLeast(ctl.get(), STOP) ||
(Thread.interrupted() &&
runStateAtLeast(ctl.get(), STOP))) &&
!hiloActual.isInterrupted())
hiloActual.interrupt();
antesDeEjecutar(hiloActual, tarea); // Hook para pre-ejecución
Throwable excepcionLanzada = null;
try {
tarea.run(); // Ejecutar la tarea
} catch (RuntimeException x) {
excepcionLanzada = x; throw x;
} catch (Error x) {
excepcionLanzada = x; throw x;
} catch (Throwable x) {
excepcionLanzada = x; throw new Error(x);
} finally {
despuesDeEjecutar(tarea, excepcionLanzada); // Hook para post-ejecución
}
} finally {
tarea = null;
trabajador.tareasCompletadas++;
trabajador.unlock(); // Liberar el bloqueo
}
}
completadoAbruptamente = false;
} finally {
// Procesa la salida del trabajador, eliminándolo del conjunto de workers
procesarSalidaTrabajador(trabajador, completadoAbruptamente);
}
}
El uso de trabajador.lock() y trabajador.unlock() alrededor de la ejecución de la tarea asegura que un hilo de trabajo no pueda ser interrumpido desde fuera mientras está activamente procesando una tarea. Solo un trabajador inactivo (que no ha llamado a lock()) puede ser interrumpido para ser terminado.
Obtención de Tareas: getTask
El método getTask es responsable de extraer tareas de la workQueue. Implementa la lógica de tiempo de vida (keepAliveTime) para hilos no de núcleo o hilos de núcleo configurados para expirar.
private Runnable obtenerTarea() {
boolean tiempoAgotado = false; // ¿Se agotó el tiempo en la última encuesta?
for (;;) {
int estadoControl = ctl.get();
int estadoEjecucion = runStateOf(estadoControl);
// Comprobar si la cola está vacía solo si es necesario para apagar/detener
if (estadoEjecucion >= SHUTDOWN && (estadoEjecucion >= STOP || workQueue.isEmpty())) {
decrementWorkerCount();
return null; // El trabajador debe terminar
}
int contadorTrabajadores = workerCountOf(estadoControl);
// ¿Están los trabajadores sujetos a eliminación por inactividad?
// Esto es cierto si se permite el timeout para hilos de núcleo o si es un hilo no de núcleo
boolean tieneTimeout = allowCoreThreadTimeOut || contadorTrabajadores > corePoolSize;
// Si el conteo de trabajadores excede el máximo o un hilo con timeout expiró,
// Y hay más de un trabajador o la cola está vacía, intentar decrementar y terminar.
if ((contadorTrabajadores > maximumPoolSize || (tieneTimeout && tiempoAgotado))
&& (contadorTrabajadores > 1 || workQueue.isEmpty())) {
if (compareAndDecrementWorkerCount(estadoControl))
return null; // Este trabajador termina
continue;
}
try {
// Esperar por una tarea: con timeout si 'tieneTimeout' es true, o bloqueando indefinidamente
Runnable r = tieneTimeout ?
workQueue.poll(keepAliveTime, TimeUnit.NANOSECONDS) :
workQueue.take();
if (r != null)
return r; // Tarea obtenida
tiempoAgotado = true; // Se agotó el tiempo de espera
} catch (InterruptedException reintento) {
tiempoAgotado = false; // Interrupción, resetear el flag de timeout
}
}
}
Terminación de Hilos de Trabajo
Cuando un hilo de trabajo ya no puede obtener tareas de la cola (por ejemplo, porque la cola está vacía y el tiempo de vida ha expirado), se invoca processWorkerExit() para gestionar su terminación y eliminación del pool.
private void procesarSalidaTrabajador(Worker trabajador, boolean completadoAbruptamente) {
if (completadoAbruptamente) // Si la salida fue abrupta, el contador de trabajadores no se ajustó
decrementWorkerCount();
final ReentrantLock bloqueoPrincipal = this.mainLock;
bloqueoPrincipal.lock();
try {
contadorTareasCompletadas += trabajador.tareasCompletadas;
hilosTrabajadores.remove(trabajador); // Eliminar el Worker del conjunto
} finally {
bloqueoPrincipal.unlock();
}
intentarTerminar(); // Intentar terminar el pool si todas las condiciones se cumplen
int estadoControl = ctl.get();
if (runStateLessThan(estadoControl, STOP)) {
if (!completadoAbruptamente) {
int minHilos = allowCoreThreadTimeOut ? 0 : corePoolSize;
if (minHilos == 0 && ! workQueue.isEmpty())
minHilos = 1; // Si no hay hilos de núcleo y la cola no está vacía, mantener al menos uno
if (workerCountOf(estadoControl) >= minHilos)
return; // No se necesita reemplazo
}
addWorker(null, false); // Añadir un nuevo trabajador si es necesario
}
}
Cierre del Thread Pool
Existen dos métodos principales para apagar un ThreadPoolExecutor:
shutdown(): Inicia un apagado ordenado. El pool deja de aceptar nuevas tareas, pero continúa ejecutando las tareas ya enviadas (las que están en ejecución y las que esperan en la cola). El pool se termina completamente una vez que todas las tareas han finalizado.shutdownNow(): Intenta detener el pool de forma inmediata. Deja de aceptar nuevas tareas y trata de interrumpir los hilos en ejecución. Devuelve una lista de las tareas que estaban en la cola y que no legaron a ejecutarse. La interrupción de tareas en curso depende de cómo la tarea maneje las interrupciones.
Implementación de shutdown y shutdownNow
public void shutdown() {
final ReentrantLock bloqueoPrincipal = this.mainLock;
bloqueoPrincipal.lock();
try {
// Verificar permisos para cerrar el pool
checkShutdownAccess();
// Cambiar el estado del pool a SHUTDOWN
advanceRunState(SHUTDOWN);
// Interrumpir hilos que están esperando tareas (inactivos)
interrumpirTrabajadoresInactivos();
onShutdown(); // Hook para ScheduledThreadPoolExecutor
} finally {
bloqueoPrincipal.unlock();
}
intentarTerminar(); // Intentar terminar el pool
}
public List<Runnable> shutdownNow() {
List<Runnable> tareasPendientes;
final ReentrantLock bloqueoPrincipal = this.mainLock;
bloqueoPrincipal.lock();
try {
checkShutdownAccess();
// Cambiar el estado del pool a STOP
advanceRunState(STOP);
// Interrumpir todos los hilos de trabajo, estén inactivos o no
interrumpirTodosLosTrabajadores();
// Vaciar la cola de tareas y obtener las tareas pendientes
tareasPendientes = vaciarCola();
} finally {
bloqueoPrincipal.unlock();
}
intentarTerminar(); // Intentar terminar el pool
return tareasPendientes;
}
Diferencias en Interrupción: interrumpirTrabajadoresInactivos vs. interrumpirTodosLosTrabajadores
shutdown() llama a interrumpirTrabajadoresInactivos(), que solo interrumpe los hilos que no están ejecutando una tarea activamente (es decir, están esperando una tarea de la cola). shutdownNow(), por otro lado, llama a interrumpirTodosLosTrabajadores(), que intenta interrumpir a todos los hilos, sin importar si están ocupados o inactivos.
private void interrumpirTrabajadoresInactivos() {
interrumpirTrabajadoresInactivos(false);
}
// Interrumpe hilos que probablemente están esperando tareas
private void interrumpirTrabajadoresInactivos(boolean soloUno) {
final ReentrantLock bloqueoPrincipal = this.mainLock;
bloqueoPrincipal.lock(); // Bloqueo externo para asegurar la estabilidad del conjunto de workers
try {
for (Worker w : hilosTrabajadores) {
Thread t = w.hiloInterno;
// Solo interrumpir si el hilo no está ya interrumpido y se puede adquirir el bloqueo del Worker
if (!t.isInterrupted() && w.tryLock()) {
try {
t.interrupt(); // Interrumpir el hilo
} catch (SecurityException ignore) {
} finally {
w.unlock();
}
}
if (soloUno)
break;
}
} finally {
bloqueoPrincipal.unlock();
}
}
// Llamado por shutdownNow: interrumpe a todos los trabajadores
private void interrumpirTodosLosTrabajadores() {
final ReentrantLock bloqueoPrincipal = this.mainLock;
bloqueoPrincipal.lock(); // Bloqueo externo para asegurar la estabilidad del conjunto de workers
try {
for (Worker w : hilosTrabajadores)
w.interrumpirSiIniciado(); // Método interno del Worker para interrumpir
} finally {
bloqueoPrincipal.unlock();
}
}
La razón por la que interruptIdleWorkers usa un w.tryLock() es para determinar si el Worker está o no ejecutando una tarea. Si tryLock() tiene éxito, significa que el Worker no ha adquirido su propio bloqueo interno (o lo ha liberado), indicando que está inactivo y puede ser interrumpido. Si falla, el Worker está ocupado con una tarea y, por lo tanto, no se le intrerumpe de inmediato. En su lugar, el hilo de trabajo terminará su tarea actual y, al intentar obtener la siguiente tarea (mediante getTask()), detectará el estado de apagado del pool y finalizará de forma natural.
AQS en Worker
La extensión de AbstractQueuedSynchronizer por la clase Worker es fundamental para gestionar el estado de los hilos y controlar las interrupciones. Worker implementa un bloqueo de exclusión mutua no reentrante. Esto significa que un hilo de trabajo solo puede adquirir el bloqueo una vez. Si intentara adquirirlo de nuevo mientras ya lo posee (por ejemplo, si una tarea invocara métodos de control del pool que también intentaran adquirir el bloqueo), fallaría. Esto previene interrupciones accidentales de tareas en ejecución. El estado inicial del AQS en Worker es -1, lo que inhibe las interrupciones hasta que el hilo de trabajo realmente comienza a ejecutar tareas, momento en el que el estado cambia a 0 (mediante unlock() en runWorker).
protected boolean estaBloqueadoExclusivamente() {
return getState() != 0;
}
protected boolean intentarAdquirir(int noUsado) {
if (compareAndSetState(0, 1)) { // Cambiar estado de 0 a 1
setExclusiveOwnerThread(Thread.currentThread());
return true;
}
return false;
}
protected boolean intentarLiberar(int noUsado) {
setExclusiveOwnerThread(null);
setState(0); // Restablecer estado a 0
return true;
}
public void lock() { acquire(1); }
public boolean tryLock() { return intentarAdquirir(1); }
public void unlock() { release(1); }
public boolean isLocked() { return estaBloqueadoExclusivamente(); }
void interrumpirSiIniciado() {
Thread t;
// Solo interrumpir si el estado no es -1 (es decir, ya se ha iniciado)
if (getState() >= 0 && (t = hiloInterno) != null && !t.isInterrupted()) {
try {
t.interrupt();
} catch (SecurityException ignore) {
}
}
}