Construcción de un Orquestador de Colas Asíncronas en Node.js con Control de Concurrencia y Full Jitter Backoff

JAVASCRIPT / NODE.JS 27 de julio de 2026 100 lecturas
Implementa un orquestador asíncrono en memoria para Node.js con control de concurrencia estricto, resiliencia con Full Jitter y capacidades operativas avanzadas.

Durante años de diseño de sistemas distribuidos y arquitecturas orientadas a eventos en Node.js, he observado repetidamente un patrón nefasto en código de producción: la confianza ciega en la naturaleza no bloqueante de V8. Cuando los desarrolladores se enfrentan a la necesidad de procesar colecciones masivas de datos o coordinar llamadas asíncronas hacia servicios de terceros, la solución inmediata suele ser el uso indiscriminado de Promise.all() combinado con iteradores como .map(). Aunque este enfoque parece elegante en entornos de desarrollo local con cargas de prueba insignificantes, en entornos de producción con alta carga se convierte en una causa principal de fallos catastróficos e impredecibles.

El problema fundamental radica en que Promise.all() no ofrece ningún tipo de estrangulamiento o regulación de flujo (throttling). Si se le pasa un array de diez mil elementos, la función instanciará diez mil promesas en el mismo 'tick' del bucle de eventos. Esto provoca una ráfaga masiva de solicitudes I/O que agota instantáneamente el pool de sockets TCP del sistema operativo, infla el consumo de memoria Heap debido al almacenamiento de miles de contextos de ejecución diferidos y, finalmente, desata respuestas de error masivas por parte de las API externas, como HTTP 429 (Too Many Requests) o HTTP 503 (Service Unavailable). Peor aún, en sistemas con recursos acotados de CPU o RAM, la recolección de basura (Garbage Collection) puede paralizar el proceso completo tratando de limpiar miles de objetos abandonados.

Para solucionar la saturación del Event Loop sin introducir dependencias pesadas como Redis, RabbitMQ o infraestructura de colas distribuida, la solución arquitectónica idónea es un orquestador de colas asíncronas en memoria. La función primordial de esta abstracción es garantizar un límite rígido de concurrencia (Concurrency Limit), asegurando que en cualquier instante de tiempo dado no existan más de 'N' tareas ejecutándose simultáneamente. Sin embargo, limitar la concurrencia es únicamente la mitad de la ecuación cuando construimos sistemas verdaderamente resilientes.

El verdadero desafío surge cuando una tarea falla debido a congestión de red, agotamiento temporal de recursos o límites de tasa impuestos por un proveedor externo. La estrategia ingenua más común es reintentar la operación inmediatamente o tras un periodo fijo. Una mejora sobre esto es el backoff exponencial estándar, donde el tiempo de espera se duplica tras cada fallo consecutivo. No obstante, el backoff exponencial rígido sufre de un problema grave conocido como el 'Thundering Herd Problem' o estampida de peticiones. Si cincuenta peticiones fallan simultáneamente al golpear un límite de tasa en el segundo X, y todas aplican un retardo exacto de dos segundos, las cincuenta peticiones reintentarán su ejecución exactamente en el segundo X+2, provocando nuevamente un pico colosal de tráfico que volverá a derribar el servicio remoto.

Para romper esta sincronización destructiva, el equipo de arquitectura de AWS popularizó el concepto de 'Full Jitter'. En lugar de utilizar un valor determinista derivado del cálculo exponencial, el algoritmo calcula el techo máximo de espera para el intento actual y selecciona aleatoriamente un entero entre cero y dicho límite máximo. Esta aleatorización distribuye los reintentos de manera uniforme a lo largo del tiempo, aplanando las crestas de tráfico y transformando las ráfagas caóticas en un flujo continuo y manejable para la infraestructura receptora.

