API
JOB
WORKER
BullMQ Queue
retry • delay • async

Redis queue trong NestJS với BullMQ: cách triển khai

Cách dùng Redis làm queue trong NestJS với BullMQ, từ luồng producer-consumer, retry, delay job đến các lưu ý production khi xử lý email, webhook và tác vụ nền.

10 phút đọc15/06/2026

Khi nào nên đưa việc xử lý sang Redis queue

Không phải tác vụ nào cũng nên xử lý ngay trong request HTTP. Những việc như gửi email, tạo thumbnail, đồng bộ webhook, gọi API bên thứ ba hoặc export file thường có 3 vấn đề:

  • Chạy lâu hơn thời gian người dùng nên phải chờ
  • Dễ fail do phụ thuộc service bên ngoài
  • Nếu request timeout thì mất luôn trạng thái xử lý

Đây là lúc queue phát huy tác dụng. API chỉ cần nhận request, validate dữ liệu, rồi đẩy một job vào Redis. Worker sẽ lấy job đó ra xử lý bất đồng bộ ở phía sau.

Sơ đồ Redis queue trong NestJS: API đẩy job vào Redis, worker xử lý và retry khi lỗi

Redis phù hợp với queue ở điểm nào

Redis rất hợp cho queue vì:

  • In-memory nên thao tác push/pop rất nhanh
  • Có data structures phù hợp cho hàng đợi
  • Hỗ trợ delayed jobs, retries, backoff thông qua BullMQ
  • Nhiều worker có thể consume cùng lúc để scale ngang

Điều cần nhớ là Redis queue không thay thế database nghiệp vụ. Redis nên giữ trạng thái tạm thời của job, còn dữ liệu quan trọng vẫn phải nằm ở PostgreSQL, MySQL hoặc storage bền vững khác.

Luồng xử lý điển hình

Ví dụ với tính năng gửi email xác nhận đơn hàng:

  1. User tạo đơn hàng qua API NestJS
  2. API ghi order vào database
  3. API thêm job send-order-email vào queue Redis
  4. Worker lấy job ra, render template rồi gửi email
  5. Nếu SMTP lỗi tạm thời, job sẽ retry theo backoff

Ưu điểm là request tạo đơn chỉ cần nhanh và ổn định. Việc gửi email có thể chậm hơn vài giây nhưng không ảnh hưởng trực tiếp đến trải nghiệm người dùng.

Cài BullMQ trong NestJS

npm install @nestjs/bullmq bullmq ioredis

Khai báo BullModule trong module gốc:

import { Module } from '@nestjs/common';
import { BullModule } from '@nestjs/bullmq';
import { OrdersModule } from './orders/orders.module';

@Module({
  imports: [
    BullModule.forRoot({
      connection: {
        host: process.env.REDIS_HOST || 'localhost',
        port: parseInt(process.env.REDIS_PORT || '6379', 10),
        password: process.env.REDIS_PASSWORD,
        keepAlive: 30000,
        connectTimeout: 10000,
        maxRetriesPerRequest: 3,
      },
      defaultJobOptions: {
        removeOnComplete: 1000,
        removeOnFail: 3000,
        attempts: 5,
        backoff: {
          type: 'exponential',
          delay: 5000,
        },
      },
    }),
    OrdersModule,
  ],
})
export class AppModule {}

Đăng ký queue và thêm job từ service

import { Module } from '@nestjs/common';
import { BullModule } from '@nestjs/bullmq';
import { OrdersService } from './orders.service';
import { OrdersController } from './orders.controller';
import { OrderEmailProcessor } from './order-email.processor';

@Module({
  imports: [
    BullModule.registerQueue({
      name: 'order-email',
    }),
  ],
  controllers: [OrdersController],
  providers: [OrdersService, OrderEmailProcessor],
})
export class OrdersModule {}

Trong service tạo đơn hàng:

import { Injectable } from '@nestjs/common';
import { InjectQueue } from '@nestjs/bullmq';
import { Queue } from 'bullmq';

@Injectable()
export class OrdersService {
  constructor(
    @InjectQueue('order-email')
    private readonly orderEmailQueue: Queue,
  ) {}

  async createOrder(payload: {
    orderId: string;
    customerEmail: string;
    total: number;
  }) {
    // 1. Lưu order vào database ở đây

    // 2. Đẩy job sang queue
    await this.orderEmailQueue.add(
      'send-order-email',
      {
        orderId: payload.orderId,
        customerEmail: payload.customerEmail,
        total: payload.total,
      },
      {
        jobId: `order-email:${payload.orderId}`,
      },
    );

    return {
      success: true,
    };
  }
}

jobId giúp tránh enqueue trùng khi cùng một order bị submit lặp lại.

Worker xử lý job

import { Injectable, Logger } from '@nestjs/common';
import { Processor, WorkerHost } from '@nestjs/bullmq';
import { Job } from 'bullmq';

