Saltar al contenido

Cómo construir una cola de tareas desde cero con Node.js y Redis

Construye paso a paso una cola de tareas para producción con Node.js y Redis: entrega fiable, reintentos, colas muertas y control de concurrencia.

6 min de lectura
Diagrama de arquitectura que muestra una cola de tareas con productores, un bróker Redis y workers consumidores

Toda aplicación termina necesitando procesamiento en segundo plano. Notificaciones por correo, redimensionado de imágenes, generación de informes, entrega de webhooks: estas tareas no pertenecen al ciclo de petición-respuesta. Una cola de tareas desacopla a los productores de los consumidores y permite un procesamiento asíncrono fiable.

Aunque existen bibliotecas como BullMQ, construir una cola desde cero te enseña las primitivas que hacen funcionar el procesamiento distribuido de trabajos. Entenderás exactamente qué ocurre cuando un trabajo falla, cómo se gestiona la concurrencia y por qué ciertas decisiones de diseño importan.

La estructura básica de la cola

Una cola de tareas necesita tres cosas: una forma de encolar trabajos, una forma de desencolarlos de manera fiable y una forma de seguir su estado. Redis proporciona todos los componentes.

tstypescript
import { createClient, RedisClientType } from "redis";
 
interface Job<T = unknown> {
  id: string;
  queue: string;
  payload: T;
  attempts: number;
  maxAttempts: number;
  createdAt: number;
  processAfter: number;
}
 
class TaskQueue {
  private client: RedisClientType;
  private prefix: string;
 
  constructor(client: RedisClientType, prefix: string = "tq") {
    this.client = client;
    this.prefix = prefix;
  }
 
  private key(queue: string, suffix: string): string {
    return `${this.prefix}:${queue}:${suffix}`;
  }
 
  async enqueue<T>(
    queue: string,
    payload: T,
    options: { delay?: number; maxAttempts?: number } = {}
  ): Promise<string> {
    const id = crypto.randomUUID();
    const now = Date.now();
 
    const job: Job<T> = {
      id,
      queue,
      payload,
      attempts: 0,
      maxAttempts: options.maxAttempts ?? 3,
      createdAt: now,
      processAfter: now + (options.delay ?? 0),
    };
 
    const multi = this.client.multi();
 
    // Store job data
    multi.set(
      this.key(queue, `job:${id}`),
      JSON.stringify(job)
    );
 
    // Add to waiting sorted set (scored by processAfter)
    multi.zAdd(this.key(queue, "waiting"), {
      score: job.processAfter,
      value: id,
    });
 
    await multi.exec();
    return id;
  }
}

Usar un sorted set para la cola de espera habilita los trabajos diferidos de forma natural: los trabajos se puntúan según su marca de tiempo processAfter, así que solo tomamos los trabajos cuya puntuación es igual o anterior al momento actual.

Consumo fiable de trabajos

El desafío crítico en cualquier cola es garantizar el procesamiento exactamente una vez. Si un worker se cae a mitad de un trabajo, el trabajo no debe perderse. El ZPOPMIN de Redis combinado con un conjunto de procesamiento proporciona esta garantía.

tstypescript
// ❌ Unreliable: job lost if worker crashes after pop
async function unsafeDequeue(client: RedisClientType, queue: string) {
  const result = await client.zPopMin(queue);
  // If process crashes here, job is gone forever
  return result;
}
tstypescript
// ✅ Reliable: job tracked in processing set
class TaskQueue {
  // ... previous code
 
  async dequeue(queue: string): Promise<Job | null> {
    const now = Date.now();
    const waitingKey = this.key(queue, "waiting");
    const processingKey = this.key(queue, "processing");
 
    // Atomically move job from waiting to processing
    const script = `
      local result = redis.call('ZRANGEBYSCORE', KEYS[1], '-inf', ARGV[1], 'LIMIT', 0, 1)
      if #result == 0 then return nil end
      local jobId = result[1]
      redis.call('ZREM', KEYS[1], jobId)
      redis.call('ZADD', KEYS[2], ARGV[2], jobId)
      return jobId
    `;
 
    const jobId = await this.client.eval(script, {
      keys: [waitingKey, processingKey],
      arguments: [String(now), String(now)],
    });
 
    if (!jobId) return null;
 
    const jobData = await this.client.get(
      this.key(queue, `job:${jobId}`)
    );
 
    if (!jobData) return null;
    return JSON.parse(jobData) as Job;
  }
 