Otro aspecto crucial en la implementación de una cola de alto rendimiento en Node.js es la gestión eficiente de los slots de concurrencia durante los periodos de espera. Si un orquestador bloquea una ranura de ejecución activa mientras espera a que se cumpla el temporizador de reintento de una tarea fallida, está desaprovechando capacidad del sistema que podría ser utilizada por otras tareas pendientes. La arquitectura refinada debe implementar la 'liberación diferida fuera de slot': al fallar una tarea, la ranura de concurrencia activa se decrementa e inmediatamente se procesa el siguiente elemento encolado. El reintento de la tarea fallida se programa mediante un temporizador asíncrono aislado que reinsertará el elemento en la cola únicamente cuando el tiempo de jitter se haya cumplido.

Adicionalmente, en entornos de nivel empresarial, un orquestador en memoria debe contemplar salvaguardas operativas indispensables: un límite superior de tiempo de espera (maxDelay) para evitar que backoffs exponenciales alcancen retardos de días o semanas, controles para pausar y reanudar el flujo de trabajo dinámicamente, y mecanismos para purgar la cola en caso de emergencias arquitectónicas. A continuación, presento la implementación optimizada, modular y orientada a producción de este orquestador en memoria.

/**
 * Orquestador de Colas Asíncronas en Memoria para Node.js
 * Proporciona control de concurrencia, resiliencia con Full Jitter Backoff,
 * limite máximo de espera y controles operativos de pausa/cancelación.
 */
class AsyncQueue {
    /**
     * @param {Object} options Configuración del orquestador.
     * @param {number} [options.concurrency=3] Límite de tareas simultáneas.
     * @param {number} [options.maxRetries=3] Máximo de reintentos por tarea.
     * @param {number} [options.baseDelay=1000] Retardo base inicial en ms.
     * @param {number} [options.maxDelay=30000] Límite máximo de retardo en ms.
     */
    constructor({ concurrency = 3, maxRetries = 3, baseDelay = 1000, maxDelay = 30000 } = {}) {
        if (concurrency <= 0) throw new TypeError('Concurrency must be greater than 0.');
        if (maxRetries < 0) throw new TypeError('Max retries cannot be negative.');
        if (baseDelay < 0 || maxDelay < 0) throw new TypeError('Delays cannot be negative.');

        this.concurrency = concurrency;
        this.maxRetries = maxRetries;
        this.baseDelay = baseDelay;
        this.maxDelay = maxDelay;

        this.queue = [];
        this.activeCount = 0;
        this.isPaused = false;
        this._drainResolvers = [];
    }

    /**
     * Encola una función asíncrona y devuelve una Promesa con su resultado.
     * @template T
     * @param {() => Promise<T>} taskFn Función asíncrona a ejecutar.
     * @returns {Promise<T>}
     */
    enqueue(taskFn) {
        if (typeof taskFn !== 'function') {
            return Promise.reject(new TypeError('Task must be a function.'));
        }

        return new Promise((resolve, reject) => {
            this.queue.push({
                taskFn,
                resolve,
                reject,
                retries: 0
            });
            this._processNext();
        });
    }

    /**
     * Pausa el procesamiento de nuevas tareas en la cola.
     */
    pause() {
        this.isPaused = true;
    }

    /**
     * Reanuda el procesamiento de la cola.
     */
    resume() {
        if (this.isPaused) {
            this.isPaused = false;
            this._processNext();
        }
    }

    /**
     * Cancela y limpia todas las tareas pendientes en cola.
     * @param {Error} [reason] Motivo del rechazo de las tareas encoladas.
     */
    clear(reason = new Error('Queue cleared by operator.')) {
        while (this.queue.length > 0) {
            const task = this.queue.shift();
            task.reject(reason);
        }
        this._checkIdle();
    }

    /**
     * Retorna una promesa que se resuelve cuando la cola se vacía y no hay tareas activas.
     * @returns {Promise<void>}
     */
    onIdle() {
        if (this.queue.length === 0 && this.activeCount === 0) {
            return Promise.resolve();
        }
        return new Promise((resolve) => {
            this._drainResolvers.push(resolve);
        });
    }

