¿Qué es una cola de trabajos?
Una cola de trabajos (job queue) es un búfer duradero que se sitúa entre el código que solicita una tarea y el código que la ejecuta. El productor escribe un registro pequeño que describe una unidad de trabajo y retorna inmediatamente. Un worker lee ese registro más tarde y se encarga de la parte lenta. La cola es lo que sobrevive a un reinicio, absorbe los picos de tráfico y proporciona un lugar donde gestionar los fallos.
Redis se ha utilizado de esta manera durante más de una década. Sus listas proporcionaron a las primeras librerías de colas primitivas atómicas de LPUSH y BRPOP, y los blocking pops permitieron que un worker permanezca inactivo hasta que llegue trabajo en lugar de realizar polling en un bucle. Más tarde, los Streams añadieron grupos de consumidores, acknowledgements y replay. Sobre estas primitivas, BullMQ envuelve todo el ciclo de vida —trabajos retrasados, reintentos, prioridades, rate limiting, programaciones repetibles y eventos— en una pequeña API de TypeScript.
Si Redis ya respalda tu caché o tus sesiones, añadir una cola es un paso sencillo. Esa conveniencia es la razón principal por la que es la opción predeterminada para los equipos de Node.js, y el motivo principal para comprender tanto lo que te ofrece como lo que no.
Una cola son tres roles, no uno
Todo sistema de colas, independientemente del broker, tiene los mismos tres roles, y mantenerlos diferenciados es lo que hace que el sistema sea mantenible.
El producer es cualquier código que añade un trabajo. Conoce el nombre del trabajo y la forma del payload, y nada más. Debe ser rápido, ya que normalmente se ejecuta dentro de una solicitud, y debe ser seguro llamarlo dos veces.
app.post("/signup", async (req, res) => {
const user = await db.user.create({ data: req.body });
await emails.add("welcome", { userId: user.id, to: user.email });
res.status(201).json({ id: user.id });
});
La queue es el estado compartido en el medio. En un despliegue de Redis, es el conjunto de claves que mantiene BullMQ: la lista de espera, el conjunto de retrasados, el conjunto activo y el conjunto de fallidos. Persiste los trabajos entre reinicios, los distribuye de forma atómica y rastrea los intentos.
El worker es un proceso de larga ejecución que extrae trabajos y ejecuta handlers. Está separado de la API por buenas razones: puede desplegarse en hardware optimizado para CPU, escalarse según la profundidad de la cola y reiniciarse sin perder solicitudes. Un solo proceso puede albergar workers para varias colas, y una sola cola puede ser atendida por muchos procesos worker. La queue es el único estado compartido, razón por la cual escala horizontalmente de manera tan limpia.
Por qué Redis es un almacenamiento común
Tres propiedades de Redis lo convierten en una cola natural.
Primero, las primitivas ya existen. Una lista con LPUSH y BRPOP es una cola, y BRPOP bloquea la conexión en lugar de consumir CPU. Esa es la razón principal por la cual las primeras librerías de colas para Node.js tenían solo unos pocos cientos de líneas.
LPUSH queue:emails "welcome:42"
BRPOP queue:emails 30
Segundo, los comandos son atómicos. Mover un trabajo de estado de espera a activo, incrementar su contador de intentos y establecer su bloqueo puede ocurrir sin condiciones de carrera, ya que Redis ejecuta los comandos uno a la vez. Esto es exactamente lo que necesita una cola para evitar que dos workers reclamen el mismo trabajo.
Tercero, probablemente Redis ya esté instalado. Es el caché y el almacenamiento de sesiones predeterminado para los servicios de Node.js. Reutilizarlo para una cola evita añadir una segunda pieza de infraestructura, un segundo conjunto de credenciales y un segundo manual de operaciones.
Los Streams llevan esta idea más allá. Mientras que en una lista solo se puede hacer push y pop, un stream es un log de solo anexado (append-only) con grupos de consumidores, confirmaciones por mensaje y un historial reproducible.
XADD jobs:emails '*' type welcome userId 42
XREADGROUP GROUP workers alice COUNT 10 STREAMS jobs:emails '>'
El inconveniente es la durabilidad. Redis es primordialmente un almacenamiento en memoria, por lo que un trabajo confirmado pero que aún no se ha escrito en disco puede perderse si la instancia falla. Puedes reducir esa ventana de riesgo con AOF y una réplica, pero no puedes hacer que Redis sea tan durable como una base de datos con write-ahead-log. Trata cada trabajo como una tarea recuperable, no como tu sistema de registro principal.
Primeros pasos con BullMQ
BullMQ necesita una conexión a Redis que tenga permitido reintentar indefinidamente. El comportamiento predeterminado de ioredis se rinde después de unos pocos fallos, lo cual es incorrecto para un worker que debe soportar una interrupción momentánea, por lo que debes configurar maxRetriesPerRequest: null.
import { Queue, Worker } from "bullmq";
import IORedis from "ioredis";
const connection = new IORedis(process.env.REDIS_URL!, {
maxRetriesPerRequest: null,
});
const emails = new Queue("emails", { connection });
const worker = new Worker(
"emails",
async (job) => {
await sendEmail(job.data.to, job.data.template);
},
{ connection, concurrency: 10 },
);
El Queue es el manejador del productor. worker es el consumidor. El nombre "emails" es la cola, y un worker solo ve los trabajos añadidos a su propia cola. Las colas separadas son la unidad de aislamiento: una cola de importación lenta no puede retrasar una cola de restablecimiento de contraseñas.
Un Queue es ligero y puede crearse donde sea que necesites añadir un trabajo. Un Worker es un proceso de larga ejecución y debe iniciarse una vez por proceso, no por solicitud. Comparte el objeto de conexión entre todas las colas y workers del proceso.
Agregar trabajos con datos y opciones
add recibe un nombre, un payload y un objeto de opciones. El nombre enruta el trabajo a un handler; el payload contiene todo lo que el worker necesita.
await emails.add("welcome", { userId, to }, {
jobId: `welcome:${userId}`,
});
El payload debe ser una captura de la intención, no un objeto vivo. Si un usuario cambia su correo electrónico entre el enqueue y la ejecución, el trabajo aún debería enviarse a la dirección con la que fue creado. Almacena ids y los pocos valores que definen la tarea en lugar de la fila completa de la base de datos, ya que la cola mantiene cada trabajo en espera en memoria.
El objeto de opciones es donde reside la mayor parte del valor de BullMQ:
attempts— el número máximo de intentos antes de que el trabajo se considere fallido.backoff— la estrategia de retraso entre intentos, comoexponential.delay— no ejecutar antes de que hayan pasado estos milisegundos desde ahora.priority— un número menor se ejecuta primero cuando hay trabajos en espera.jobId— un id estable que evita duplicados al encolar el mismo trabajo lógico.removeOnComplete— cuántos trabajos completados conservar, otruepara eliminarlos inmediatamente.removeOnFail— si se deben conservar los fallos para su inspección. Consérvalos.
Procesamiento de trabajos y retorno de resultados
Un handler recibe el trabajo, realiza la tarea y puede devolver un valor. El valor de retorno se almacena en el trabajo y puede leerse posteriormente, lo que convierte a la cola en un RPC asíncrono sencillo.
const worker = new Worker(
"reports",
async (job) => {
const pdf = await renderPdf(job.data.reportId);
return { url: pdf.url, bytes: pdf.bytes };
},
{ connection },
);
El llamador puede entonces consultar el trabajo para obtener su resultado o escuchar el flujo de eventos hasta que se complete.
const job = await reports.getJob(jobId);
if (await job?.isCompleted()) {
return job.returnvalue;
}
if (await job?.isFailed()) {
throw new Error(job.failedReason);
}
Devuelve valores pequeños. Un resultado se almacena en Redis como cualquier otro dato, por lo que devolver un PDF de un megabyte saturaría la cola. Devuelve una URL o un id y permite que el llamador recupere los bytes desde el almacenamiento de objetos.
Reintentos, exponential backoff y el failed set
Los fallos transitorios son normales. Una base de datos hace un failover, una API te aplica un rate-limit, un contenedor es reprogramado. Reintentar es la respuesta correcta, pero reintentar inmediatamente no lo es.
await queue.add("sync", { accountId }, {
attempts: 5,
backoff: { type: "exponential", delay: 2_000 },
});
Esto produce retrasos de aproximadamente 2s, 4s, 8s y 16s, aplicando el jitter propio de BullMQ para que muchos jobs que fallaron simultáneamente no se reintenten al mismo tiempo. Sin jitter, una flota de workers que se vieran afectados por la misma caída recrearían el pico de tráfico en el momento exacto de la recuperación.
Cuando un job agota sus intentos, BullMQ no lo elimina. Lo mueve al failed set, junto con el mensaje de error y el stack trace. Ese set es la superficie operativa de la cola: configura alertas cuando crezca y construye una ruta de re-ejecución (replay path).
const failed = await queue.getFailed(0, 20);
for (const job of failed) {
// Fix the underlying cause first, then requeue.
await job.retry();
}
Un job que falla siempre es un “poison message”. Reintentarlo infinitamente es peor que no reintentarlo en absoluto, ya que consume un slot de worker en cada intento y deja sin recursos al trabajo saludable. Limita los intentos e inspecciona el failed set en lugar de ignorarlo.
Trabajos retardados y repetibles
Un trabajo retardado (delayed job) se programa para ejecutarse más tarde sin necesidad de un programador externo. Un trabajo repetible se ejecuta siguiendo un patrón cron y pertenece a la cola.
await queue.add("reminder", { userId }, { delay: 600_000 });
await queue.add(
"digest",
{ region: "eu" },
{
repeat: { pattern: "0 7 * * *", tz: "Europe/Berlin" },
jobId: "digest:eu",
},
);
El jobId estable en un trabajo repetible es fundamental. Evita que el programador apile una nueva copia si hay una todavía en ejecución, y permite que cada instancia de la aplicación registre el mismo horario sin crear duplicados. Es preferible usar UTC o una zona horaria explícita para el patrón: un resumen que se envía a las 07:00 UTC no es lo mismo que uno a las 07:00 locales, y esa diferencia se traduce en un ticket de soporte cada vez que hay un cambio de horario de verano.
Los trabajos retardados se almacenan en un sorted set indexado por su hora de ejecución, por lo que un trabajo retardado no ocupa un worker. Permanece en Redis hasta que llega su momento, lo que hace que los retardos de horas o días sean muy eficientes.
Concurrencia, rate limiting y backpressure
La concurrencia de un worker es la cantidad de jobs que procesa simultáneamente. Aumentarla incrementa el throughput hasta que el worker agota la CPU, las conexiones a la base de datos o la memoria, momento en el cual la situación empeora.
const worker = new Worker("sync", handler, {
connection,
concurrency: 10,
limiter: { max: 50, duration: 1_000 },
});
El limiter limita cuántos jobs inicia el worker por ventana de tiempo. Esta es la primera línea de defensa para evitar saturar una API externa. Si un proveedor permite 50 solicitudes por segundo, un worker con una concurrencia de 200 provocará que te apliquen rate-limiting; un limitador de 50 por segundo evitará que esto suceda.
El backpressure es lo que evita que la cola crezca sin control. Si los productores añaden jobs más rápido de lo que los workers pueden procesarlos, la cola se convierte en un backlog creciente y su latencia pasa a ser de horas. Monitorea la profundidad de la cola, pausa los productores al superar un umbral y escala los workers automáticamente. Una cola que solo crece es una caída del servicio que nadie ha notado todavía.
Progreso y eventos
Los trabajos largos deben reportar su progreso para que la UI pueda mostrar una barra y el operador pueda diferenciar entre un proceso lento y uno bloqueado.
const worker = new Worker("imports", async (job) => {
const rows = await loadRows(job.data.fileId);
for (let i = 0; i < rows.length; i += 500) {
await insertChunk(rows.slice(i, i + 500));
await job.updateProgress(Math.round((i / rows.length) * 100));
}
return { rows: rows.length };
}, { connection });
Los eventos son la forma en que el resto del sistema observa la cola. Un worker emite completed, failed, progress y stalled para los trabajos que ejecuta. QueueEvents escucha el mismo flujo desde fuera del worker, lo que permite que un proceso de API reaccione cuando un trabajo finaliza sin necesidad de ser quien lo ejecutó.
worker.on("progress", (job, progress) => {
console.log(`job ${job.id} at ${progress}%`);
});
worker.on("completed", (job) => {
console.log(`job ${job.id} finished`);
});
Apagado gradual (graceful shutdown) y trabajos estancados
Un worker interrumpido a mitad de un trabajo deja dicho trabajo en un estado ambiguo. BullMQ gestiona esto mediante un bloqueo (lock): mientras un trabajo está activo, el worker renueva un lock en Redis. Si el worker muere y el lock expira, el trabajo se marca como stalled (estancado), se mueve de nuevo a espera y se procesa otra vez. Esta es una de las razones por las que la entrega es “at-least-once” (al menos una vez).
Dado que un estancamiento provoca que el trabajo se ejecute de nuevo, gestiona SIGTERM deliberadamente para que los trabajos en curso finalicen en lugar de ser interrumpidos.
process.on("SIGTERM", async () => {
await worker.close(); // stop accepting, wait for in-flight jobs
await connection.quit();
process.exit(0);
});
Asigna al despliegue un periodo de gracia lo suficientemente largo para el trabajo más lento, y limita los timeouts de los trabajos para que ninguno pueda superar dicho periodo. Un handler que pueda ejecutarse durante una hora debería registrar sus puntos de control (checkpoint) para que, en caso de re-ejecución, se reanude en lugar de reiniciarse.
Redis Streams directamente
BullMQ es una capa sobre las primitivas de Redis y, a veces, lo que necesitas son las primitivas. Los Streams son la herramienta adecuada cuando varios consumidores deben ver cada mensaje, cuando necesitas reproducir el historial o cuando quieres confirmaciones (acknowledgement) sin un framework de jobs.
XADD jobs:emails '*' type welcome userId 42
XREADGROUP GROUP workers alice COUNT 10 BLOCK 5000 STREAMS jobs:emails '>'
XACK jobs:emails workers 1760000000000-0
El grupo de consumidores rastrea qué mensajes han sido entregados y cuáles han sido confirmados. Un mensaje que es entregado pero nunca confirmado permanece en la lista de pendientes del grupo, evitando que un consumidor que haya fallado lo pierda. XAUTOCLAIM reasigna los mensajes pendientes de un consumidor inactivo a uno activo.
Usa streams directamente cuando los datos sean un registro de eventos y los consumidores sean lectores independientes. Usa BullMQ cuando los datos sean una unidad de trabajo con una política de reintentos, una prioridad y un resultado. Reimplementar reintentos, retrasos y un conjunto de fallos sobre streams es exactamente el trabajo que BullMQ ya realizó.
Idempotencia: la entrega es at-least-once
El hecho más importante sobre las colas es que la entrega es at-least-once (al menos una vez), no exactamente una vez. Un worker puede fallar después de realizar el trabajo pero antes de confirmarlo, y el job se ejecutará de nuevo. Por diseño, un bloqueo provoca una re-ejecución. La entrega exactly-once a través de una red es efectivamente imposible, por lo que las colas optan por at-least-once y trasladan esa responsabilidad a ti.
Tu handler debe ser idempotente: ejecutarlo dos veces debe producir el mismo estado final que ejecutarlo una sola vez. La técnica habitual es utilizar una clave de deduplicación derivada del trabajo, escrita de forma atómica antes del efecto secundario.
export async function handleCharge(job) {
const key = `charged:${job.data.orderId}`;
const first = await connection.set(key, "1", "NX", "EX", 86_400);
if (first === null) return { skipped: true, reason: "already_processed" };
await stripe.charges.create(
{ amount: job.data.amount, source: job.data.token },
{ idempotencyKey: job.data.orderId },
);
}
Observa que la clave proviene del id del pedido, no del id del job. Esto hace que el handler sea seguro incluso si el productor encola el mismo trabajo lógico dos veces. Combina la clave de deduplicación con la propia idempotency key del proveedor, ya que esa clave solo protege la llamada a la API, no la lógica circundante.
Monitoreo de la cola
Una cola es invisible a menos que la hagas visible. Cuatro señales cubren la mayor parte de lo que necesitas.
- Queue depth (Profundidad de la cola) — cuántos trabajos están esperando. Una profundidad creciente significa que los workers no pueden mantener el ritmo.
- Oldest waiting job (Trabajo en espera más antiguo) — la profundidad indica cuántos, la antigüedad indica qué tan grave es. Diez mil trabajos que se procesan en un segundo están bien; diez que han esperado una hora, no.
- Failure rate (Tasa de fallos) — fallos por minuto, desglosados por nombre del trabajo. Un pico después de un deploy apunta directamente al cambio realizado.
- Job duration (Duración del trabajo) — un histograma por nombre de trabajo. Un p95 creciente significa que una dependencia se está volviendo lenta.
BullMQ expone estos conteos directamente, y un listener QueueEvents o un exportador pueden enviarlos a tu sistema de métricas.
const counts = await queue.getJobCounts(
"wait", "active", "completed", "failed", "delayed",
);
console.log(counts);
Para una vista visual, Bull Board o el dashboard de BullMQ montan una pequeña web UI sobre las mismas claves. Añade un correlation id a cada payload e inclúyelo en los logs para que un trabajo pueda rastrearse desde la solicitud que lo creó a través de cada reintento.
Mantener los payloads pequeños y versionados
Una cola es una interfaz persistente entre dos despliegues. Un productor que ejecute la versión 1 del código puede escribir un trabajo que un worker que ejecute la versión 2 debe leer. Este es el mismo problema de compatibilidad que ocurre con una API, y es fácil ignorarlo hasta que un despliegue rompe el backlog.
Dos hábitos mantienen los payloads compatibles. Primero, añade campos en lugar de renombrarlos o eliminarlos, y asigna a los campos nuevos un valor por defecto sensato en el handler. Un worker que tolere la ausencia de locale puede procesar trabajos encolados antes de que el campo existiera. Segundo, incluye una versión en el payload cuando la estructura pueda cambiar significativamente, y utiliza una bifurcación (branch) basada en ella en el handler.
await queue.add("import", { version: 2, importId, mapping });
Mantén los payloads pequeños por una segunda razón: Redis almacena cada trabajo en espera en memoria. Un payload que incluya una fila completa de la base de datos se multiplica a través de miles de trabajos y encarece la cola. Haz referencia a los datos mediante su id y deja que el worker los recupere. La única excepción es un valor que deba quedar congelado al momento de encolar, como el destinatario de un email o el precio cotizado a un cliente, el cual pertenece al payload precisamente porque no debe cambiar.
Componiendo trabajo con flows
Algunos trabajos son, en realidad, varios trabajos. Un reporte podría obtener datos, renderizar un PDF y enviarlo por correo electrónico, y es posible que quieras que cada paso se reintente de forma independiente. Los flows de BullMQ expresan esto como un árbol de trabajos padres e hijos.
import { FlowProducer } from "bullmq";
const flow = new FlowProducer({ connection });
await flow.add({
name: "report",
queueName: "reports",
data: { reportId },
children: [
{ name: "fetch", queueName: "reports", data: { reportId } },
{ name: "render", queueName: "reports", data: { reportId } },
],
});
Un trabajo hijo se ejecuta primero, y el padre se vuelve ejecutable solo cuando todos sus hijos se completan. Esto te permite implementar fan-out y fan-in sin tener que coordinar el estado en tu propia base de datos. Si un hijo falla, el padre espera o falla según las opciones del flow, de modo que la política de reintentos permanece vinculada al paso que realmente falló. Mantén los flows superficiales y explícitos; un árbol profundo de trabajos interdependientes es más difícil de razonar que un pipeline pequeño con etapas claras.
Elegir entre colas de Redis, Kafka y RabbitMQ
Redis es la opción predeterminada ideal cuando la unidad de trabajo es un job: una tarea nombrada con un payload, una política de reintentos y un resultado. Es rápido, familiar y ya está presente en la mayoría de los stacks.
Kafka es la elección correcta cuando los datos son un log. Si varios consumidores independientes deben leer cada evento, si necesitas reproducir el historial desde un offset, o si el throughput se mide en millones de mensajes por segundo, un log particionado encaja mejor que una cola de jobs.
RabbitMQ es la elección correcta cuando el routing es la parte compleja. Los exchanges y las binding keys permiten que un mensaje se distribuya (fan out) a muchas colas mediante patrones, contando con acknowledgements por mensaje, prioridades y dead-letter exchanges como funcionalidades nativas. Es más pesado de operar que Redis y beneficia a los equipos que necesitan esa flexibilidad.
La recomendación honesta es empezar con lo que ya tienes implementado. Una cola de Redis es mucho mejor que no tener ninguna cola por estar esperando a evaluar brokers. Cambia cuando una limitación específica —durabilidad, replay o routing— realmente se convierta en un problema.
Mejores prácticas
- Mantén los manejadores de solicitudes limitados a una escritura y un encolado, y devuelve
202 Accepted. - Comparte una única conexión de Redis por proceso, utilizando
maxRetriesPerRequest: null. - Deriva
jobIda partir del trabajo para que los encolados duplicados se colapsen en un solo job. - Haz que cada manejador sea idempotente, ya que la entrega es “at-least-once”.
- Utiliza un backoff exponencial con jitter y establece un límite de intentos.
- Conserva los fallos configurando
removeOnFail: falsey limita los jobs completados conremoveOnComplete. - Separa las colas según la carga de trabajo para que los jobs lentos no bloqueen a los urgentes.
- Limita la concurrencia y añade un limiter para cualquier dependencia con rate-limit.
- Reporta el progreso de los jobs largos para que los lentos y los bloqueados se distingan entre sí.
- Apaga los workers mediante
SIGTERMy otorga a los despliegues un periodo de gracia equivalente. - Monitoriza la profundidad, la antigüedad del job más viejo, la tasa de fallos y la duración en un dashboard.
Errores comunes
- Asumir que un job se ejecuta exactamente una vez y terminar cobrando dos veces a un cliente.
- Usar un job id aleatorio y acumular schedules repetibles duplicados.
- Reintentar un poison message infinitamente y saturar la cola.
- Ejecutar tareas pesadas en el proceso de la API y llamarlo cola.
- Configurar una concurrencia tan alta que la base de datos alcance su límite de conexiones.
- Retornar un objeto enorme como resultado de un job y saturar Redis.
- Olvidar
maxRetriesPerRequest: nully perder jobs durante una reconexión. - Ignorar el failed set en todos los dashboards hasta que ocurre el primer incidente.
- Matar workers con
SIGKILLy forzar que todos los jobs en curso se detengan. - Incluir una fila completa de la base de datos en el payload en lugar de solo un id.
Próximos pasos
La guía de Redis cubre el almacenamiento que sustenta la cola: lists, streams, TTLs, persistencia y el modelo de comandos de un solo hilo que hace posibles los claims atómicos. Para profundizar en la disciplina de extraer el trabajo de la ruta de solicitud, lee Batch Processing. Cuando el enrutamiento y los acknowledgements por mensaje se vuelvan la parte compleja, continúa con RabbitMQ, y cuando necesites un log reproducible en lugar de una cola de trabajos, la guía de Kafka es la siguiente parada.