@Processor('order-email')
@Injectable()
export class OrderEmailProcessor extends WorkerHost {
  private readonly logger = new Logger(OrderEmailProcessor.name);

  async process(
    job: Job<{
      orderId: string;
      customerEmail: string;
      total: number;
    }>,
  ) {
    this.logger.log(`Processing job ${job.name} for order ${job.data.orderId}`);

    if (job.name === 'send-order-email') {
      // Gọi mail provider ở đây
      // throw error nếu provider fail để BullMQ retry
      await fakeSendEmail(job.data.customerEmail, job.data.orderId, job.data.total);
    }

    return {
      delivered: true,
    };
  }
}

async function fakeSendEmail(email: string, orderId: string, total: number) {
  console.log(`Send email to ${email} for order ${orderId}, total=${total}`);
}

Nếu fakeSendEmail throw error, BullMQ sẽ retry theo attemptsbackoff đã cấu hình ở forRoot.

Delay job, retry job và rate limit

Queue không chỉ dùng cho fire-and-forget. Bạn có thể áp dụng thêm:

Delay job

Ví dụ gửi email nhắc thanh toán sau 30 phút:

await this.orderEmailQueue.add(
  'payment-reminder',
  { orderId, customerEmail },
  {
    delay: 30 * 60 * 1000,
  },
);

Retry với backoff

Rất hữu ích khi gọi API bên ngoài có khả năng fail tạm thời như SMTP, webhook hoặc OCR service.

await this.orderEmailQueue.add(
  'sync-webhook',
  { orderId },
  {
    attempts: 6,
    backoff: {
      type: 'exponential',
      delay: 3000,
    },
  },
);

Giới hạn concurrency ở worker

Nếu downstream chỉ chịu được ít request đồng thời, nên giảm concurrency:

import { Processor, WorkerHost } from '@nestjs/bullmq';

@Processor('order-email', {
  concurrency: 5,
})
export class OrderEmailProcessor extends WorkerHost {
  async process(job: Job) {
    // ...
  }
}

Những use case Redis queue làm rất tốt

  • Gửi email, SMS, push notification
  • Xử lý webhook theo kiểu retry-safe
  • Resize ảnh, generate PDF, export CSV
  • Đồng bộ dữ liệu sang CRM, ERP hoặc bên vận chuyển
  • Tách các tác vụ nặng ra khỏi request của API chính

Nếu tác vụ cần throughput cao nhưng độ trễ không phải real-time tuyệt đối, Redis queue thường là lựa chọn thực dụng và triển khai nhanh.

Những chỗ dễ làm sai

1. Ghi dữ liệu nghiệp vụ sau khi enqueue

Không nên add job trước rồi mới insert database. Nếu insert fail nhưng job đã vào queue, worker sẽ xử lý trên dữ liệu chưa tồn tại hoặc không nhất quán.

Thứ tự an toàn hơn là:

  1. Ghi database trước
  2. Commit thành công
  3. Enqueue job sau

2. Không có idempotency

Webhook và queue đều có thể retry. Nếu worker không idempotent, một job chạy lại có thể:

  • Gửi trùng email
  • Trừ kho hai lần
  • Gọi API đối tác hai lần

Hãy dùng jobId, unique key hoặc trạng thái xử lý trong database để chống duplicate.

3. Không dọn job cũ

Nếu để tất cả completed jobs tồn tại mãi, Redis sẽ phình to theo thời gian. removeOnCompleteremoveOnFail là cấu hình gần như bắt buộc.

4. Dùng chung một Redis cho mọi thứ mà không kiểm soát

Nếu cùng một Redis vừa làm cache, vừa làm queue, vừa làm session, cần theo dõi memory và eviction policy rất kỹ. Queue bị eviction là một lỗi rất khó debug.

Lưu ý production

  • Bật keepAlive cho kết nối Redis, nhất là khi dùng ElastiCache
  • Log đủ các event completed, failed, stalled
  • Theo dõi queue depth: số job waiting, active, failed
  • Tách worker thành process riêng nếu load lớn
  • Không coi Redis là nơi lưu trữ bền vững duy nhất cho dữ liệu nghiệp vụ

Nếu queue là critical path, hãy có dashboard hoặc metrics để biết worker đang tắc ở đâu thay vì chờ user báo lỗi.

Kết luận

Redis không chỉ dùng để cache. Trong nhiều hệ thống NestJS, giá trị thực dụng nhất của Redis nằm ở queue: đẩy việc nặng ra khỏi request, retry an toàn khi phụ thuộc bên ngoài lỗi, và scale worker độc lập với API.

Nếu bạn đang có các tác vụ như gửi email, gọi webhook hoặc export file chạy trực tiếp trong controller, đó thường là dấu hiệu nên tách chúng sang Redis queue với BullMQ.