Processamento em Background com BullMQ e Redis
BullMQ (evolução do Bull): filas tipadas, Workers com concurrency, retry com backoff exponencial, jobs recorrentes com cron, Bull Board para monitoramento e separação de processos API/Worker em Docker.
Quando um endpoint envia e-mail, gera PDF, processa vídeo ou faz chamadas a APIs externas de forma síncrona, o usuário fica bloqueado esperando a conclusão. Com 50 requisições simultâneas desse tipo, a API fica congestionada. A solução é desacoplar: a API publica um job na fila e responde imediatamente com 202 Accepted. O Worker processa em background, no seu próprio ritmo.
BullMQ é a versão moderna do Bull (arquitetura reescrita, TypeScript nativo, melhor garantia de entrega). Com Redis como broker, ele persiste os jobs, distribui entre múltiplos workers, gerencia retries com backoff, e executa jobs agendados com expressões cron.
npm install bullmq ioredisTipagem dos Jobs
Uma das vantagens do BullMQ sobre o Bull original é o suporte nativo a TypeScript Generics. Ao tipar a fila e o worker com a interface do job, o TypeScript garante que o producer (quem publica) e o consumer (worker) estejam em sincronia: se você adicionar um campo obrigatório à interface, o compilador aponta todos os lugares onde o job é criado sem esse campo.
// Tipos de jobs disponíveis nas filas — TypeScript garante consistência
export interface SendWelcomeEmailJobData {
userId: string;
email: string;
name: string;
}
export interface GenerateReportJobData {
reportId: string;
userId: string;
startDate: string;
endDate: string;
format: 'pdf' | 'csv' | 'xlsx';
}
export interface ProcessImageJobData {
fileKey: string;
userId: string;
operations: Array<'resize' | 'compress' | 'watermark'>;
}Configuração da Fila e Producer
A fila (Queue) é responsável apenas por receber e armazenar jobs — ela nunca processa nada. Isso permite que qualquer Use Case publique um job na fila sem se preocupar com a execução. O backoff exponencial é fundamental para lidar com falhas transitórias (ex: serviço de e-mail temporariamente indisponível): a primeira retentativa acontece após 1s, a segunda após 2s, a terceira após 4s — dando tempo para o serviço externo se recuperar sem sobrecarregá-lo.
import { Queue } from 'bullmq';
import { redisConnection } from '../redis/RedisClient';
import type { SendWelcomeEmailJobData } from './types';
// Queue: apenas adiciona jobs — não processa
// Pode ser usado em qualquer Use Case
export const emailQueue = new Queue<SendWelcomeEmailJobData>('emails', {
connection: redisConnection,
defaultJobOptions: {
attempts: 3, // Tenta 3x antes de marcar como failed
backoff: {
type: 'exponential', // Aguarda 2^n segundos entre tentativas
delay: 1000, // Delay inicial: 1s, 2s, 4s
},
removeOnComplete: {
age: 24 * 60 * 60, // Remove jobs completed após 24h
count: 1000, // Mantém no máximo 1000 jobs completed
},
removeOnFail: {
age: 7 * 24 * 60 * 60, // Mantém jobs failed por 7 dias para debug
},
},
});
// Uso nos Use Cases:
// await emailQueue.add('send-welcome', { userId, email, name });
// await emailQueue.add('send-password-reset', { ... }, { priority: 1 }); // Alta prioridade
// await emailQueue.add('send-newsletter', { ... }, { delay: 60_000 }); // Delay de 1minWorker: O Processador
O Worker é o único componente que processa jobs. O parâmetro concurrency: 5 significa que ele pode processar até 5 jobs simultaneamente no mesmo processo Node.js — isso funciona bem para tarefas I/O-bound como envio de e-mail. Para tarefas CPU-bound (processamento de imagem, geração de PDF), use concurrency: 1 e escale horizontalmente com múltiplas instâncias do worker. Os eventos completed, failed e stalled são essenciais para observabilidade — integre com o logger para rastrear a saúde das filas.
import { Worker, type Job } from 'bullmq';
import { redisConnection } from '../redis/RedisClient';
import type { SendWelcomeEmailJobData } from '../queues/types';
import { logger } from '../logger';
export function createEmailWorker() {
const worker = new Worker<SendWelcomeEmailJobData>(
'emails',
async (job: Job<SendWelcomeEmailJobData>) => {
logger.info(`Processando job ${job.name}`, {
jobId: job.id,
attempt: job.attemptsMade + 1,
data: { userId: job.data.userId },
});
switch (job.name) {
case 'send-welcome': {
const mailProvider = container.resolve<IMailProvider>('MailProvider');
await mailProvider.sendWelcomeEmail({
to: job.data.email,
name: job.data.name,
});
break;
}
case 'send-password-reset': {
// ...
break;
}
default:
throw new Error(`Job type desconhecido: ${job.name}`);
}
},
{
connection: redisConnection,
concurrency: 5, // Processa até 5 jobs em paralelo por worker
limiter: {
max: 100, // Máximo de 100 jobs por período
duration: 60_000, // Período de 1 minuto (rate limiting)
},
}
);
// Eventos para observabilidade
worker.on('completed', (job) => {
logger.info(`Job ${job.name} concluído`, { jobId: job.id, durationMs: Date.now() - job.timestamp });
});
worker.on('failed', (job, err) => {
logger.error(`Job ${job?.name} falhou (tentativa ${job?.attemptsMade})`, err, {
jobId: job?.id,
});
});
worker.on('stalled', (jobId) => {
logger.warn(`Job ${jobId} ficou stalled — recolocado na fila`);
});
return worker;
}Jobs Recorrentes com Cron
O BullMQ suporta jobs recorrentes via expressões cron, tornando-o uma alternativa ao node-cron com a vantagem de persistência no Redis. Se a aplicação reiniciar, o job agendado não é perdido. O jobId fixo com sufixo -recurring é uma convenção importante: sem ele, cada restart da aplicação criaria um novo job recorrente duplicado no Redis.
// Jobs que rodam automaticamente em um horário específico
await reportQueue.add(
'generate-daily-summary',
{ type: 'daily-summary' },
{
repeat: {
pattern: '0 9 * * *', // Todo dia às 9h (cron expression)
},
jobId: 'daily-summary-recurring', // ID fixo evita duplicatas no restart
}
);
await maintenanceQueue.add(
'cleanup-expired-sessions',
{},
{
repeat: { every: 60 * 60 * 1000 }, // A cada 1 hora
jobId: 'cleanup-sessions-recurring',
}
);Processo Separado para o Worker
Rodar o Worker no mesmo processo da API HTTP é um antipadrão: um bug no processamento de um job pode derrubar a API inteira, e o consumo de memória e CPU dos workers compete com o atendimento das requisições HTTP. O correto é ter dois processos separados — a API (server.ts) e o Worker (worker.ts) — geralmente em dois containers distintos no docker-compose.yml. O Worker implementa graceful shutdown: ao receber SIGTERM, espera os jobs em andamento terminarem antes de fechar, evitando jobs interrompidos no meio do processamento.
import 'reflect-metadata';
import './infra/di/container';
import { createEmailWorker } from './infra/workers/EmailWorker';
import { createReportWorker } from './infra/workers/ReportWorker';
import { logger } from './infra/logger';
// Worker roda em processo separado da API
// Isolamento: um crash no worker não derruba a API
const workers = [
createEmailWorker(),
createReportWorker(),
];
logger.info(`${workers.length} workers iniciados`);
// Graceful shutdown: espera jobs ativos terminarem
async function shutdown() {
logger.info('Encerrando workers...');
await Promise.all(workers.map((w) => w.close()));
process.exit(0);
}
process.on('SIGTERM', shutdown);
process.on('SIGINT', shutdown);
// No docker-compose.yml, adicione um serviço separado:
// worker:
// build: .
// command: node dist/worker.js
// environment: [mesmas env vars da API]
// restart: on-failureO Bull Board (@bull-board/express) adiciona uma UI web para visualizar e monitorar as filas: jobs pendentes, em processamento, concluídos e falhos. É indispensável em produção para debugar jobs presos ou com muitas falhas. Restrinja o acesso com autenticação — nunca exponha publicamente.
Conclusão
BullMQ transforma tarefas lentas em operações assíncronas de baixo impacto na API: a resposta ao usuário é imediata (202), o trabalho pesado acontece em background com retry automático, rate limiting e observabilidade. Separar o processo do Worker do processo da API garante isolamento de falhas — um job travado não afeta a experiência dos outros usuários.