使用队列
队列(Queue)在后端开发中有很多用途,主要目的是提高系统的伸缩性和解决一些性能挑战。
有很多例子可以用到队列:
-
平滑处理特别消耗性能的任务。
比如一些占用性能的任务。可以将这些任务添加到队列中,而不是同步执行它们,再以受控的方式从队列中提取出任务并交给任务消费者(Consumer)处理。
-
拆分可能会阻塞 Node.js 事件循环的庞大任务。例如解码转码这样的 CPU 密集型任务,可以交给队列来处理释放主进程的压力。
-
为不同服务之间提供通信。例如,在一个进程中将任务排入队列,然后再在另外一个进程或服务中消费它们。还可以监听状态事件来得知作业的生命周期,以及完成情况。如果任务失败,还可以重新启动。
在 Nodejs 中,Bull 是可以提供高性能、高可用的队列实现库。Nest 对它做了一层封装,提供了@nestjs/bull 包。
安装和使用
bash1pnpm install --save @nestjs/bull bull
在 AppModule 中配置:
typescript1import { Module } from '@nestjs/common' 2import { BullModule } from '@nestjs/bull' 3 4@Module({ 5 imports: [ 6 BullModule.forRoot({ 7 redis: { 8 host: 'localhost', 9 port: 6379, 10 }, 11 }), 12 ], 13}) 14export class AppModule {}
这一步主要是配置 redis 的信息。
Bull 是用 Redis 保存 job 数据的,所以开启 Redis 是前提条件。
队列任务相对于整个系统而言是一个非常独立的功能,比较好的实现是将它封装成一个动态 Module 来使用。
bash1nest generate module nestBull --no-spec
接着实现动态 module 的 register 方法:
typescript1import { BullModule } from '@nestjs/bull' 2import { DynamicModule, Module, Provider } from '@nestjs/common' 3 4@Module({}) 5export class NestBullModule { 6 static register(queueName: string, consumer: Provider): DynamicModule { 7 const testQueue = BullModule.registerQueue({ 8 name: queueName, 9 }) 10 11 return { 12 module: NestBullModule, 13 imports: [testQueue], 14 providers: [consumer, ...testQueue.providers], 15 exports: [consumer, ...testQueue.exports], 16 } 17 } 18}
register 方法会动态返回整个 Module。
哪个模块想用 Queue 功能,直接调用 register 注册一下,传入queueName 以及 consumer 即可。
注意这个代码:
typescript1 providers: [consumer, ...testQueue.providers], 2 exports: [consumer, ...testQueue.exports],
BullModule.registerQueue方法会返回 providers 属性和 exports 属性,我们来看一下这些属性:

可以看到返回的是动态的 provider 以及它的 token。
把它们从模块导出后就可以被其他模块的 handler 注入了。
至于 consumer,这是一个 class,用来消费 QUEUE 的任务以及实现一些生命周期方法。
这里就做一个最简单的实现:
typescript1import { Processor, Process } from '@nestjs/bull' 2import { Job } from 'bull' 3 4@Processor('test') 5export class TestConsumer { 6 @Process('doSomething') 7 async handle(job: Job) { 8 console.log('Job Start') 9 console.log('data', job.data) 10 } 11}
@Processor(QUEUE_NAME)表示这是一个 Bull 的 Processor,队列名为:test@Process('doSomething')表示这是一个叫 doSomething 的队列处理方法。@Process装饰器是必须的,队列开始时会调用它装饰的方法。doSomething是可缺省的,因为一个队列有可能有多种处理方法,也可能只有一个默认的,这些都要在派发队列任务时指定。
现在我们已经有了一个可以注入其他模块的 BullModule,并且已经写好了Consumer。
接着尝试在 appModule 中注册该队列。
typescript1import { Module } from '@nestjs/common' 2import { AppController } from './app.controller' 3import { AppService } from './app.service' 4import { BullModule } from '@nestjs/bull' 5import { NestBullModule } from './nest-bull/nest-bull.module' 6import { TestConsumer } from './test-consumer/test-consumer' 7 8@Module({ 9 imports: [ 10 BullModule.forRoot({ 11 redis: { 12 host: 'localhost', 13 port: 6379, 14 }, 15 }), 16 NestBullModule.register('test', TestConsumer), 17 ], 18 controllers: [AppController], 19 providers: [AppService], 20}) 21export class AppModule {}
注册时提供队列名以及 Consumer 对象。
接着在 appService 中注入该 Queue 并往队列中加一个任务试试。
typescript1import { InjectQueue } from '@nestjs/bull' 2import { Injectable } from '@nestjs/common' 3import { Queue } from 'bull' 4 5@Injectable() 6export class AppService { 7 constructor( 8 @InjectQueue('test') 9 private readonly testQueue: Queue, 10 ) {} 11 12 async getHello(): Promise<string> { 13 const queue = await this.testQueue.add('doSomething', { foo: 'bar' }) 14 console.log('——————🚀🚀🚀🚀🚀 —— getHello —— queue.id:', queue.id) 15 return 'Hello World!' 16 } 17}
@InjectQueue 装饰器是用来注入我们写好的队列的。
调用 add 方法可以往队列中添加数据。
这里的第一个参数doSomething是指定用什么方法进行消费的,对应 Consumer 中的@Process('doSomething')。
现在打开localhost:3000调试一下:

