Mensageria com BullMQ e Redis: Processamento Assíncrono que Funciona

Conteúdo técnico toda semana
Receba artigos sobre arquitetura, padrões de projeto e engenharia de software. Direto no seu e-mail, sem enrolação.
Sem spam. Cancele a qualquer momento com 1 clique.
Neste artigo
Um endpoint que envia email, gera PDF e notifica um webhook externo leva 4 segundos para responder. O usuário espera. O load balancer faz timeout. O retry do cliente duplica a operação. Tudo isso porque o processamento acontece de forma síncrona dentro do ciclo de request/response.
BullMQ resolve esse problema com uma abstração de filas sobre Redis que é simples de operar e previsível de debugar. Diferente de soluções como RabbitMQ ou Amazon SQS, BullMQ roda no mesmo ecossistema Node.js, usa Redis como único backing store e oferece retry, prioridade, rate limiting e dead letter queue sem infraestrutura adicional.
Por que BullMQ e não as alternativas
A escolha de mensageria depende de três variáveis: volume de mensagens, complexidade de roteamento e quanto de infraestrutura você quer gerenciar.
| Critério | BullMQ + Redis | RabbitMQ (AMQP) | Amazon SQS | Kafka |
|---|---|---|---|---|
| Protocolo | Redis commands | AMQP 0-9-1 | HTTP/SQS API | Kafka protocol |
| Ordering garantido | Por fila (FIFO) | Por fila | FIFO opcional (custo extra) | Por partição |
| Retry nativo | Sim, com backoff configurável | Sim, com dead letter exchange | Sim, com redrive policy | Não nativo (consumer controla) |
| Rate limiting | Nativo por fila | Plugin | Não nativo | Não nativo |
| Infra necessária | Redis (que você já tem) | Broker dedicado | Conta AWS | Cluster Kafka + ZooKeeper |
| Caso ideal | Jobs assíncronos em aplicações Node.js com até ~50k jobs/min | Roteamento complexo entre serviços heterogêneos | Serverless ou integração AWS nativa | Streaming de eventos com replay |
Se sua aplicação é Node.js, já usa Redis para cache ou sessão, e o volume fica abaixo de 50 mil jobs por minuto: BullMQ é a escolha direta. Acima disso, ou com necessidade de replay de eventos, Kafka entra na conversa. Para roteamento fan-out complexo entre serviços em linguagens diferentes, RabbitMQ faz mais sentido.
Setup: Redis, BullMQ e a primeira fila
Instale as dependências. BullMQ precisa de Redis 5.0+ (usa Streams internamente) e Node.js 18+.
npm install bullmq ioredis
npm install -D @types/node typescriptA conexão com Redis deve ser centralizada. Não crie múltiplas instâncias de IORedis espalhadas pelo código.
// src/lib/redis-connection.ts
import IORedis from "ioredis";
// maxRetriesPerRequest: null é obrigatório para BullMQ.
// O default do ioredis é 20, e BullMQ usa blocking commands (BRPOPLPUSH)
// que não devem ter limite de retry por request.
export const redisConnection = new IORedis(
process.env.REDIS_URL || "redis://localhost:6379",
{
maxRetriesPerRequest: null,
enableReadyCheck: false,
}
);Agora, a primeira fila. O padrão que funciona é separar a definição da fila (producer) do worker (consumer) em módulos distintos. Isso permite escalar workers independentemente do servidor HTTP.
// src/queues/email.queue.ts
import { Queue } from "bullmq";
import { redisConnection } from "../lib/redis-connection";
export interface EmailJobData {
to: string;
subject: string;
templateId: string;
variables: Record<string, string>;
}
export const emailQueue = new Queue<EmailJobData>("email", {
connection: redisConnection,
defaultJobOptions: {
attempts: 3,
// Backoff exponencial: 1s, 2s, 4s entre tentativas.
// Evita sobrecarregar o serviço de email durante instabilidade.
backoff: {
type: "exponential",
delay: 1000,
},
// Remove jobs completos após 24h para não inflar a memória do Redis.
removeOnComplete: { age: 86400 },
// Mantém jobs falhados por 7 dias para análise.
removeOnFail: { age: 604800 },
},
});Workers: onde o processamento acontece
O worker é um processo separado. Pode rodar no mesmo servidor ou em máquinas dedicadas. BullMQ garante que cada job é processado por exatamente um worker (at-least-once delivery com locking via Redis).
// src/workers/email.worker.ts
import { Worker, Job } from "bullmq";
import { redisConnection } from "../lib/redis-connection";
import { EmailJobData } from "../queues/email.queue";
// Simula integração com serviço de email real
async function sendEmail(data: EmailJobData): Promise<{ messageId: string }> {
const response = await fetch("https://api.resend.com/emails", {
method: "POST",
headers: {
Authorization: `Bearer ${process.env.RESEND_API_KEY}`,
"Content-Type": "application/json",
},
body: JSON.stringify({
from: "[email protected]",
to: data.to,
subject: data.subject,
html: `<p>Template: ${data.templateId}</p>`,
}),
});
if (!response.ok) {
throw new Error(`Resend API retornou ${response.status}`);
}
return response.json();
}
const emailWorker = new Worker<EmailJobData>(
"email",
async (job: Job<EmailJobData>) => {
const result = await sendEmail(job.data);
// O retorno fica acessível via job.returnvalue para debugging
return { messageId: result.messageId, processedAt: new Date().toISOString() };
},
{
connection: redisConnection,
// concurrency controla quantos jobs este worker processa em paralelo.
// Para I/O bound (chamadas HTTP), 5-10 é seguro.
// Para CPU bound (geração de PDF), use 1 e escale com mais workers.
concurrency: 5,
limiter: {
// Rate limit: máximo 100 emails por minuto neste worker.
// Respeita os limites da API do provedor de email.
max: 100,
duration: 60000,
},
}
);
emailWorker.on("completed", (job) => {
console.log(`Job ${job.id} concluído: ${job.returnvalue?.messageId}`);
});
emailWorker.on("failed", (job, err) => {
console.error(`Job ${job?.id} falhou após ${job?.attemptsMade} tentativas:`, err.message);
});Rode o worker como processo separado: npx tsx src/workers/email.worker.ts. Em produção, use o mesmo entrypoint que seu servidor HTTP ou um processo dedicado gerenciado por PM2, Docker ou Kubernetes. Se você já usa Docker, o post sobre Docker para devs cobre como rodar múltiplos serviços com docker-compose.
Enfileirando jobs a partir da API
O producer é leve. Adicionar um job à fila leva menos de 1ms (é um XADD no Redis).
// src/routes/user.routes.ts
import { Router, Request, Response } from "express";
import { emailQueue } from "../queues/email.queue";
const router = Router();
router.post("/users", async (req: Request, res: Response) => {
const { email, name } = req.body;
// Persiste o usuário no banco (omitido por brevidade)
const userId = "user_123";
// Enfileira o email de boas-vindas.
// O endpoint responde em <50ms independente de quanto tempo o email leva.
await emailQueue.add(
"welcome-email", // nome do job (útil para filtrar no dashboard)
{
to: email,
subject: `Bem-vindo, ${name}`,
templateId: "welcome-v2",
variables: { name, userId },
},
{
// jobId previne duplicação: se o mesmo userId tentar enfileirar
// o welcome-email novamente, BullMQ ignora silenciosamente.
jobId: `welcome-${userId}`,
// delay de 5 minutos: envia o email após o usuário ter tempo
// de completar o onboarding, não imediatamente após o cadastro.
delay: 5 * 60 * 1000,
}
);
res.status(201).json({ id: userId });
});
export { router as userRouter };Esse padrão de desacoplar a resposta HTTP do processamento pesado é o mesmo que aparece em pipelines de conteúdo automatizado, onde a geração de conteúdo com IA leva segundos mas o endpoint precisa responder rápido.
Filas com prioridade e jobs agendados
BullMQ suporta prioridade numérica (menor número = maior prioridade) e cron jobs nativos.
// src/queues/notification.queue.ts
import { Queue } from "bullmq";
import { redisConnection } from "../lib/redis-connection";
interface NotificationJobData {
userId: string;
channel: "push" | "sms" | "email";
payload: Record<string, unknown>;
}
export const notificationQueue = new Queue<NotificationJobData>("notification", {
connection: redisConnection,
});
// Job de alta prioridade: alerta de segurança
await notificationQueue.add(
"security-alert",
{ userId: "user_456", channel: "sms", payload: { type: "login_suspicious" } },
{ priority: 1 } // processado antes de jobs com priority > 1
);
// Job de baixa prioridade: newsletter semanal
await notificationQueue.add(
"weekly-digest",
{ userId: "user_456", channel: "email", payload: { type: "digest" } },
{ priority: 10 }
);
// Job recorrente: limpeza de sessões expiradas a cada hora
await notificationQueue.upsertJobScheduler(
"cleanup-sessions",
{ pattern: "0 * * * *" }, // cron: a cada hora cheia
{
name: "session-cleanup",
data: { userId: "system", channel: "email", payload: { type: "cleanup" } },
}
);Se você implementa notificações em tempo real com WebSockets, a fila de notificações serve como buffer entre o evento e o push via socket: o worker consome o job e dispara o evento WebSocket.
Stalled jobs na prática: reproduzindo o problema e o que fazer
Um job fica "stalled" quando o worker que o pegou deixa de renovar o lock no Redis. A causa clássica não é Redis nem rede: é código CPU-bound que bloqueia o event loop do worker, porque a renovação do lock roda no mesmo loop. O script abaixo reproduz isso com BullMQ 6.3.4 num Redis 7 descartável (docker run -d -p 6380:6379 redis:7-alpine).
// stalled.mjs
import { Queue, Worker, QueueEvents } from 'bullmq'
const connection = { host: '127.0.0.1', port: 6380 }
const queue = new Queue('demo-stalled', { connection })
const events = new QueueEvents('demo-stalled', { connection })
await queue.obliterate({ force: true })
const worker = new Worker('demo-stalled', async (job) => {
console.log(`[worker] iniciou job ${job.id} (${job.name})`)
if (job.name === 'bloqueia-event-loop') {
const fim = Date.now() + 3000
while (Date.now() < fim) { /* CPU-bound: sem await, o lock não renova */ }
} else {
await new Promise((r) => setTimeout(r, 300))
}
return 'ok'
}, { connection, lockDuration: 1000, stalledInterval: 500, maxStalledCount: 0, concurrency: 1 })
events.on('stalled', ({ jobId }) => console.log(`[events] stalled: job ${jobId}`))
events.on('failed', ({ jobId, failedReason }) => console.log(`[events] failed: job ${jobId}: ${failedReason}`))
events.on('completed', ({ jobId }) => console.log(`[events] completed: job ${jobId}`))
worker.on('error', (err) => console.log(`[worker] error: ${err.message}`))
await queue.add('cooperativo', { n: 1 })
await queue.add('bloqueia-event-loop', { n: 2 })
await queue.add('cooperativo', { n: 3 })
setTimeout(async () => {
console.log('[fim] contagem:', JSON.stringify(await queue.getJobCounts('completed', 'failed', 'active', 'waiting')))
await worker.close(); await events.close(); await queue.close(); process.exit(0)
}, 7000)Saída real (node stalled.mjs):
[worker] iniciou job 1 (cooperativo)
[worker] iniciou job 2 (bloqueia-event-loop)
[events] completed: job 1
[worker] error: Missing lock for job 2. moveToFinished
[worker] error: Missing lock for job 2. moveToFinished
[worker] iniciou job 3 (cooperativo)
[events] completed: job 3
[events] stalled: job 2
[events] failed: job 2: job stalled more than allowable limit
[fim] contagem: {"completed":2,"failed":1,"active":0,"waiting":0}Leitura da sequência: o job 2 travou o loop por 3 s com lockDuration de 1 s; quando o handler terminou, o worker tentou concluir e recebeu Missing lock for job 2, porque o lock já tinha expirado; o verificador de stalled (stalledInterval) marcou o job e, como maxStalledCount é 0, ele foi para failed em vez de ser reprocessado. Com o padrão (maxStalledCount: 1), o mesmo job seria executado de novo, ou seja, processado duas vezes: é assim que efeitos colaterais duplicados aparecem em produção.
O que fazer, em ordem:
- Tirar CPU-bound do worker: trabalho pesado vai para um
worker_threads, para um processo separado (useWorkerThreads/sandboxed processors do BullMQ) ou é quebrado em passos comawaitentre eles, para o loop renovar o lock. - Dimensionar o lock para o trabalho real:
lockDurationmaior que o pior caso legítimo do handler, estalledIntervalproporcional; o padrão de 30 s serve para a maioria das filas de I/O. - Tornar o handler idempotente: enquanto stalled existir, reprocessamento existe. Chave de idempotência por
job.id(ou por um identificador do domínio) antes de qualquer efeito colateral, emaxStalledCountbaixo para jobs que não podem repetir. - Observar: escute
stallednoQueueEventse alerte; um stalled por dia é bug de código, não de infraestrutura.
O que NÃO fazer
Anti-pattern 1: processar dentro do request handler
// ERRADO: o endpoint fica bloqueado até o PDF ser gerado
router.post("/reports", async (req, res) => {
const pdf = await generatePDF(req.body); // leva 8 segundos
await uploadToS3(pdf); // mais 2 segundos
await sendEmail(req.body.email, pdf); // mais 1 segundo
res.json({ url: pdf.url }); // 11 segundos depois
});O timeout do load balancer (tipicamente 30s) vai derrubar esse request em cenários de carga. O correto:
// CORRETO: enfileira e responde imediatamente
router.post("/reports", async (req, res) => {
const job = await reportQueue.add("generate", req.body);
// Retorna o ID do job para o cliente consultar o status via polling ou webhook
res.status(202).json({ jobId: job.id, statusUrl: `/reports/status/${job.id}` });
});O status code 202 Accepted comunica ao cliente que o processamento foi aceito, mas não concluído. Isso se integra com padrões de API Gateway onde o gateway pode cachear o status do job.
Anti-pattern 2: não tratar idempotência
// ERRADO: se o worker crashar após cobrar mas antes de confirmar,
// o retry vai cobrar o usuário duas vezes
const paymentWorker = new Worker("payment", async (job) => {
await chargeCard(job.data.cardToken, job.data.amount);
await updateOrderStatus(job.data.orderId, "paid");
});A correção exige uma chave de idempotência:
// CORRETO: verifica se a cobrança já foi feita antes de processar
const paymentWorker = new Worker("payment", async (job) => {
const existing = await db.payment.findUnique({
where: { idempotencyKey: job.data.idempotencyKey },
});
if (existing) {
// Job já foi processado. Retorna sem cobrar novamente.
return { alreadyProcessed: true, paymentId: existing.id };
}
const charge = await chargeCard(job.data.cardToken, job.data.amount);
// Persiste com a chave de idempotência ANTES de confirmar o job.
await db.payment.create({
data: {
idempotencyKey: job.data.idempotencyKey,
chargeId: charge.id,
amount: job.data.amount,
orderId: job.data.orderId,
},
});
await updateOrderStatus(job.data.orderId, "paid");
return { paymentId: charge.id };
});Esse padrão se conecta com a ideia de CQRS e Event Sourcing, onde cada evento tem um ID único que previne processamento duplicado.
Anti-pattern 3: ignorar jobs mortos
Sem dead letter queue configurada, jobs que falham após todas as tentativas desaparecem silenciosamente. Configure um handler para failed com attemptsMade === opts.attempts:
// src/workers/email.worker.ts (adicione ao worker existente)
emailWorker.on("failed", async (job, err) => {
if (job && job.attemptsMade === job.opts.attempts) {
// Job esgotou todas as tentativas. Registra em tabela de dead letters
// para investigação manual ou reprocessamento posterior.
await db.deadLetterJob.create({
data: {
queue: "email",
jobId: job.id ?? "unknown",
jobName: job.name,
payload: JSON.stringify(job.data),
error: err.message,
failedAt: new Date(),
},
});
}
});Monitoramento com BullMQ Dashboard
BullMQ oferece o @bull-board para visualização. A instalação é direta:
// src/dashboard.ts
import { createBullBoard } from "@bull-board/api";
import { BullMQAdapter } from "@bull-board/api/bullMQAdapter";
import { ExpressAdapter } from "@bull-board/express";
import { emailQueue } from "./queues/email.queue";
import { notificationQueue } from "./queues/notification.queue";
export function setupDashboard() {
const serverAdapter = new ExpressAdapter();
serverAdapter.setBasePath("/admin/queues");
createBullBoard({
queues: [
new BullMQAdapter(emailQueue),
new BullMQAdapter(notificationQueue),
],
serverAdapter,
});
return serverAdapter.getRouter();
}Proteja essa rota com autenticação. O dashboard expõe dados sensíveis dos jobs. O post sobre segurança em APIs Node.js cobre middleware de autenticação que se aplica aqui.
Checklist de produção
Antes de fazer deploy de filas BullMQ, verifique:
- Redis persistência: use
appendonly yesno Redis. Sem persistência, um restart do Redis perde todos os jobs pendentes. - Concurrency do worker: I/O bound (HTTP, banco) tolera 5-20. CPU bound (crypto, PDF) exige 1 por worker, escale horizontalmente.
- Backoff configurado: exponencial para APIs externas, fixo para operações internas com falha previsível.
- removeOnComplete ativo: sem isso, Redis acumula jobs completos indefinidamente e a memória cresce sem limite.
- Graceful shutdown: chame
worker.close()noSIGTERMpara que jobs em andamento terminem antes do processo morrer. - Idempotência: se o job tem efeito colateral (cobrança, envio de email), a lógica do worker deve ser idempotente.
- Alertas em dead letter: jobs que esgotam tentativas precisam gerar alerta, não sumir silenciosamente.
Para deploys seguros dessas mudanças, estratégias de deploy como rolling updates garantem que workers antigos terminam seus jobs antes de serem substituídos.
FAQ
BullMQ funciona com Redis Cluster? Sim, mas com restrições. Todas as filas de um mesmo nome precisam estar no mesmo shard (BullMQ usa hash tags para isso). A documentação oficial do BullMQ cobre a configuração de prefix para garantir co-localização. Redis Cluster é necessário acima de ~16GB de dados de fila ou quando você precisa de alta disponibilidade com failover automático.
Posso usar o mesmo processo para servidor HTTP e worker? Pode, e para aplicações com menos de 1000 jobs por hora, funciona. Acima disso, separe. O worker compete por CPU e memória com o event loop do servidor HTTP. Se um job pesado bloqueia o event loop (operação CPU bound sem worker thread), suas rotas HTTP ficam lentas. O runtime Node.js executa JavaScript em uma única thread: o event loop do V8 não distingue entre "código do worker BullMQ" e "código do request handler Express".
Como migrar de Bull (v3) para BullMQ (v4+)?
A API é diferente. Queue, Worker e QueueScheduler são classes separadas no BullMQ (QueueScheduler foi removido na v4, suas responsabilidades foram absorvidas pelo Worker). A migração exige reescrever producers e consumers, mas a estrutura de dados no Redis é compatível: você pode rodar Bull e BullMQ lado a lado durante a transição, desde que usem nomes de fila diferentes.
BullMQ garante entrega exatamente uma vez (exactly-once)? Não. BullMQ garante at-least-once. Se o worker processa o job mas crasha antes de confirmar (ack), o job será reprocessado. A responsabilidade de idempotência é do código do worker, não da fila. Isso é verdade para qualquer sistema de mensageria que prioriza durabilidade.
Qual o limite prático de jobs por segundo? Depende do tamanho do payload e da latência do Redis. Com Redis local e payloads de 1KB, BullMQ sustenta ~10k jobs/segundo por fila em benchmarks. Com Redis remoto (latência de 2-5ms), esse número cai para ~2-5k jobs/segundo. Se você precisa de mais, particione em múltiplas filas por tipo de job.
Posição editorial
BullMQ é a melhor escolha para processamento assíncrono em aplicações Node.js que já usam Redis. A barreira de entrada é baixa, o modelo mental é simples (fila, producer, worker) e a operação é previsível. Mas existe uma armadilha: equipes adotam filas para tudo e acabam com 30 filas diferentes, cada uma com configuração própria, sem monitoramento centralizado e sem padrão de retry. Filas são infraestrutura, não feature. Trate-as como migrations de banco: versionadas, testadas e com rollback planejado. Se você não sabe explicar por que um job específico precisa ser assíncrono, ele não precisa.

Escrito por
Marcos Soares
Fullstack Developer · CEO da Agência Poti
Fullstack Developer e CEO da Agência Poti. Mais de 20 anos construindo arquiteturas cloud-native com React, Next.js e sistemas distribuídos. Parceiro comercial do estúdio iellou design. Fundador do Vivo de Código.
Comentários
Participe da discussão
Seja o primeiro a comentar!
Continue Aprofundando

API Gateway Patterns: Autenticação, Rate Limiting e Caching com Node.js

Fetch, Retry e Timeout: O Código que Falta entre Sua API e a Realidade

Arquitetando Comunicação em Tempo Real em Escala: Redis Pub/Sub, Load Balancing e Milhares de Conexões Simultâneas
Guias de integração relacionados
Conteúdo técnico toda semana
Receba artigos sobre arquitetura, padrões de projeto e engenharia de software. Direto no seu e-mail, sem enrolação.
Sem spam. Cancele a qualquer momento com 1 clique.