  async acknowledge(queue: string, jobId: string): Promise<void> {
    const multi = this.client.multi();
 
    multi.zRem(this.key(queue, "processing"), jobId);
    multi.del(this.key(queue, `job:${jobId}`));
    multi.incr(this.key(queue, "stats:completed"));
 
    await multi.exec();
  }
}

El script Lua se ejecuta de forma atómica en Redis: ningún otro cliente puede interferir entre la comprobación del conjunto de espera y el traslado del trabajo al conjunto de procesamiento. Esto elimina la condición de carrera en la que dos workers podrían tomar el mismo trabajo.

Lógica de reintentos y colas de mensajes muertos

Los trabajos fallan. Las redes agotan el tiempo de espera, las APIs externas devuelven errores, los datos llegan mal formados. Una cola robusta reintenta los fallos transitorios y envía los fallos persistentes a una dead letter queue para su investigación.

tstypescript
class TaskQueue {
  // ... previous code
 
  async fail(
    queue: string,
    jobId: string,
    error: string
  ): Promise<"retried" | "dead-lettered"> {
    const jobData = await this.client.get(
      this.key(queue, `job:${jobId}`)
    );
 
    if (!jobData) throw new Error(`Job ${jobId} not found`);
 
    const job: Job = JSON.parse(jobData);
    job.attempts += 1;
 
    // Remove from processing set
    await this.client.zRem(this.key(queue, "processing"), jobId);
 
    if (job.attempts < job.maxAttempts) {
      // Exponential backoff: 1s, 4s, 9s, 16s...
      const delay = Math.pow(job.attempts, 2) * 1000;
      job.processAfter = Date.now() + delay;
 
      const multi = this.client.multi();
      multi.set(
        this.key(queue, `job:${jobId}`),
        JSON.stringify(job)
      );
      multi.zAdd(this.key(queue, "waiting"), {
        score: job.processAfter,
        value: jobId,
      });
      await multi.exec();
 
      return "retried";
    }
 
    // Max attempts exceeded: dead letter queue
    const multi = this.client.multi();
    multi.lPush(
      this.key(queue, "dead"),
      JSON.stringify({ ...job, error, failedAt: Date.now() })
    );
    multi.del(this.key(queue, `job:${jobId}`));
    multi.incr(this.key(queue, "stats:dead-lettered"));
    await multi.exec();
 
    return "dead-lettered";
  }
}

El backoff exponencial con la fórmula attempts² × 1000ms produce esperas cada vez más largas: 1 segundo, 4 segundos, 9 segundos. Esto evita tormentas de reintentos cuando un servicio descendente tiene problemas y le da tiempo para recuperarse.

Construcción del bucle del worker

Un worker consulta continuamente si hay trabajos, los procesa y gestiona el éxito o el fallo. El intervalo de sondeo equilibra la capacidad de respuesta con la carga sobre Redis.

tstypescript
type JobHandler<T = unknown> = (payload: T) => Promise<void>;
 
class Worker {
  private queue: TaskQueue;
  private queueName: string;
  private handler: JobHandler;
  private running: boolean = false;
  private concurrency: number;
  private activeJobs: number = 0;
 
  constructor(
    queue: TaskQueue,
    queueName: string,
    handler: JobHandler,
    concurrency: number = 5
  ) {
    this.queue = queue;
    this.queueName = queueName;
    this.handler = handler;
    this.concurrency = concurrency;
  }
 
  async start(): Promise<void> {
    this.running = true;
    console.log(
      `Worker started for queue "${this.queueName}" ` +
      `(concurrency: ${this.concurrency})`
    );
 
    while (this.running) {
      if (this.activeJobs >= this.concurrency) {
        await this.sleep(100);
        continue;
      }
 
      const job = await this.queue.dequeue(this.queueName);
 
      if (!job) {
        await this.sleep(1000); // No jobs available, wait
        continue;
      }
 
      this.activeJobs++;
      this.processJob(job).finally(() => {
        this.activeJobs--;
      });
    }
  }
 
  private async processJob(job: Job): Promise<void> {
    try {
      await this.handler(job.payload);
      await this.queue.acknowledge(this.queueName, job.id);
    } catch (error) {
      const message =
        error instanceof Error ? error.message : "Unknown error";
      const result = await this.queue.fail(
        this.queueName,
        job.id,
        message
      );
      console.warn(
        `Job ${job.id} failed (${result}): ${message}`
      );
    }
  }
 
  stop(): void {
    this.running = false;
  }
 
  private sleep(ms: number): Promise<void> {
    return new Promise(resolve => setTimeout(resolve, ms));
  }
}