可以看到队列任务已经被 Consumer 消费了。
同时 add 方法还支持指定延时出列时间、指定固定时间出列等功能:
typescript1await this.testQueue.add( 2 'doSomething', 3 { 4 foo: 'bar', 5 }, 6 { 7 delay: 1000, 8 }, 9)
在 consumer 中指定生命周期可以看到 job 的执行情况:
typescript1import { Processor, Process, OnQueueActive, OnQueueCompleted } from '@nestjs/bull' 2import { Job } from 'bull' 3 4@Processor('test') 5export class TestConsumer { 6 @Process('doSomething') 7 async handle(job: Job) { 8 console.log('Job Start') 9 console.log('data', job.data) 10 } 11 12 // 进行中 13 @OnQueueActive() 14 onActive(job: Job) { 15 console.log(`Processing job ${job.id} of type ${job.name} with data ${job.data}...`) 16 } 17 // 完成 18 @OnQueueCompleted() 19 onCompleted(job: Job) { 20 console.log(`Completed job ${job.id} of type ${job.name} with data ${job.data}...`) 21 } 22}

还有控制 job 的 API:
typescript1await this.testQueue.pause(); 2await this.testQueue.resume(); 3...
这些 API 都可以在Nest 中翻到,基本上涵盖了开发所需。
搭配 Bull-Board 可视化界面
虽然 Bull 提供了 API 可以用于查看 Job 的执行情况,并且可以控制 job 的运行。
但是 Job 一旦多了,在终端中查看每个 Job 显然是一件头疼的事情,即使在使用 winston 等日志工具的情况下。
有没有一款工具能够查看到 Bull 中每个 Job 的运行情况,如果某个 Job 运行失败,我还可以看到失败信息并且手动重新执行它呢?答案是肯定的。
@bull-board/nestjs这个包就是社区提供的优秀界面操作工具。
直接下载:
bash1pnpm install --save @bull-board/nestjs @bull-board/api
再下载适配器:
bash1$ pnpm install --save @bull-board/express 2//or 3$ pnpm install --save @bull-board/fastify
在 AppModule 中引入并注册 BullBoard:
typescript1import { Module } from '@nestjs/common' 2import { AppController } from './app.controller' 3import { AppService } from './app.service' 4import { BullModule } from '@nestjs/bull' 5import { NestBullModule } from './nest-bull/nest-bull.module' 6import { TestConsumer } from './test-consumer/test-consumer' 7import { BullBoardModule } from '@bull-board/nestjs' 8import { ExpressAdapter } from '@bull-board/express' 9 10@Module({ 11 imports: [ 12 BullModule.forRoot({ 13 redis: { 14 host: 'localhost', 15 port: 6379, 16 }, 17 }), 18 BullBoardModule.forRoot({ 19 route: '/bull/queues', 20 adapter: ExpressAdapter, 21 }), 22 NestBullModule.register('test', TestConsumer), 23 ], 24 controllers: [AppController], 25 providers: [AppService], 26}) 27export class AppModule {}
然后在register方法中把board也注册进来。
typescript1import { BullBoardModule } from '@bull-board/nestjs' 2import { BullAdapter } from '@bull-board/api/bullAdapter' 3import { BullModule } from '@nestjs/bull' 4import { DynamicModule, Module, Provider } from '@nestjs/common' 5 6@Module({}) 7export class NestBullModule { 8 static register(queueName: string, consumer: Provider): DynamicModule { 9 const testQueue = BullModule.registerQueue({ 10 name: queueName, 11 }) 12 const testBoard = BullBoardModule.forFeature({ 13 name: queueName, 14 adapter: BullAdapter, 15 }) 16 17 return { 18 module: NestBullModule, 19 imports: [testBoard, testQueue], 20 providers: [consumer, ...testQueue.providers], 21 exports: [consumer, ...testQueue.exports], 22 } 23 } 24}
现在打开http://localhost:3000/bull/queues/

所有 Job 的执行情况都会在该界面中展示出来,是不是特别直观?
如果想在 process 中打 log 查看代码执行情况,可以使用job.log方法。
typescript1 @Process('doSomething') 2 async handle(job: Job) { 3 await job.log('Job Start'); 4 await job.log(`data${job.data}`); 5 return 'success'; 6 }
这些 log 也会在 board 的 Logs 中出现。

是不是非常方便呢?
总结
后端开发中很多占用资源的任务我们可以使用 QUEUE 来处理。
nodejs 中比较有名的实现是 bull.js 这个库。nestjs 对它做了一层封装。
使用 @nestjs/bull 非常简单,在这篇博客中,仅仅只是用它封装了一个独立的 Module。
此外,我们还可以搭配 bull-board 来可视化所有 job 的运行情况。