FORMA

队列(Bull)

异步任务可用 BullMQ@nestjs/bullmq)等方案,需 Redis。

官方推荐使用 BullMQ@nestjs/bullmq)处理基于 Redis 的异步任务。下文示例使用较早的 @nestjs/bull(Bull 3.x),模块名与 API 与 BullMQ 包不同;新项目请直接按 官方 Queues 文档 使用 BullMQ

核心模型:QueueProducerConsumer(Worker)。

模块注册与配置(@nestjs/bull 示例)

使用前注册 BullModule(Bull 3.x):

typescript
// app.module.ts
import { Module } from "@nestjs/common";
import { BullModule } from "@nestjs/bull";

@Module({
  imports: [
    BullModule.forRootAsync({
      useFactory: (configService: ConfigService) => ({
        redis: {
          host: configService.get("REDIS_HOST", "localhost"),
          port: configService.get<number>("REDIS_PORT", 6379),
        },
        prefix: "my-queue", // 自定义队列键名前缀
        defaultJobOptions: {
          // 全局默认作业选项
          attempts: 3, // 失败时重试3次
          removeOnComplete: true, // 完成后自动移除
        },
      }),
      inject: [ConfigService],
    }),
    BullModule.registerQueueAsync({
      name: "email", // 队列名称,用于后续注入
    }),
  ],
})
export class AppModule {}
  • 使用 registerQueueAsync 能按需创建队列。name 字段是队列的唯一标识,在后续的生产和消费环节都需要用到。

📤 生产者 (Producer):派发作业

在服务中通过 @InjectQueue() 注入特定的队列实例,并使用 add() 方法向队列中派发任务。

typescript
// email.service.ts
import { Injectable } from "@nestjs/common";
import { InjectQueue } from "@nestjs/bull";
import { Queue } from "bull";

export interface EmailJobData {
  to: string;
  subject: string;
  template: string;
}

@Injectable()
export class EmailService {
  constructor(@InjectQueue("email") private emailQueue: Queue) {}

  async sendWelcomeEmail(to: string) {
    // Job 名称是 send-welcome,用于区分不同类型的任务
    await this.emailQueue.add("send-welcome", {
      to,
      subject: "Welcome!",
      template: "welcome",
    } as EmailJobData);
  }

  async scheduleBirthdayEmail(to: string, date: Date) {
    // 延迟执行:在设定时间之后才被消费
    const delay = date.getTime() - Date.now();
    await this.emailQueue.add("send-birthday", { to }, { delay });
  }

  async sendHighPriorityNotification(to: string) {
    await this.emailQueue.add("send-notify", { to }, { priority: 1 }); // 数字越小优先级越高
  }
}

add() 方法接受三个参数,核心是作业名称作业数据Job Options 则提供了丰富的控制选项:

  • delay: 延迟执行,单位是毫秒(ms),常用于实现定时或延时任务。
  • priority: 设置作业优先级,队列会优先处理高优先级任务。
  • jobId: 自定义作业ID,可通过此ID对作业进行后续操作。
  • attempts: 任务失败时的重试次数,这是保障队列可靠性的关键之一。
  • backoff: 重试时的间隔策略,例如每次重试的等待时间可以递增。
  • removeOnComplete / removeOnFail: 作业完成或失败后是否自动从Redis中移除,有助于管理内存。

关于作业的重要说明:作业本身(及其承载的数据)会被序列化后存储在 Redis 中。有时即便作业被成功处理,它仍可能因“等待回调”而在视觉上残留或延迟完成,此时有效的解决办法是确保在 @Process() 方法结束时始终调用 callbackresolve/return,让Bull明确知晓作业已处理完毕。

️ 消费者 (Consumer):处理作业

消费者通过 @Processor()@Process() 装饰器定义。有两种主要的组织结构:

1. 分类式 (One Processor):轻量选择,共享资源

一个 @Processor 类根据作业名称来消费同一个队列中的多种作业,适合快速上手以及那些需要共享资源的任务。

typescript
// email.processor.ts
import { Processor, Process } from "@nestjs/bull";
import { Job } from "bull";

@Processor("email")
export class EmailProcessor {
  @Process("send-welcome")
  async handleWelcome(job: Job<EmailJobData>) {
    const { to, subject, template } = job.data;
    // 具体发送逻辑
    await this.sendEmail(to, subject, template);
    // 可选: 更新进度
    await job.progress(100);
  }

