使用队列

队列(Queue)在后端开发中有很多用途,主要目的是提高系统的伸缩性和解决一些性能挑战。

有很多例子可以用到队列:

  1. 平滑处理特别消耗性能的任务。

    比如一些占用性能的任务。可以将这些任务添加到队列中,而不是同步执行它们,再以受控的方式从队列中提取出任务并交给任务消费者(Consumer)处理。

  2. 拆分可能会阻塞 Node.js 事件循环的庞大任务。例如解码转码这样的 CPU 密集型任务,可以交给队列来处理释放主进程的压力。

  3. 为不同服务之间提供通信。例如,在一个进程中将任务排入队列,然后再在另外一个进程或服务中消费它们。还可以监听状态事件来得知作业的生命周期,以及完成情况。如果任务失败,还可以重新启动。

在 Nodejs 中,Bull 是可以提供高性能、高可用的队列实现库。Nest 对它做了一层封装,提供了@nestjs/bull 包。

安装和使用

bash
1pnpm install --save @nestjs/bull bull

在 AppModule 中配置:

typescript
1import { 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 来使用。

bash
1nest generate module nestBull --no-spec

接着实现动态 module 的 register 方法:

typescript
1import { 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 即可。

注意这个代码:

typescript
1      providers: [consumer, ...testQueue.providers],
2      exports: [consumer, ...testQueue.exports],

BullModule.registerQueue方法会返回 providers 属性和 exports 属性,我们来看一下这些属性:

image-20231202201228614

可以看到返回的是动态的 provider 以及它的 token

把它们从模块导出后就可以被其他模块的 handler 注入了。

至于 consumer,这是一个 class,用来消费 QUEUE 的任务以及实现一些生命周期方法。

这里就做一个最简单的实现:

typescript
1import { 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 中注册该队列。

typescript
1import { 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 并往队列中加一个任务试试。

typescript
1import { 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调试一下:

image-20231202204219801

可以看到队列任务已经被 Consumer 消费了。

同时 add 方法还支持指定延时出列时间、指定固定时间出列等功能:

typescript
1await this.testQueue.add(
2  'doSomething',
3  {
4    foo: 'bar',
5  },
6  {
7    delay: 1000,
8  },
9)

consumer 中指定生命周期可以看到 job 的执行情况:

typescript
1import { 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}

image-20231202205120559

还有控制 job 的 API:

typescript
1await 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这个包就是社区提供的优秀界面操作工具。

直接下载:

bash
1pnpm install --save @bull-board/nestjs @bull-board/api

再下载适配器:

bash
1$ pnpm install --save @bull-board/express
2//or
3$ pnpm install --save @bull-board/fastify

AppModule 中引入并注册 BullBoard

typescript
1import { 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也注册进来。

typescript
1import { 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/

Dec-02-2023 21-14-17

所有 Job 的执行情况都会在该界面中展示出来,是不是特别直观?

如果想在 process 中打 log 查看代码执行情况,可以使用job.log方法。

typescript
1  @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 也会在 boardLogs 中出现。

image-20231202211958589

是不是非常方便呢?

总结

后端开发中很多占用资源的任务我们可以使用 QUEUE 来处理。

nodejs 中比较有名的实现是 bull.js 这个库。nestjs 对它做了一层封装。

使用 @nestjs/bull 非常简单,在这篇博客中,仅仅只是用它封装了一个独立的 Module

此外,我们还可以搭配 bull-board 来可视化所有 job 的运行情况。

代码示例