    /**
     * Despacha tareas pendientes respetando el límite de concurrencia y el estado de pausa.
     * @private
     */
    _processNext() {
        if (this.isPaused) return;

        while (this.activeCount < this.concurrency && this.queue.length > 0) {
            const taskItem = this.queue.shift();
            this.activeCount++;
            this._executeTask(taskItem);
        }

        this._checkIdle();
    }

    /**
     * Verifica si la cola se ha completado y notifica a los oyentes de onIdle.
     * @private
     */
    _checkIdle() {
        if (this.queue.length === 0 && this.activeCount === 0 && this._drainResolvers.length > 0) {
            const resolvers = [...this._drainResolvers];
            this._drainResolvers = [];
            resolvers.forEach(resolve => resolve());
        }
    }

    /**
     * Ejecuta la tarea asíncrona gestionando reintentos y cálculo de Full Jitter.
     * @private
     * @param {Object} taskItem
     */
    async _executeTask(taskItem) {
        try {
            const result = await taskItem.taskFn();
            taskItem.resolve(result);
        } catch (error) {
            if (taskItem.retries < this.maxRetries) {
                taskItem.retries++;

                // Cálculo de Exponential Backoff acotado por maxDelay
                const calculatedMax = this.baseDelay * Math.pow(2, taskItem.retries - 1);
                const cappedMax = Math.min(this.maxDelay, calculatedMax);
                
                // Algoritmo Full Jitter: rand(0, min(maxDelay, base * 2^attempt))
                const jitterDelay = Math.floor(Math.random() * cappedMax);

                console.warn(
                    `[AsyncQueue] Error en tarea (Intento ${taskItem.retries}/${this.maxRetries}). ` +
                    `Reintentando en ${jitterDelay}ms. Detalle: ${error.message}`
                );

                // Liberación diferida: El slot actual se decrementa inmediatamente.
                // El reintento se reinserta en la cola tras expirar el temporizador.
                setTimeout(() => {
                    this.queue.push(taskItem);
                    this._processNext();
                }, jitterDelay);
            } else {
                console.error(
                    `[AsyncQueue] Tarea fallida tras ${this.maxRetries} reintentos: ${error.message}`
                );
                taskItem.reject(error);
            }
        } finally {
            this.activeCount--;
            this._processNext();
        }
    }

    /**
     * Retorna métricas del estado actual de la cola.
     */
    get stats() {
        return {
            pendingInQueue: this.queue.length,
            currentlyExecuting: this.activeCount,
            concurrencyLimit: this.concurrency,
            isPaused: this.isPaused,
            isIdle: this.queue.length === 0 && this.activeCount === 0
        };
    }
}

// ==========================================
// Ejemplo de Uso Práctico
// ==========================================
async function main() {
    const queue = new AsyncQueue({
        concurrency: 2,
        maxRetries: 3,
        baseDelay: 100,
        maxDelay: 2000
    });

    const simulateNetworkRequest = async (id) => {
        // Simulación de fallo aleatorio (60% de probabilidad)
        if (Math.random() < 0.6) {
            throw new Error(`HTTP 429 Rate Limit en item #${id}`);
        }
        return { id, data: `Resultado de payload #${id}`, timestamp: new Date().toISOString() };
    };

    const items = [101, 102, 103, 104, 105];

    console.log('=== Iniciando procesador de tareas ===');

    items.forEach(id => {
        queue.enqueue(() => simulateNetworkRequest(id))
            .then(res => console.log(`[ÉXITO] ID ${id}:`, res))
            .catch(err => console.error(`[FALLO FINAL] ID ${id}:`, err.message));
    });

    await queue.onIdle();
    console.log('=== Procesamiento finalizado ===');
    console.log('Métricas finales:', queue.stats);
}

main();
¿Qué te pareció?
🔥 Brillante 0
💡 Me sirvió 0
🚀 A otro nivel 0

¿Te resultó útil este snippet? Explora más código y soluciones en AndresSY.dev.

Volver a Snippets