  @Process("send-birthday")
  async handleBirthday(job: Job<EmailJobData>) {
    // 成员方法可使用共享属性
  }
}

2. 分离式 (Multiple Consumers):资源隔离,专人专事

每种作业可以各自拥有独立的 @Processor 类,每个类内部只处理该作业。这种结构让关注点分离,适合业务复杂的情况。

typescript
// welcome.processor.ts
@Processor("email")
export class WelcomeEmailProcessor {
  @Process("send-welcome")
  async handle(job: Job) {
    /* ... */
  }
}

// birthday.processor.ts
@Processor("email")
export class BirthdayEmailProcessor {
  @Process("send-birthday")
  async handle(job: Job) {
    /* ... */
  }
}

⚙️ Worker 配置优化

可以通过 @Processor 的第二个参数来配置底层 Worker 的行为,例如控制并发数或限流。

typescript
// 并发5个任务,每秒最多处理10个任务
@Processor('email', { concurrency: 5, limiter: { max: 10, duration: 1000 } })
export class EmailProcessor { ... }
  • concurrency: 决定同一时间该处理器可并行处理的任务数。注意,即使在同一 @Processor 内,也可为不同 @Process 方法设置不同的并发数,实现内部差异化。

🎤 事件监听 (Listeners):追踪作业全生命周期

Bull 提供了完整的事件系统,用于监控作业从入队到完成(或失败)的整个生命周期。

typescript
// email.listener.ts
import { Processor, OnQueueCompleted, OnQueueFailed, OnQueueProgress } from "@nestjs/bull";

@Processor("email")
export class EmailListener {
  @OnQueueCompleted()
  onCompleted(job: Job, result: any) {
    console.log(`Job ${job.id} completed with result: ${result}`);
  }

  @OnQueueFailed()
  onFailed(job: Job, err: Error) {
    console.error(`Job ${job.id} failed: ${err.message}`);
    // 可在此集成 Sentry 等错误追踪系统
  }

  @OnQueueProgress()
  onProgress(job: Job, progress: number) {
    console.log(`Job ${job.id} progress: ${progress}`);
  }
}

每个 @Process 中调用 job.progress() 时,@OnQueueProgress 对应的事件就会触发,便于追踪进度。

如果需要对某个事件的“全局”版本进行监听(不限于特定 @Processor),可使用 OnGlobalQueueCompletedOnGlobalQueueFailed 等带 Global 字样的装饰器。这对跨队列的统一监控非常有用。

⚡ 高级特性一览:延迟、重试、移除与重复

  • 延迟与优先级:在生产者添加作业时,通过 delaypriority 选项可轻松实现延时任务和优先级排序,这在功能说明部分已收到过。
  • 移除作业:若需动态管理队列(如用户取消预约),可通过 queue.getJob(jobId) 获取作业对象并调用 job.remove()
  • 重复作业 (Repeatable Jobs):Bull 支持按 Cron 表达式或特定时间间隔创建重复执行的作业。
    1. 从生产端使用 queue.add() 时,在 Job Options 中增加 repeat 字段,例如 { cron: '0 9 * * *', tz: 'Asia/Shanghai' }
    2. 从消费端来看,某队列可在多个实例中仅由一个实例的 @Process 进行处理,这依赖于 Redis 的分布式锁机制,有效避免多个消费者同时消费同一个重复作业,从而达到分布式调度的效果。

✨ 最佳实践总结

  • 命名策略:建议使用 kebab-case 为队列命名。作业名称与处理函数名称保持一致,利于全局搜索和维护。
  • JobOptions 选择原则:比较关键的配置包括 priority (优先级)、delay (延迟)、attempts (重试次数) 与 backoff (退避策略)。
  • 数据大小限制:Redis 对所有键值的大小和数量均有实际限制。务必不要将大块数据 (如 base64 图片) 直接放进 job.data,应改为存入对象存储(如 OSS)的元数据与文件路径。设计优雅的 Job 结构能有效降低对 Redis 的存储压力。
  • 监控与告警:务必为关键业务队列实现 @OnQueueFailed 或者附加全局监听,以便及时发现处理失败的任务。

参考文献

以下链接在编写时均可正常访问:

资料说明
NestJS 文档官方
Queues本章主题
Request lifecycle执行顺序