队列任务v1.3.6+

队列是一种强大的设计模式,可帮助您应对常见的应用程序扩展和性能挑战。队列可以帮助您解决的一些问题:

平滑处理峰值。可以在任意时间启动资源密集型任务,然后将这些任务添加到队列中,而不是同步执行。让任务进程以受控方式从队列中提取任务。也可以轻松添加新的队列消费者以扩展后端任务处理。

分解可能会阻塞 Node.js 事件循环的单一任务。比如用户请求需要像音频转码这样的 CPU 密集型工作,就可以将此任务委托给其他进程,从而释放面向用户的进程以保持响应。

提供跨各种服务的可靠通信渠道。例如,您可以在一个进程或服务中排队任务(作业),并在另一个进程或服务中使用它们。在任何流程或服务的作业生命周期中完成、错误或其他状态更改时,您都可以收到通知(通过监听状态事件)。当队列生产者或消费者失败时,它们的状态被保留,并且当节点重新启动时任务处理可以自动重新启动。

本项目基于@midwayjs/bullmq做了对应封装,快速创建队列任务。

快速上手

打开管理后台选择队列任务菜单,点击新增即可快速创建队列任务。

参数说明

  • 任务名称不能重复不能包含空格、$等特殊字符,可以为中文,执行时会根据任务名称更新队列进度
  • 执行策略分为 立即执行、延时执行、定时执行、放弃执行。
    当设置为延时执行时,必须输入延时时间(ms),队列会在到时间后自动执行。
    当设置为 定时执行时,必须输入到秒级的cron表达式。
  • 去重策略 分为可重复、去重、 覆盖,当具有多个进程时,同一队列同时触发如何执行。注意:队列对应bullmq的Queue,指的不是同一个任务。例如:任务类型为sql的任务或自定义任务处理器为test的任务属于同一队列。
    可重复指,同时触发时可重复执行
    去重指第一次声明的任务再次声明时后续的不会执行
    覆盖指同时等待触发时,后续任务的会覆盖。
    更多可参考bullmq文档(https://docs.bullmq.io/guide/jobs/deduplication)
  • 任务类型分为 sql、url请求、自定义,详情参考下方内容任务说明
任务类型:sql

本项目内置了sql执行任务,创建时选择任务类型为 sql,并且输入需要执行的sql语句,即可创建相关队列任务,对应的队列类文件为src/app/admin/processor/sql.processor.ts,对应的 Queuenamesql-task

任务类型:url请求

本项目内置了url请求执行任务,创建时选择任务类型为 url请求,并且输入需要请求的地址方法请求头post请求体GET参数语句,即可创建相关队列任务,对应的队列类文件为src/app/admin/processor/curl.processor.ts,对应的 Queuenamecurl-task

任务类型:自定义

自定义任务,需要在src\app\admin\processor文件夹声明对应的任务类,

  • 编写任务处理器
    使用 @Processor 装饰器装饰一个类,用于快速定义一个任务处理器。

@Processor 装饰器需要传递一个 Queue(队列)的名字,在框架启动时,如果没有名为 test-task 的队列,则会自动创建。

比如,我们在 src/app/admin/processor/test.processor.ts文件中编写如下代码。

复制代码
//src/app/admin/processor/test.processor.ts
import { Processor, IProcessor } from '@midwayjs/bullmq';

@Processor('test-task')
export class TestProcessor implements IProcessor {
  async execute(data: any) {
    // 处理任务逻辑
    console.log('processing job:', data);
  }
}

如果任务处理器需要参数,需要在src/app/admin/processor/default.options.ts中声明

复制代码
//src/app/admin/processor/default.options.ts
import { Processor } from '@midwayjs/bullmq';
export const processorOptions: Record<
  string,
  {
    jobOptions?: Parameters<typeof Processor>[1];
    workerOptions?: Parameters<typeof Processor>[2];
    queueOptions?: Parameters<typeof Processor>[3];
  }
> = {
  //...
  'test-task':{
    queueOptions:{
      defaultJobOptions:{
       attempts:10
      }
    }
  }
};

然后引用

复制代码
//src/app/admin/processor/test.processor.ts
import { Processor, IProcessor } from '@midwayjs/bullmq';
import { processorOptions } from './default.options.js';

@Processor('test-task', processorOptions['test-task'].jobOptions, processorOptions['test-task'].workerOptions, processorOptions['test-task'].queueOptions)
export class TestProcessor implements IProcessor {
  async execute(data: any) {
    // 处理任务逻辑
    console.log('processing job:', data);
  }
}

如果不这样设置,当分别部署网站和worker进程(进程说明可参考下边的部署说明)时,投递任务会获取不到queueOptions

  • 创建任务
    任务类型选择自定义,自定义任务处理器写test-task

一些概念

BullMQ 将整个队列分为以下几个部分:
Queue:队列,管理任务,上边的sql、url请求、test-task都是一个queue
Job:每个任务对象,可以对任务进行启停控制,每次创建的队列任务都是对应的job,任务名称会设置为jobId
Worker:任务处理器,实际的逻辑执行部分
QueueEvents:队列事件,用于监听任务状态变化
FlowProducer:任务流生产者,用于创建任务依赖关系

部署和启动

本项目设置了MODE 环境变量,允许的值为 ALL | ONLY_WORKER | ONLY_API,当设置为ALL时对应进程会加载koa等接口依赖,监听http端口,也会自动创建队列worker处理器;当设置为ONLY_WORKER时对应进程用于执行队列worker会自动创建队列worker处理器,但不会加载koa等接口依赖,不监听http端口;当设置为ONLY_API时对应进程仅加载koa等接口依赖,监听http端口,不会加载processor,不创建队列worker处理器,但可调用jobService投递任务。

开发调试

.env文件设置MODEALL 方便调试

线上部署

线上部署时建议分别部署网站进程 和 worker进程,以防止队列任务和api请求相互影响。

  • 网站进程
    注释掉.env文件的MODE,在进程启动时设置cross-env MODE=ONLY_API,如果是宝塔部署建议做以下设置
  • worker进程
    注释掉.env文件的MODE,在进程启动时设置cross-env MODE=ONLY_WORKER,如果是宝塔部署建议做以下设置

深入

更多说明请参考:midway文档bullmq文档