O que é uma fila de jobs?
Uma fila de jobs (job queue) é um buffer durável que fica entre o código que solicita o trabalho e o código que o executa. O produtor escreve um pequeno registro descrevendo uma unidade de trabalho e retorna imediatamente. Um worker lê esse registro posteriormente e executa a parte demorada. A fila é o que sobrevive a um restart, absorve picos de demanda e oferece um lugar para onde enviar falhas.
O Redis tem sido utilizado dessa forma por mais de uma década. Suas listas forneceram às primeiras bibliotecas de filas primitivos atômicos de LPUSH e BRPOP, e os blocking pops permitiram que um worker “dormisse” até que o trabalho chegasse, em vez de fazer polling em um loop. Posteriormente, os Streams adicionaram grupos de consumidores, acknowledgements e replay. Sobre esses primitivos, o BullMQ encapsula todo o ciclo de vida — jobs atrasados, retries, prioridades, rate limiting, agendamentos repetíveis e eventos — em uma pequena API TypeScript.
Se o Redis já sustenta seu cache ou suas sessões, adicionar uma fila é um passo simples. Essa conveniência é o principal motivo de ele ser a escolha padrão para times de Node.js, e a principal razão para entender tanto o que ele oferece quanto o que ele não oferece.
Uma fila são três papéis, não um
Todo sistema de fila, independentemente do broker, possui os mesmos três papéis, e mantê-los distintos é o que torna o sistema sustentável.
O producer é qualquer código que adiciona um job. Ele conhece o nome do job e o formato do payload, e nada mais. Ele deve ser rápido, pois geralmente é executado dentro de uma requisição, e deve ser seguro de ser chamado duas vezes.
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 });
});
A queue é o estado compartilhado no meio. Em uma implantação com Redis, é o conjunto de chaves que o BullMQ mantém: a lista de espera, o conjunto de atrasados (delayed set), o conjunto de ativos e o conjunto de falhas. Ela persiste jobs entre reinicializações, distribui-os atomicamente e rastreia as tentativas.
O worker é um processo de longa duração que consome jobs e executa handlers. Ele é separado da API por bons motivos: pode ser implantado em hardware otimizado para CPU, escalado de acordo com a profundidade da fila e reiniciado sem derrubar requisições. Um único processo pode hospedar workers para várias filas, e uma única fila pode ser atendida por muitos processos worker. A queue é o único estado compartilhado, e é por isso que ela escala horizontalmente de forma tão limpa.
Por que o Redis é um backing store comum
Três propriedades do Redis o tornam uma fila natural.
Primeiro, as primitivas já existem. Uma lista com LPUSH e BRPOP é uma fila, e BRPOP bloqueia a conexão em vez de consumir CPU. Esse é todo o motivo pelo qual as primeiras bibliotecas de fila para Node.js tinham apenas algumas centenas de linhas.
LPUSH queue:emails "welcome:42"
BRPOP queue:emails 30
Segundo, os comandos são atômicos. Mover um job de espera para ativo, incrementar sua contagem de tentativas e definir seu lock podem acontecer sem race conditions, porque o Redis executa os comandos um por um. É exatamente disso que uma fila precisa para evitar que dois workers reivindiquem o mesmo job.
Terceiro, o Redis provavelmente já está lá. Ele é o cache e store de sessão padrão para serviços Node.js. Reutilizá-lo para uma fila evita a necessidade de uma segunda peça de infraestrutura, um segundo conjunto de credenciais e um segundo runbook.
Streams expandem ainda mais essa ideia. Enquanto uma lista só pode sofrer push e pop, um stream é um log append-only com consumer groups, confirmações (acknowledgements) por mensagem e um histórico reproduzível.
XADD jobs:emails '*' type welcome userId 42
XREADGROUP GROUP workers alice COUNT 10 STREAMS jobs:emails '>'
O problema é a durabilidade. O Redis é, primariamente, um store em memória, portanto, um job confirmado mas ainda não gravado no disco pode ser perdido se a instância cair. Você pode diminuir essa janela com AOF e uma réplica, mas não consegue tornar o Redis tão durável quanto um banco de dados com write-ahead-log. Trate um job como um trabalho recuperável, e não como o seu sistema de registro oficial.
Primeiros passos com BullMQ
O BullMQ precisa de uma conexão Redis que tenha permissão para tentar a reconexão indefinidamente. O comportamento padrão do ioredis desiste de um comando após algumas falhas, o que é inadequado para um worker que precisa resistir a instabilidades momentâneas; portanto, configure 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 },
);
O Queue é o handle do produtor. O worker é o consumidor. O nome "emails" é a fila, e um worker só enxerga jobs adicionados à sua própria fila. Filas separadas são a unidade de isolamento: uma fila de importação lenta não pode atrasar uma fila de redefinição de senha.
Um Queue é leve e pode ser criado onde quer que você precise adicionar um job. Um Worker é um processo de longa duração e deve ser iniciado apenas uma vez por processo, não por requisição. Compartilhe o objeto de conexão entre todas as filas e workers do processo.
Adicionando jobs com dados e opções
add recebe um nome, um payload e um objeto de opções. O nome roteia o job para um handler; o payload carrega tudo o que o worker precisa.
await emails.add("welcome", { userId, to }, {
jobId: `welcome:${userId}`,
});
O payload deve ser um snapshot da intenção, não um objeto vivo. Se um usuário alterar seu e-mail entre o enqueue e a execução, o job ainda deve ser enviado para o endereço com o qual foi criado. Armazene ids e os poucos valores que definem o trabalho em vez da linha inteira do banco de dados, pois a fila mantém cada job aguardando na memória.
O objeto de opções é onde reside a maior parte do valor do BullMQ:
attempts— o número máximo de tentativas antes que o job seja considerado falho.backoff— a estratégia de delay entre as tentativas, comoexponential.delay— não execute antes de X milissegundos a partir de agora.priority— um número menor é executado primeiro quando há jobs aguardando.jobId— um id estável que remove a duplicidade de enqueues do mesmo trabalho lógico.removeOnComplete— quantos jobs concluídos manter, outruepara remover imediatamente.removeOnFail— se deve manter as falhas para inspeção. Mantenha-as.
Processando jobs e retornando resultados
Um handler recebe o job, executa o trabalho e pode retornar um valor. O valor de retorno é armazenado no job e pode ser lido posteriormente, transformando a fila em um RPC assíncrono simples.
const worker = new Worker(
"reports",
async (job) => {
const pdf = await renderPdf(job.data.reportId);
return { url: pdf.url, bytes: pdf.bytes };
},
{ connection },
);
O chamador pode então fazer o polling do job para obter seu resultado ou monitorar o stream de eventos para a conclusão.
const job = await reports.getJob(jobId);
if (await job?.isCompleted()) {
return job.returnvalue;
}
if (await job?.isFailed()) {
throw new Error(job.failedReason);
}
Retorne valores pequenos. Um resultado é armazenado no Redis como qualquer outro dado, portanto, retornar um PDF de um megabyte sobrecarrega a fila. Retorne uma URL ou um id e deixe que o chamador busque os bytes no object storage.
Retries, exponential backoff e o failed set
Falhas transitórias são normais. Um banco de dados sofre failover, uma API aplica rate-limit, um container é reagendado. Tentar novamente (retry) é a resposta correta, mas tentar imediatamente não é.
await queue.add("sync", { accountId }, {
attempts: 5,
backoff: { type: "exponential", delay: 2_000 },
});
Isso produz atrasos de aproximadamente 2s, 4s, 8s e 16s, com o jitter do próprio BullMQ aplicado para que muitos jobs que falharam juntos não tentem novamente ao mesmo tempo. Sem o jitter, uma frota de workers que atingissem a mesma interrupção recriaria o pico de carga no momento da recuperação.
Quando um job esgota suas tentativas, o BullMQ não o deleta. Ele o move para o failed set, junto com a mensagem de erro e o stack trace. Esse set é a superfície operacional da fila: configure alertas para quando ele crescer e construa um caminho de replay.
const failed = await queue.getFailed(0, 20);
for (const job of failed) {
// Fix the underlying cause first, then requeue.
await job.retry();
}
Um job que falha todas as vezes é uma “poison message”. Tentá-lo infinitamente é pior do que não tentar, pois ele consome um slot de worker a cada tentativa e prejudica o processamento de jobs saudáveis. Limite as tentativas e inspecione o failed set em vez de ignorá-lo.
Jobs atrasados e repetíveis
Um job atrasado (delayed job) é agendado para posteriormente sem a necessidade de um scheduler separado. Já um job repetível roda com base em um padrão cron e pertence à fila.
await queue.add("reminder", { userId }, { delay: 600_000 });
await queue.add(
"digest",
{ region: "eu" },
{
repeat: { pattern: "0 7 * * *", tz: "Europe/Berlin" },
jobId: "digest:eu",
},
);
O jobId estável em um job repetível é fundamental. Ele evita que o scheduler empilhe uma nova cópia caso uma ainda esteja em execução, e permite que cada instância do app registre o mesmo agendamento sem criar duplicatas. Prefira UTC ou um fuso horário explícito para o padrão: um resumo enviado às 07:00 UTC não é a mesma coisa que um às 07:00 local, e essa diferença resulta em um ticket de suporte a cada mudança de horário de verão.
Jobs atrasados são armazenados em um sorted set indexado pelo horário de execução, portanto, um job atrasado não ocupa um worker. Ele permanece no Redis até que chegue o momento de ser executado, o que torna atrasos de horas ou dias economicamente viáveis.
Concorrência, rate limiting e backpressure
A concorrência de um worker é a quantidade de jobs que ele processa simultaneamente. Aumentá-la eleva o throughput até que o worker esgote a CPU, as conexões com o banco de dados ou a memória, ponto em que a situação piora.
const worker = new Worker("sync", handler, {
connection,
concurrency: 10,
limiter: { max: 50, duration: 1_000 },
});
O limiter limita quantos jobs o worker inicia por janela de tempo. Essa é a primeira linha de defesa para evitar a sobrecarga de uma API downstream. Se um provedor permite 50 requisições por segundo, um worker com concorrência 200 fará você sofrer rate-limiting; um limiter de 50 por segundo não.
Backpressure é o que impede que a própria fila cresça indefinidamente. Se os produtores adicionam jobs mais rápido do que os workers conseguem processá-los, a fila torna-se um backlog crescente e sua latência passa a ser de horas. Monitore a profundidade da fila, pause os produtores acima de um determinado limite e escale os workers automaticamente. Uma fila que apenas cresce é uma queda de sistema que ninguém percebeu ainda.
Progresso e eventos
Jobs longos devem reportar o progresso para que a UI possa exibir uma barra e o operador consiga diferenciar se o processo está lento ou travado.
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 });
Eventos são a forma como o restante do sistema observa a fila. Um worker emite completed, failed, progress e stalled para os jobs que ele executa. QueueEvents escuta esse mesmo stream de fora do worker, permitindo que um processo de API reaja à finalização de um job sem necessariamente ter sido quem o executou.
worker.on("progress", (job, progress) => {
console.log(`job ${job.id} at ${progress}%`);
});
worker.on("completed", (job) => {
console.log(`job ${job.id} finished`);
});
Graceful shutdown e jobs travados (stalled)
Um worker encerrado no meio de um job deixa esse job em um estado ambíguo. O BullMQ resolve isso com um lock: enquanto um job está ativo, o worker renova um lock no Redis. Se o worker morrer e o lock expirar, o job é marcado como stalled, movido de volta para a fila de espera e processado novamente. Este é um dos motivos pelos quais a entrega é do tipo at-least-once.
Como um stall leva a uma nova execução, trate SIGTERM deliberadamente para que os jobs em andamento sejam finalizados em vez de serem interrompidos.
process.on("SIGTERM", async () => {
await worker.close(); // stop accepting, wait for in-flight jobs
await connection.quit();
process.exit(0);
});
Dê ao deploy um período de carência (grace period) longo o suficiente para o job mais lento e limite os timeouts dos jobs para que nenhum job individual dure mais que esse período. Um handler que possa rodar por uma hora deve salvar checkpoints do seu progresso, para que uma nova execução retome de onde parou em vez de reiniciar do zero.
Redis Streams diretamente
O BullMQ é uma camada sobre as primitivas do Redis e, às vezes, você precisa dessas primitivas. Streams são a ferramenta certa quando vários consumidores devem ver cada mensagem, quando você precisa reproduzir o histórico ou quando deseja confirmação (acknowledgement) sem a necessidade de um 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
O grupo de consumidores rastreia quais mensagens foram entregues e quais foram confirmadas. Uma mensagem que foi entregue, mas nunca confirmada, permanece na lista de pendências do grupo, garantindo que um consumidor que travou não a perca. XAUTOCLAIM reatribui mensagens pendentes de um consumidor inativo para um que esteja ativo.
Use streams diretamente quando os dados forem um log de eventos e os consumidores forem leitores independentes. Use BullMQ quando os dados forem uma unidade de trabalho com política de tentativa (retry), prioridade e um resultado. Reimplementar retries, delays e um conjunto de falhas (failed set) sobre streams é exatamente o trabalho que o BullMQ já realizou.
Idempotência: a entrega é at-least-once
O fato mais importante sobre filas é que a entrega é at-least-once (pelo menos uma vez), e não exatamente uma vez. Um worker pode travar após realizar o trabalho, mas antes de confirmá-lo, e o job será executado novamente. Um travamento leva a uma nova execução por design. A entrega exactly-once através de uma rede é efetivamente impossível, então as filas optam por at-least-once e transferem a responsabilidade para você.
Seu handler deve ser idempotente: executá-lo duas vezes deve produzir o mesmo estado final que executá-lo uma única vez. A técnica usual é utilizar uma chave de deduplicação derivada do trabalho, gravada atomicamente antes do efeito colateral.
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 },
);
}
Observe que a chave vem do id do pedido, não do id do job. Isso torna o handler seguro mesmo que o produtor enfileire o mesmo trabalho lógico duas vezes. Combine a chave de deduplicação com a própria idempotency key do provedor, pois essa chave protege apenas a chamada da API, não a lógica ao redor dela.
Monitorando a fila
Uma fila é invisível a menos que você a torne visível. Quatro sinais cobrem a maior parte do que você precisa.
- Queue depth (Profundidade da fila) — quantos jobs estão aguardando. Uma profundidade crescente significa que os workers não estão conseguindo acompanhar a demanda.
- Oldest waiting job (Job mais antigo aguardando) — a profundidade diz quantos, a idade diz o quão grave é a situação. Dez mil jobs que são processados em um segundo estão ok; dez que esperaram por uma hora não estão.
- Failure rate (Taxa de falha) — falhas por minuto, divididas por nome do job. Um pico após um deploy aponta para a alteração realizada.
- Job duration (Duração do job) — um histograma por nome do job. Um p95 crescente significa que alguma dependência está ficando lenta.
O BullMQ expõe essas contagens diretamente, e um listener QueueEvents ou um exportador pode encaminhá-las para o seu sistema de métricas.
const counts = await queue.getJobCounts(
"wait", "active", "completed", "failed", "delayed",
);
console.log(counts);
Para uma visualização gráfica, o Bull Board ou o dashboard do BullMQ montam uma pequena web UI sobre as mesmas chaves. Adicione um correlation id a cada payload e inclua-o nos logs para que um job possa ser rastreado desde a requisição que o criou até cada tentativa de reprocessamento (retry).
Mantendo payloads pequenos e versionados
Uma fila é uma interface persistente entre dois deploys. Um produtor executando a versão 1 do código pode escrever um job que um worker executando a versão 2 deve ler. Este é o mesmo problema de compatibilidade de uma API, e é fácil ignorá-lo até que um deploy quebre um backlog.
Dois hábitos mantêm os payloads compatíveis. Primeiro, adicione campos em vez de renomeá-los ou removê-los, e atribua um valor padrão sensato aos novos campos no handler. Um worker que tolera a ausência de locale consegue processar jobs enfileirados antes de o campo existir. Segundo, inclua uma versão no payload quando a estrutura puder mudar significativamente e utilize ramificações (branching) no handler com base nela.
await queue.add("import", { version: 2, importId, mapping });
Mantenha os payloads pequenos por um segundo motivo: o Redis armazena cada job em espera na memória. Um payload que incorpora uma linha inteira do banco de dados se multiplica por milhares de jobs e torna a fila cara. Referencie os dados por id e deixe que o worker os busque. A única exceção é um valor que deve ser congelado no momento do enfileiramento, como o destinatário de um e-mail ou o preço cotado para um cliente, que deve estar no payload precisamente porque não deve mudar.
Compondo fluxos de trabalho com flows
Alguns jobs são, na verdade, vários jobs. Um relatório pode buscar dados, renderizar um PDF e enviá-lo por e-mail, e você pode querer que cada etapa seja repetida independentemente. Os flows do BullMQ expressam isso como uma árvore de jobs pais e filhos.
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 } },
],
});
Um job filho é executado primeiro, e o pai torna-se executável apenas quando todos os seus filhos forem concluídos. Isso permite implementar fan-out e fan-in sem a necessidade de coordenar o estado no seu próprio banco de dados. Se um filho falhar, o pai aguarda ou falha de acordo com as opções do flow, mantendo a política de retry vinculada à etapa que realmente apresentou erro. Mantenha os flows rasos e explícitos; uma árvore profunda de jobs interdependentes é mais difícil de analisar do que um pipeline pequeno com etapas claras.
Escolhendo entre filas Redis, Kafka e RabbitMQ
O Redis é a escolha padrão ideal quando a unidade de trabalho é um job: uma tarefa nomeada com um payload, uma política de tentativa (retry) e um resultado. Ele é rápido, familiar e já está presente na maioria das stacks.
O Kafka é a escolha certa quando os dados são um log. Se vários consumidores independentes precisam ler cada evento, se você precisa reproduzir o histórico a partir de um offset, ou se a vazão é medida em milhões de mensagens por segundo, um log particionado se adapta melhor do que uma fila de jobs.
O RabbitMQ é a escolha certa quando o roteamento é a parte complexa. Exchanges e binding keys permitem que uma mensagem seja distribuída para várias filas por padrão, com confirmações (acknowledgements) por mensagem, prioridades e dead-letter exchanges como recursos nativos. Ele é mais pesado de operar do que o Redis e beneficia equipes que precisam dessa flexibilidade.
A orientação sincera é começar com o que você já utiliza. Uma fila Redis é muito melhor do que nenhuma fila enquanto você espera para avaliar brokers. Mude quando uma limitação específica — durabilidade, replay ou roteamento — realmente se tornar um problema.
Melhores práticas
- Mantenha os request handlers limitados a uma escrita e um enqueue, e retorne
202 Accepted. - Compartilhe uma única conexão Redis por processo, utilizando
maxRetriesPerRequest: null. - Derive
jobIda partir do trabalho para que enqueues duplicados sejam colapsados em um único job. - Torne cada handler idempotente, pois a entrega é do tipo at-least-once.
- Use exponential backoff com jitter e defina um limite de tentativas.
- Mantenha o registro de falhas configurando
removeOnFail: falsee limite os jobs concluídos comremoveOnComplete. - Separe as filas por carga de trabalho para que jobs lentos não causem starvation em jobs urgentes.
- Limite a concorrência e adicione um limiter para qualquer dependência com rate-limit.
- Reporte o progresso de jobs longos para que jobs lentos e travados sejam distinguidos.
- Desligue os workers no
SIGTERMe defina um grace period correspondente para os deploys. - Monitore a profundidade da fila, a idade do job mais antigo, a taxa de falhas e a duração em um dashboard.
Erros comuns
- Assumir que um job é executado exatamente uma vez e acabar cobrando o cliente duas vezes.
- Usar um job id aleatório e empilhar agendamentos repetíveis duplicados.
- Tentar reprocessar uma “poison message” infinitamente e causar o esgotamento da fila.
- Executar tarefas pesadas no processo da API e chamar isso de fila.
- Definir a concorrência tão alta que o banco de dados atinge o limite de conexões.
- Retornar um objeto enorme como resultado de um job e inflar o Redis.
- Esquecer o
maxRetriesPerRequest: nulle perder jobs em uma reconexão. - Deixar o conjunto de falhas fora de todos os dashboards até acontecer o primeiro incidente.
- Derrubar workers com
SIGKILLe forçar a interrupção de todos os jobs em execução. - Colocar uma linha inteira do banco de dados no payload em vez de apenas um id.
Próximos passos
O guia de Redis aborda a camada de armazenamento por trás da fila: lists, streams, TTLs, persistência e o modelo de comandos single-threaded que torna as claims atômicas possíveis. Para entender a disciplina mais ampla de remover o processamento do caminho da requisição, leia sobre Batch Processing. Quando o roteamento e as confirmações (acknowledgements) por mensagem se tornarem a parte complexa, continue para RabbitMQ, e quando você precisar de um log reproduzível em vez de uma fila de jobs, o guia de Kafka é a próxima parada.