El control de concurrencia garantiza que el worker procese varios trabajos simultáneamente sin sobrecargar el sistema. Cada trabajo se ejecuta de forma independiente: un trabajo lento no bloquea a los demás.

Recuperación de trabajos estancados

Si un worker se cae, sus trabajos permanecen en el conjunto de procesamiento indefinidamente. Un proceso de recuperación independiente debe detectar los trabajos estancados y volver a encolarlos.

tstypescript
class TaskQueue {
  // ... previous code
 
  async recoverStalledJobs(
    queue: string,
    stallTimeout: number = 30000
  ): Promise<number> {
    const processingKey = this.key(queue, "processing");
    const cutoff = Date.now() - stallTimeout;
 
    // Find jobs that have been processing longer than stallTimeout
    const stalledIds = await this.client.zRangeByScore(
      processingKey,
      "-inf",
      String(cutoff)
    );
 
    let recovered = 0;
 
    for (const jobId of stalledIds) {
      const jobData = await this.client.get(
        this.key(queue, `job:${jobId}`)
      );
 
      if (!jobData) {
        // Job data missing, just clean up
        await this.client.zRem(processingKey, jobId);
        continue;
      }
 
      const job: Job = JSON.parse(jobData);
      job.attempts += 1;
 
      if (job.attempts >= job.maxAttempts) {
        // Exceeded retries, dead letter it
        const multi = this.client.multi();
        multi.zRem(processingKey, jobId);
        multi.lPush(
          this.key(queue, "dead"),
          JSON.stringify({
            ...job,
            error: "Stalled and exceeded max attempts",
            failedAt: Date.now(),
          })
        );
        multi.del(this.key(queue, `job:${jobId}`));
        await multi.exec();
      } else {
        // Re-enqueue for retry
        const multi = this.client.multi();
        multi.zRem(processingKey, jobId);
        multi.set(
          this.key(queue, `job:${jobId}`),
          JSON.stringify(job)
        );
        multi.zAdd(this.key(queue, "waiting"), {
          score: Date.now(),
          value: jobId,
        });
        await multi.exec();
        recovered++;
      }
    }
 
    return recovered;
  }
}

Ejecuta el barrido de recuperación con un temporizador: cada 30 segundos es razonable para la mayoría de las aplicaciones. El tiempo de espera de estancamiento debe ser mayor que la duración del trabajo más largo que esperes, para no recuperar trabajos que aún se están procesando activamente.

Uniendo todas las piezas

Así es como los productores y consumidores interactúan a través de la cola en un escenario de aplicación real.

tstypescript
async function main() {
  const redis = createClient({ url: "redis://localhost:6379" });
  await redis.connect();
 
  const taskQueue = new TaskQueue(redis);
 
  // Producer: enqueue email jobs
  await taskQueue.enqueue("emails", {
    to: "user@example.com",
    subject: "Welcome",
    template: "onboarding",
  });
 
  await taskQueue.enqueue(
    "emails",
    {
      to: "admin@example.com",
      subject: "Daily Report",
      template: "report",
    },
    { delay: 60000 } // Delay 1 minute
  );
 
  // Consumer: process email jobs
  const worker = new Worker(
    taskQueue,
    "emails",
    async (payload) => {
      const emailPayload = payload as {
        to: string;
        subject: string;
        template: string;
      };
      console.log(`Sending email to ${emailPayload.to}`);
      // await sendEmail(emailPayload);
    },
    3 // Process 3 emails concurrently
  );
 
  // Start stall recovery on interval
  setInterval(() => {
    taskQueue.recoverStalledJobs("emails").then(count => {
      if (count > 0) console.log(`Recovered ${count} stalled jobs`);
    });
  }, 30000);
 
  await worker.start();
}

Conclusiones clave

Construir una cola de tareas desde cero revela la complejidad que se esconde tras las APIs de procesamiento de trabajos de aspecto sencillo. Los desafíos centrales no consisten en encolar datos, sino en las garantías de fiabilidad. Las transiciones atómicas de trabajos entre estados evitan duplicados. El backoff exponencial evita tormentas de reintentos. Las dead letter queues evitan la pérdida silenciosa de datos. La recuperación de estancados evita que los trabajos desaparezcan cuando los workers se caen.

Entender estas primitivas te convierte en un mejor usuario de las bibliotecas de colas de producción y te da la base para depurar los problemas que inevitablemente surgen en los sistemas distribuidos de procesamiento de trabajos.

Wilfredo Rujel

Wilfredo Rujel

Ingeniero de Software Full Stack

Compartir esta publicaciónX