队列(Bull)
异步任务可用 BullMQ(@nestjs/bullmq)等方案,需 Redis。
官方推荐使用 BullMQ(@nestjs/bullmq)处理基于 Redis 的异步任务。下文示例使用较早的 @nestjs/bull(Bull 3.x),模块名与 API 与 BullMQ 包不同;新项目请直接按 官方 Queues 文档 使用 BullMQ。
核心模型:Queue、Producer、Consumer(Worker)。
模块注册与配置(@nestjs/bull 示例)
使用前注册 BullModule(Bull 3.x):
// 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() 方法向队列中派发任务。
// 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()方法结束时始终调用callback或resolve/return,让Bull明确知晓作业已处理完毕。
️ 消费者 (Consumer):处理作业
消费者通过 @Processor() 和 @Process() 装饰器定义。有两种主要的组织结构:
1. 分类式 (One Processor):轻量选择,共享资源
一个 @Processor 类根据作业名称来消费同一个队列中的多种作业,适合快速上手以及那些需要共享资源的任务。
// 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 类,每个类内部只处理该作业。这种结构让关注点分离,适合业务复杂的情况。
// 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 的行为,例如控制并发数或限流。
// 并发5个任务,每秒最多处理10个任务
@Processor('email', { concurrency: 5, limiter: { max: 10, duration: 1000 } })
export class EmailProcessor { ... }
concurrency: 决定同一时间该处理器可并行处理的任务数。注意,即使在同一@Processor内,也可为不同@Process方法设置不同的并发数,实现内部差异化。
🎤 事件监听 (Listeners):追踪作业全生命周期
Bull 提供了完整的事件系统,用于监控作业从入队到完成(或失败)的整个生命周期。
// 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),可使用OnGlobalQueueCompleted、OnGlobalQueueFailed等带Global字样的装饰器。这对跨队列的统一监控非常有用。
⚡ 高级特性一览:延迟、重试、移除与重复
- 延迟与优先级:在生产者添加作业时,通过
delay和priority选项可轻松实现延时任务和优先级排序,这在功能说明部分已收到过。 - 移除作业:若需动态管理队列(如用户取消预约),可通过
queue.getJob(jobId)获取作业对象并调用job.remove()。 - 重复作业 (Repeatable Jobs):Bull 支持按 Cron 表达式或特定时间间隔创建重复执行的作业。
- 从生产端使用
queue.add()时,在Job Options中增加repeat字段,例如{ cron: '0 9 * * *', tz: 'Asia/Shanghai' }。 - 从消费端来看,某队列可在多个实例中仅由一个实例的
@Process进行处理,这依赖于 Redis 的分布式锁机制,有效避免多个消费者同时消费同一个重复作业,从而达到分布式调度的效果。
- 从生产端使用
✨ 最佳实践总结
- 命名策略:建议使用
kebab-case为队列命名。作业名称与处理函数名称保持一致,利于全局搜索和维护。 JobOptions选择原则:比较关键的配置包括priority(优先级)、delay(延迟)、attempts(重试次数) 与backoff(退避策略)。- 数据大小限制:Redis 对所有键值的大小和数量均有实际限制。务必不要将大块数据 (如 base64 图片) 直接放进
job.data,应改为存入对象存储(如 OSS)的元数据与文件路径。设计优雅的Job结构能有效降低对 Redis 的存储压力。 - 监控与告警:务必为关键业务队列实现
@OnQueueFailed或者附加全局监听,以便及时发现处理失败的任务。
参考文献
以下链接在编写时均可正常访问:
| 资料 | 说明 |
|---|---|
| NestJS 文档 | 官方 |
| Queues | 本章主题 |
| Request lifecycle | 执行顺序 |
相关文章
认证与授权
认证(Authentication):确认「你是谁」(如 JWT、Session)。 - 授权(Authorization):确认「你能做什么」(如 RBAC、策略检查)。
缓存
@nestjs/cache-manager 统一缓存 API,存储实现可插拔。
提供者与服务
Service 是最常见的 Provider,封装业务逻辑。见 module。
守卫 (Guards) 与授权
守卫决定是否放行请求,常用于认证与授权。见 auth。
配置管理(Config 模块)
@nestjs/config 加载 .env 并提供 ConfigService。见 工程化 env。
数据库集成(以 TypeORM 为例)
Nest 通过 @nestjs/typeorm 等包集成 ORM;生产环境用 migration,慎用 synchronize。亦可选用 Prisma、MikroORM 等(见 官方 Database)。
Series
new
17 / 19