Arquitetura & EngenhariaBackend

Filas e Processamento Assíncrono

Arquitetura de processamento assíncrono com Bull, Redis, retries e Dead Letter Queue (DLQ).

O processamento assíncrono no ecossistema Urbis (em especial nas APIs do backend e fluxos do Viabiliza) é orquestrado através do Bull sobre instâncias de Redis.

O módulo QueueModule encapsula o registro dinâmico de filas, vínculo de consumers com decoradores NestJS e tratamento avançado de erros com Dead Letter Queue (DLQ).


🧩 Conceitos Fundamentais

  • Fila (Queue): Fila de mensagens baseada em Redis para armazenar jobs pendentes.
  • Consumer: Classe decorada com @Processor responsável por escutar e processar tarefas de uma fila específica.
  • Job: Estrutura contendo o payload e metadados de execução da tarefa.

Convenções de Nomenclatura:

  • Nomes de Fila: Formato kebab-case (ex: email-notification, report-generation).
  • Classes de Consumer: Formato PascalCase (ex: EmailNotificationConsumer, ReportGenerationConsumer).

🛠️ Como Registrar uma Fila no Módulo

Para registrar uma nova fila em um módulo NestJS, utilize o método QueueModule.registerQueue:

import { Module } from '@nestjs/common';
import { QueueModule } from '../queue/queue.module';

@Module({
  imports: [
    QueueModule.registerQueue('notification-queue'),
  ],
})
export class NotificationsModule {}

Registrando Fila com Consumer Vinculado:

import { Module } from '@nestjs/common';
import { QueueModule } from '../queue/queue.module';
import { NotificationConsumer } from './notification.consumer';

@Module({
  imports: [
    QueueModule.registerQueueWithConsumers('notification-queue', [NotificationConsumer]),
  ],
})
export class NotificationsModule {}

⚡ Criando um Consumer

Utilize o decorator @Processor e @Process para definir a rotina de consumo:

import { Processor, Process, OnQueueCompleted, OnQueueFailed } from '@nestjs/bull';
import { Job } from 'bull';

@Processor('notification-queue')
export class NotificationConsumer {
  @Process()
  async handleJob(job: Job<any>) {
    console.log('Processando job de notificação:', job.data);
    // Lógica assíncrona aqui (envio de email, emissão de certidão, etc.)
  }

  @OnQueueCompleted()
  async handleComplete(job: Job) {
    console.log(`Job ${job.id} concluído com sucesso.`);
  }

  @OnQueueFailed()
  async handleFailed(job: Job, error: any) {
    console.error(`Falha no job ${job.id}:`, error.message);
    if (job.attemptsMade >= job.opts.attempts) {
      console.error(`Job ${job.id} esgotou todas as tentativas. Encaminhado para DLQ.`);
    }
  }
}

📨 Adicionando Jobs com Política de Retry e Backoff

No controller ou service, injete a fila com @InjectQueue e passe as opções de política de retentativas:

import { Controller, Post, Body } from '@nestjs/common';
import { InjectQueue } from '@nestjs/bull';
import { Queue } from 'bull';

@Controller('notifications')
export class NotificationsController {
  constructor(@InjectQueue('notification-queue') private readonly notificationQueue: Queue) {}

  @Post('send')
  async addNotificationJob(@Body() payload: any) {
    await this.notificationQueue.add(payload, {
      attempts: 5,        // Número máximo de retentativas
      timeout: 10000,     // Timeout de 10 segundos por execução
      backoff: {
        type: 'exponential',
        delay: 2000,      // Espera 2s, 4s, 8s... entre tentativas
      },
      removeOnComplete: true,
      removeOnFail: false, // Mantém na fila em caso de falha para inspeção
    });

    return { success: true, message: 'Job enfileirado com sucesso!' };
  }
}