Midway 集成 Bull 任务队列:分布式 Job 调度、定时任务与 Redis 实战指南
发布时间:2026/9/29 3:19:49来源:尧图网络
后端微服务云原生【免费下载链接】midway A Node.js Serverless Framework for front-end/full-stack developers. Build the application for next decade. Works on AWS, Alibaba Cloud, Tencent Cloud and traditional VM/Container. Super easy integrate with React and Vue. 项目地址https://gitcode.com/gh_mirrors/mi/midway点击查看免费下载导读本文围绕 Midway 提供的midwayjs/bull组件系统讲解如何在 Midway 应用中引入基于 Redis 的分布式任务队列从安装配置、Processor编写处理器到手动/定时执行任务、任务选项与状态管理、队列管理与组件日志再到 Redis 集群接入、Bull UI 可视化与常见问题排障。读完本文你将掌握在 Midway 中落地可靠的异步/定时任务系统的完整方案。说明从 v3.6.0 开始Midway 原有的midwayjs/task任务调度模块已被废弃任务队列统一收敛到 Bull 组件历史文档见 legacy/task 文档。同时请注意bull 是分布式任务管理系统必须依赖 Redis。为什么选择任务队列队列是一种强大的设计模式可帮助应用应对常见的扩展与性能挑战midwayjs/bull正是 Bull 之上的抽象/封装Bull 是广受欢迎、性能良好的 Node.js 队列系统实现能让 Bull Queues 以非常友好的方式集成进 Midway 应用。队列能解决的核心问题包括平滑处理峰值在任意时间启动资源密集型任务并将它们加入队列而非同步执行让任务进程以受控方式从队列中提取任务也可以轻松添加新的队列消费者来扩展后端处理能力。分解阻塞事件循环的单一任务例如用户请求需要音频转码这类 CPU 密集型工作时把任务委托给其他进程释放面向用户的进程以保持响应。提供跨服务的可靠通信渠道在一个进程/服务中排队任务Job在另一个进程/服务中消费在任何流程或服务的任务生命周期中完成、错误或其他状态变更时通过监听状态事件收到通知当生产者或消费者失败时其状态被保留节点重启后任务处理可自动恢复。Bull 使用 Redis 保存任务数据因此队列架构完全分布式、与平台无关你可以在一个或多个节点/进程中运行部分生产者和消费者在其他节点上运行其余部分。相关信息一览描述可用于标准项目✅可用于 Serverless❌可用于一体化✅包含独立主框架✅包含独立日志✅安装组件$ npm i midwayjs/bull4 --save或者在package.json中增加如下依赖后重新安装{ dependencies: { midwayjs/bull: ^4.0.0 } }使用组件将 bull 组件配置到代码中即在Configuration的imports中引入import { Configuration } from midwayjs/core; import * as bull from midwayjs/bull; Configuration({ imports: [ // ... bull ] }) export class MainConfiguration { //... }从源码结构看组件在 configuration.ts 中注册了bull命名空间与默认配置并注入BullFramework完成队列的加载与管理。一些概念Bull 将整个队列分为三个部分Queue队列管理任务Job任务对象每个任务对象可对任务进行启停控制Processor任务处理实际逻辑执行部分。对应到midwayjs/bull的源码中BullFramework维护了queueMap: Mapstring, BullQueueBullQueue继承自 Bull 原生 Queue 并扩展了addJobToQueue/runJob等方法见 framework.ts。基础配置bull 强依赖 Redis在config.default.ts文件中配置。最简单的默认队列配置// src/config/config.default.ts export default { // ... bull: { // 默认的队列配置 defaultQueueOptions: { redis: redis://127.0.0.1:32768, } }, }有账号密码的情况// src/config/config.default.ts export default { // ... bull: { defaultQueueOptions: { redis: { port: 6379, host: 127.0.0.1, password: foobared, }, } }, }所有的队列都会复用该配置。组件内置的默认值位于 config/config.default.ts包括export const bull { defaultQueueOptions: { prefix: {midway-bull}, defaultJobOptions: { removeOnComplete: 3, removeOnFail: 10, }, }, defaultConcurrency: 1, clearRepeatJobWhenStart: true, };即默认队列 key 前缀为{midway-bull}该写法与 Redis 集群 hash slot 有关见下文「常见问题」任务成功/失败后分别默认保留 3 条/10 条记录默认并发数为 1启动时默认清理旧的重复任务。编写任务处理器使用Processor装饰器装饰一个类即可快速定义一个任务处理器此处不使用 Job避免语义歧义。Processor需要传递一个 Queue队列名字框架启动时若不存在名为该名字的队列会自动创建。例如在src/queue/test.queue.ts中编写// src/queue/test.queue.ts import { Processor, IProcessor } from midwayjs/bull; Processor(test) export class TestProcessor implements IProcessor { async execute() { // ... } }启动时框架会自动查找并初始化上述处理器代码同时自动创建名为test的 Queue。其底层逻辑见 framework.ts 的run()方法遍历BULL_PROCESSOR_KEY上登记的处理器模块调用ensureQueue创建/复用队列再通过addProcessor把execute方法注册为队列处理函数。装饰器本身在 decorator.ts 中实现会把处理器类注册为请求作用域Scope(Request)的 Provider。Processor装饰器还支持第二、三个参数签名如下见 decorator.tsProcessor(queueName: string, jobOptions?: JobOptions, queueOptions?: QueueOptions): ClassDecorator; Processor(queueName: string, concurrency?: number, jobOptions?: JobOptions, queueOptions?: QueueOptions): ClassDecorator;即可以传入并发数、任务默认选项与队列选项当第二个参数不是数字时会被识别为jobOptions此时并发数默认为 1。执行任务定义完 Processor 之后由于并未指定 Processor 如何执行还需要手动触发。通过获取对应的队列即可方便地执行任务。手动执行任务可以在项目启动后执行例如在onServerReady中import { Configuration, Inject } from midwayjs/core; import * as bull from midwayjs/bull; Configuration({ imports: [ // ... bull ] }) export class MainConfiguration { Inject() bullFramework: bull.Framework; //... async onServerReady() { // 获取 Processor 相关的队列 const testQueue this.bullFramework.getQueue(test); // 立即添加这个任务 await testQueue?.addJobToQueue(); } }BullQueue.addJobToQueue在 framework.ts 中的实现为this.add(data || {}, options)即对 Bull 原生add的封装同文件还保留了runJob作为已弃用别名。增加执行参数执行时可以附加默认参数这些参数会以job.data的形式传入executeProcessor(test) export class TestProcessor implements IProcessor { async execute(params) { // params.aaa 1 } } // invoke const testQueue this.bullFramework.getQueue(test); // 立即添加这个任务 await testQueue?.addJobToQueue({ aaa: 1, bbb: 2, });任务状态和管理执行addJobToQueue后可以获取到一个Job对象// invoke const testQueue this.bullFramework.getQueue(test); const job await testQueue?.addJobToQueue();通过这个 Job 对象可以管理进度// 更新进度 await job.progress(60); // 获取进度 const progress await job.process(); // 60获取任务状态const state await job.getState(); // state delayed 延迟状态 // state completed 完成状态在组件测试用例 test/index.test.ts 中可以看到真实的状态流转验证runJob(params, { delay: 1000 })后getState()为delayed等待 1.2s 后变为completed。更多 Job API 可查阅 Bull 官方 REFERENCE 文档。延迟执行执行任务时也有一些额外选项比如延迟 1s 执行const testQueue this.bullFramework.getQueue(test); // 立即添加这个任务 await testQueue?.addJobToQueue({}, { delay: 1000 });中间件和错误处理Bull 组件包含可独立启动的 Framework拥有自己的 App 对象和 Context 结构因此可以对 bull 的 App 配置独立的中间件和错误过滤器Configuration({ imports: [ // ... bull ] }) export class MainConfiguration { App(bull) bullApp: bull.Application; //... async onReady() { this.bullApp.useMiddleare( /*中间件*/); this.bullApp.useFilter( /*过滤器*/); } }上下文任务处理器在请求作用域中执行其有着特殊的 Context 对象结构export interface Context extends IMidwayContext { jobId: JobId; job: Job, from: new (...args) IProcessor; }该接口定义在 interface.ts 中。在addProcessor的实现里框架会调用app.createAnonymousContext({ jobId, job, from })创建匿名上下文因此可以直接从 ctx 中访问当前的 Job 对象// src/queue/test.queue.ts import { Processor, IProcessor, Context } from midwayjs/bull; Processor(test) export class TestProcessor implements IProcessor { Inject() ctx: Context; async execute() { // ctx.jobId xxxx } }更多任务选项除了上面的 delay 之外还有更多的执行选项选项类型描述prioritynumber可选的优先级值。范围从 1最高优先级到 MAX_INT最低优先级。请注意使用优先级对性能有轻微影响因此请谨慎使用。delaynumber等待可以处理此作业的时间量毫秒。请注意为了获得准确的延迟服务器和客户端都应该同步它们的时钟。attemptsnumber在任务完成之前尝试尝试的总次数。repeatRepeatOpts根据 cron 规范的重复任务配置可查看 Bull 文档中 queue.add 的 RepeatOpts 说明以及下文「重复执行的任务」。backoffnumber | BackoffOpts任务失败时自动重试的回退设置。lifoboolean如果为 true则将任务添加到队列的右端而不是左端默认为 false。timeoutnumber任务因超时错误而失败的毫秒数。jobIdnumber | string覆盖任务 id - 默认情况下任务 id 是唯一整数但您可以使用此设置覆盖它。如果您使用此选项则由您来确保 jobId 是唯一的。如果您尝试添加一个 id 已经存在的任务它将不会被添加。removeOnCompleteboolean | number如果为 true则在成功完成后删除任务。如果设置数字则为指定要保留的任务数量。默认行为是任务信息保留在已完成列表中。removeOnFailboolean | number如果为 true则在所有尝试后都失败时删除任务。如果设置数字指定要保留的任务数量。默认行为是将任务信息保留在失败列表中。stackTraceLimitnumber限制将在堆栈跟踪中记录的堆栈跟踪行的数量。重复执行的任务除了手动执行也可以通过Processor装饰器的参数快速配置任务的重复执行import { Processor, IProcessor } from midwayjs/bull; import { FORMAT } from midwayjs/core; Processor(test, { repeat: { cron: FORMAT.CRONTAB.EVERY_PER_5_SECOND } }) export class TestProcessor implements IProcessor { Inject() logger; async execute() { // ... } }框架启动时会读取jobOptions.repeat并自动把重复任务加入队列见 framework.ts 的run()方法中if (options.jobOptions?.repeat) { await this.addJobToQueue(...) }分支。测试夹具 hello.task.ts 正是使用了FORMAT.CRONTAB.EVERY_PER_5_SECOND来验证自动重复执行。常用 Cron 表达式Cron 表达式格式Bull 采用 6 段式含秒* * * * * * ┬ ┬ ┬ ┬ ┬ ┬ │ │ │ │ │ | │ │ │ │ │ └ day of week (0 - 7) (0 or 7 is Sun) │ │ │ │ └───── month (1 - 12) │ │ │ └────────── day of month (1 - 31) │ │ └─────────────── hour (0 - 23) │ └──────────────────── minute (0 - 59) └───────────────────────── second (0 - 59, optional)常见表达式每隔5秒执行一次*/5 * * * * *每隔1分钟执行一次0 */1 * * * *每小时的20分执行一次0 20 * * * *每天 0 点执行一次0 0 0 * * *每天的两点35分执行一次0 35 2 * * *Midway 在框架侧提供了一批常用表达式位于midwayjs/core的FORMAT.CRONTAB定义见 core/src/util/format.tsimport { FORMAT } from midwayjs/core; // 每分钟执行的 cron 表达式 FORMAT.CRONTAB.EVERY_MINUTE内置的其他表达式表达式对应时间CRONTAB.EVERY_SECOND每秒钟CRONTAB.EVERY_MINUTE每分钟CRONTAB.EVERY_HOUR每小时整点CRONTAB.EVERY_DAY每天 0 点CRONTAB.EVERY_DAY_ZERO_FIFTEEN每天 0 点 15 分CRONTAB.EVERY_DAY_ONE_FIFTEEN每天 1 点 15 分CRONTAB.EVERY_PER_5_SECOND每隔 5 秒CRONTAB.EVERY_PER_10_SECOND每隔 10 秒CRONTAB.EVERY_PER_30_SECOND每隔 30 秒CRONTAB.EVERY_PER_5_MINUTE每隔 5 分钟CRONTAB.EVERY_PER_10_MINUTE每隔 10 分钟CRONTAB.EVERY_PER_30_MINUTE每隔 30 分钟高级配置清理之前的任务默认情况下框架会自动清理前一次未调度的重复执行任务保持每次的重复执行任务队列为最新对应配置项clearRepeatJobWhenStart默认值为true。如果某些环境不需要清理可以单独关闭// src/config/config.prod.ts export default { // ... bull: { clearRepeatJobWhenStart: false, }, }注意如果不清理前一次队列若为 10s 执行、现在修改为 20s 执行则两个定时都会存储在 Redis 中导致代码重复执行。日常开发中不清理很容易出现代码重复执行问题但集群部署场景下多台服务器轮流重启可能会导致定时任务被意外清理请评估开关时机。也可以在启动时手动清理所有任务// src/configuration.ts import { Configuration, App, Inject } from midwayjs/core; import * as koa from midwayjs/koa; import { join } from path; import * as bull from midwayjs/bull; Configuration({ imports: [koa, bull], importConfigs: [join(__dirname, ./config)], }) export class MainConfiguration { App() app: koa.Application; Inject() bullFramework: bull.Framework; async onReady() { // 在这个阶段装饰器队列还未创建使用 API 提前手动创建队列装饰器会复用同名队列 const queue this.bullFramework.createQueue(user); // 通过队列手动执行清理 await queue.obliterate({ force: true }); } }清理任务历史记录开启 Redis 后默认情况下 bull 会记录所有成功和失败的任务 key可能导致 Redis 的 key 暴涨可以配置成功或失败后清理的选项。默认情况为成功时保留任务记录 3 条失败保留任务记录 10 条见前文默认配置。也可通过参数调整。例如在装饰器配置import { FORMAT } from midwayjs/core; import { IProcessor, Processor } from midwayjs/bull; Processor(user, { repeat: { cron: FORMAT.CRONTAB.EVERY_MINUTE, }, removeOnComplete: 3, // 成功后移除任务记录最多保留最近 3 条记录 removeOnFail: 10, // 失败后移除任务记录 }) export class UserService implements IProcessor { execute(data: any) { // ... } }也可以在全局 config 中配置// src/config/config.default.ts export default { // ... bull: { defaultQueueOptions: { // 默认的任务配置 defaultJobOptions: { // 保留 10 条记录 removeOnComplete: 10, }, }, }, }从源码看Processor传参中的repeat、delay会被单独取出其余选项作为defaultJobOptions与队列选项合并见 framework.ts 的run()方法。Redis 集群可以使用 bull 提供的createClient方式接入自定义的 Redis 实例从而接入 Redis 集群// src/config/config.default import Redis from ioredis; const clusterOptions { enableReadyCheck: false, // 一定要是false retryDelayOnClusterDown: 300, retryDelayOnFailover: 1000, retryDelayOnTryAgain: 3000, slotsRefreshTimeout: 10000, maxRetriesPerRequest: null // 一定要是null } const redisClientInstance new Redis.Cluster([ { port: 7000, host: 127.0.0.1 }, { port: 7002, host: 127.0.0.1 }, ], clusterOptions); export default { bull: { defaultQueueOptions: { createClient: (type, opts) { return redisClientInstance; }, // 这些任务存储的 key都是相同开头以便区分用户原有 redis 里面的配置 prefix: {midway-bull}, }, } }队列管理队列是廉价的每个 Job 都会绑定一个队列在一些情况下也可以手动对队列进行管理操作。手动创建队列除了使用Processor简单定义队列还可以使用 API 创建import { Configuration, Inject } from midwayjs/core; import * as bull from midwayjs/bull; Configuration({ imports: [ // ... bull ] }) export class MainConfiguration { Inject() bullFramework: bull.Framework; async onReady() { const testQueue this.bullFramework.createQueue(test, { redis: { port: 6379, host: 127.0.0.1, password: foobared, }, prefix: {midway-bull}, }); // ... } }通过createQueue手动创建队列后队列依旧会自动保存写入queueMap。如果在启动时Processor使用了该队列名则会自动使用已经创建好的队列ensureQueue优先复用// 会自动使用上面手动创建的同名队列 Processor(test) export class TestProcessor implements IProcessor { async execute(params) { } }获取队列可以根据队列名获取队列const testQueue bullFramework.getQueue(test);也可以通过装饰器获取import { InjectQueue, BullQueue } from midwayjs/bull; import { Provide } from midwayjs/core; Provide() export class UserService { InjectQueue(test) testQueue: BullQueue; async invoke() { await this.testQueue.pause(); // ... } }InjectQueue装饰器在 decorator.ts 中通过自定义属性装饰器实现并在 configuration.ts 中注册属性处理器直接返回framework.getQueue(queueName)。队列常用操作暂停队列await testQueue.pause();继续队列await testQueue.resume();队列事件// Local events pass the job instance... testQueue.on(progress, function (job, progress) { console.log(Job ${job.id} is ${progress * 100}% ready!); }); testQueue.on(completed, function (job, result) { console.log(Job ${job.id} completed! Result: ${result}); job.remove(); });组件日志组件有着自己的日志默认会将ctx.logger记录在midway-bull.log中。可以单独配置这个 logger 对象export default { midwayLogger: { clients: { // ... bullLogger: { fileLogName: midway-bull.log, }, }, }, }日志的输出格式也可以单独配置export default { bull: { // ... contextLoggerFormat: info { const { jobId, from } info.ctx; return ${info.timestamp} ${info.LEVEL} ${info.pid} [${jobId} ${from.name}] ${info.message}; }, } }组件默认配置中已经内置了上述格式[jobId from.name]形式的上下文格式见 config/config.default.ts。框架内部把frameworkLoggerName设为bullLogger见 framework.ts因此任务执行的日志会统一写入midway-bull.log便于分布式场景下按队列与任务 ID 检索。关于 Redis 版本请尽可能选择最新的版本 5目前在低版本 Redis 上发现有定时任务创建失败的问题。Bull UI在分布式场景中可以利用 Bull UI 简化管理。和 bull 组件类似需要独立安装和启用$ npm i midwayjs/bull-board4 --save或者在package.json中增加如下依赖后重新安装{ dependencies: { midwayjs/bull-board: ^4.0.0 } }将 bull-board 组件配置到代码中import { Configuration } from midwayjs/core; import * as bull from midwayjs/bull; import * as bullBoard from midwayjs/bull-board; Configuration({ imports: [ // ... bull, bullBoard, ] }) export class MainConfiguration { //... }默认的访问路径为http://127.1:7001/ui。可以通过配置修改基础路径// src/config/config.prod.ts export default { // ... bullBoard: { basePath: /ui, }, }此外组件提供了BullBoardManager可以添加动态创建的队列import { Configuration, Inject } from midwayjs/core; import * as bull from midwayjs/bull; import * as bullBoard from midwayjs/bull-board; Configuration({ imports: [ // ... bull, bullBoard ] }) export class MainConfiguration { Inject() bullFramework: bull.Framework; Inject() bullBoardManager: bullBoard.BullBoardManager; async onServerReady() { const testQueue this.bullFramework.createQueue(test, { // ... }); this.bullBoardManager.addQueue(new bullBoard.BullAdapter(testQueue)); } }常见问题1、EVALSHA 错误该问题基本明确会出现在 Redis 的集群版本上。原因是 Redis 会对 key 做 hash 来确定存储的 slot集群下midwayjs/bull的 key 命中了不同的 slot。解决方案将队列的prefix配置用{}包裹强制 Redis 只计算{}内的 hash例如prefix: {midway-bull}这也是组件默认前缀采用{midway-bull}的原因。2、EVAL inside MULTI is not allowed 错误表现为queue.createBulk()、job.moveToFailed()等任务队列 API 调用无效并出现如下错误ReplyError: EXECABORT Transaction discarded because of previous errors. at parseError (project_dir/node_modules/redis-parser/lib/parser.js:179:12) ... previousErrors: [ ReplyError: ERR EVAL inside MULTI is not allowed ... ] }该问题常出现于使用阿里云 Redis 服务。由于这些 API 依赖的 Redis Lua 脚本中使用了 EVAL 或者 EVALSHA阿里云 Redis 使用代理模式连接时会对 Lua 脚本调用做额外限制包括不允许在 MULTI 事务中执行 EVAL 命令其文档提及可通过参数script_check_enable关闭校验但经验证无效。解决方案在阿里云控制台操作开启直连地址将服务切换到直连模式客户端切换成集群模式参考上文「Redis 集群」章节切换配置方式。小结midwayjs/bull以非常低的接入成本把 Bull 的分布式队列能力带入 Midway 应用通过Processor定义处理器、Framework.getQueue/addJobToQueue投递任务、FORMAT.CRONTAB编排定时任务配合默认配置与 Redis 集群接入方案即可构建健壮的异步任务体系。若需要可视化管理可叠加midwayjs/bull-board遇到 Redis 集群或云服务商代理模式带来的 Lua 脚本限制也可参照上文方案逐一排查。赞分享后端微服务云原生【免费下载链接】midway A Node.js Serverless Framework for front-end/full-stack developers. Build the application for next decade. Works on AWS, Alibaba Cloud, Tencent Cloud and traditional VM/Container. Super easy integrate with React and Vue. 项目地址https://gitcode.com/gh_mirrors/mi/midway点击查看免费下载相关推荐gh_mirrors/ji/jira_clone后端任务调度Bull队列与定时任务实现gh_mirrors/ji/jira_clone后端任务调度Bull队列与定时任务实现 你是否还在为项目中的任务调度问题烦恼当用户提交大量请求时系统是否经前端后端企业应用Bull 队列完全指南基于 Redis 的 Node.js 分布式任务与消息队列实战Bull 队列完全指南基于 Redis 的 Node.js 分布式任务与消息队列实战 本文围绕开源仓库 bull一个基于 Redis 的 Node.js 任任务调度后端从定时任务到分布式调度Spring Boot集成XXL-Job实战指南从定时任务到分布式调度Spring Boot集成XXL Job实战指南 在现代应用开发中定时任务和分布式调度是系统稳定性和效率的关键组成部分。Spring示例工程后端教程上一篇Botty 像素机器人上手指南4 小时跑通你的第一条暗黑2重制版自动刷图流水线下一篇QQ空间历史说说完整备份三步上手GetQzonehistory让青春不设保质期创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
网站建设高端定制企业官网