Ir para o conteúdo
Backend

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

Marcos Soares
Atualizado em 
15 minutos de leitura
Ilustracao 3D de tubos de vidro translucido com capsulas luminosas verdes representando filas de mensageria BullMQ e Redis
Ouça este artigo
0:00Mensageria 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érioBullMQ + RedisRabbitMQ (AMQP)Amazon SQSKafka
ProtocoloRedis commandsAMQP 0-9-1HTTP/SQS APIKafka protocol
Ordering garantidoPor fila (FIFO)Por filaFIFO opcional (custo extra)Por partição
Retry nativoSim, com backoff configurávelSim, com dead letter exchangeSim, com redrive policyNão nativo (consumer controla)
Rate limitingNativo por filaPluginNão nativoNão nativo
Infra necessáriaRedis (que você já tem)Broker dedicadoConta AWSCluster Kafka + ZooKeeper
Caso idealJobs assíncronos em aplicações Node.js com até ~50k jobs/minRoteamento complexo entre serviços heterogêneosServerless ou integração AWS nativaStreaming 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+.

Bash
npm install bullmq ioredis
npm install -D @types/node typescript

A conexão com Redis deve ser centralizada. Não crie múltiplas instâncias de IORedis espalhadas pelo código.

TypeScript
// 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.

TypeScript
// 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).

TypeScript
// 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).

TypeScript
// 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.

TypeScript
// 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).

JAVASCRIPT
// 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):

Text
[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:

  1. 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 com await entre eles, para o loop renovar o lock.
  2. Dimensionar o lock para o trabalho real: lockDuration maior que o pior caso legítimo do handler, e stalledInterval proporcional; o padrão de 30 s serve para a maioria das filas de I/O.
  3. 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, e maxStalledCount baixo para jobs que não podem repetir.
  4. Observar: escute stalled no QueueEvents e 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

TypeScript
// 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:

TypeScript
// 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

TypeScript
// 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:

TypeScript
// 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:

TypeScript
// 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:

TypeScript
// 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:

  1. Redis persistência: use appendonly yes no Redis. Sem persistência, um restart do Redis perde todos os jobs pendentes.
  2. Concurrency do worker: I/O bound (HTTP, banco) tolera 5-20. CPU bound (crypto, PDF) exige 1 por worker, escale horizontalmente.
  3. Backoff configurado: exponencial para APIs externas, fixo para operações internas com falha previsível.
  4. removeOnComplete ativo: sem isso, Redis acumula jobs completos indefinidamente e a memória cresce sem limite.
  5. Graceful shutdown: chame worker.close() no SIGTERM para que jobs em andamento terminem antes do processo morrer.
  6. Idempotência: se o job tem efeito colateral (cobrança, envio de email), a lógica do worker deve ser idempotente.
  7. 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.

Marcos Soares

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

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.