Gestión de Hilos en Java: Profundizando en los `ThreadPoolExecutor`

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:

  1. Executor: La interfaz más básica, que define un único método execute(Runnable command) para ejecutar tareas. Es la abstracción fundamental para enviar tareas.
  2. ExecutorService: Extiende Executor y 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).
  3. AbstractExecutorService: Una clase abstracta que proporciona implementaciones predeterminadas para varios métodos de ExecutorService, facilitando la creación de implementaciones concretas de pools de hilos.
  4. ScheduledExecutorService: Extiende ExecutorService y 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:

  1. AbortPolicy: (Política por defecto) Lanza una RejectedExecutionException, 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
);

  1. 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();

  1. 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();

  1. 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();

  1. 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:

  1. void execute(Runnable command): Heredado de la interfaz Executor. Se usa para tareas que no devuelven un resultado y no se espera que lancen excepciones verificadas (checked exceptions).
  2. Future<T> submit(Callable<T> task) o Future<?> submit(Runnable task) o <T> Future<T> submit(Runnable task, T result): Estos métodos devuelven un objeto Future, que permite obtener el resultado de la tarea (si es un Callable), verificar su estado o cancelarla. También encapsulan y permiten manejar excepciones lanzadas durante la ejecución. Internamente, todos los métodos submit delegan la ejecución a execute.
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 un UncaughtExceptionHandler. Con submit(), las excepciones se encapsulan en un ExecutionException y se lanzan al llamar a Future.get().
  • Cancelación de Tareas: submit() devuelve un objeto Future, lo que permite intentar cancelar la tarea usando Future.cancel(). Con execute(), no hay un objeto Future directamente 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:

  1. 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.
  2. Cola de Tareas: Si el número de hilos ya es igual o mayor que corePoolSize, la tarea se intenta añadir a la workQueue. Si la adición es exitosa, un hilo de núcleo inactivo (si lo hay) la tomará de la cola para ejecutarla.
  3. Hilos Adicionales: Si la workQueue está llena, el pool verifica si el número actual de hilos es menor que maximumPoolSize. 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.
  4. Política de Rechazo: Si el pool ha alcanzado maximumPoolSize y la workQueue también está llena, la tarea se considera "rechazada" y se aplica la RejectedExecutionHandler configurada.

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:

  1. 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.
  2. 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) {
        }
    }
}

Etiquetas: java concurrencia hilos threadpool executor

Publicado el 7-21 19:45