NestJS中调用this.aQueue.add()添加任务无限等待问题求助
问题解决:NestJS Bull队列任务无限等待的修复方案
核心问题分析
- 处理器未返回任务结果:
@Process装饰的方法没有返回值,导致Bull认为任务未完成,不会触发completed事件。 - 函数上下文错误:
setJobData内部使用普通函数,this指向丢失,无法正确访问requestQueue,导致Promise永远无法resolve。 - 不合理的事件监听方式:手动绑定队列
completed事件会重复注册监听,引发内存泄漏,且逻辑冗余。 OnQueueCompleted的误用:该装饰器的方法返回值不会作为任务结果传递,任务结果应直接在@Process方法中返回。
分步修复代码
1. 修复任务处理器(app.processor.ts)
将任务结果直接在@Process方法中返回,移除多余的OnQueueCompleted逻辑:
import { Job } from 'bull'; import { Processor, Process, OnQueueActive } from '@nestjs/bull'; import { ResponseScheme } from './app.schemes'; @Processor('request') export class RequestConsumer { @Process('request') async process_request(job: Job): Promise<ResponseScheme> { console.log(`Job ${job.id} proceed`); // 直接返回任务结果,Bull会自动触发completed事件 const response = new ResponseScheme(); response.answer = job.data.answer; return response; } @OnQueueActive() onActive(job: Job) { console.log(`Data ${JSON.stringify(job.data)} were sent`); } }
2. 简化服务层逻辑(app.service.ts)
使用Bull内置的job.waitUntilFinished()方法替代手动监听事件,同时修复参数类型问题:
import { Injectable } from '@nestjs/common'; import { Job, Queue } from 'bull'; import { InjectQueue } from '@nestjs/bull'; import { plainToClass } from 'class-transformer'; import { RequestScheme, ResponseScheme } from './app.schemes'; @Injectable() export class RequestService { constructor( @InjectQueue('request') private requestQueue: Queue ) {} async sendData(data: RequestScheme): Promise<ResponseScheme> { data = plainToClass(RequestScheme, data); console.log("data in service", data); // 确保delay是数字类型,若data.wait是字符串则转换 const delay = typeof data.wait === 'string' ? parseInt(data.wait, 10) : data.wait; // 添加任务并获取Job实例 const jobInstance = await this.requestQueue.add( 'request', data, { delay: delay || 0 } ); console.log(`Job: ${jobInstance.id}`); // 等待任务完成,直接获取结果 const result = await jobInstance.waitUntilFinished(); return result as ResponseScheme; } }
3. 修复模块配置(app.module.ts)
确保处理器被加入providers数组,否则任务无法被消费:
import { Module } from '@nestjs/common'; import { BullModule } from '@nestjs/bull'; import { RequestController } from './app.controller'; import { RequestService } from './app.service'; import { RequestConsumer } from './app.processor'; @Module({ imports: [ BullModule.forRoot({ redis: { host: 'localhost', port: 6379, maxRetriesPerRequest: null } }), BullModule.registerQueue({ name: 'request' }) ], controllers: [RequestController], providers: [RequestService, RequestConsumer], // 加入RequestConsumer exports: [RequestService] }) export class AppModule {}
额外注意事项
- 确保Redis服务正常运行,否则队列无法连接,任务会一直处于等待状态。
data.wait参数必须是有效数字,若前端传入字符串需做类型转换,避免Bull解析错误。- 避免在请求上下文中等候长时间运行的任务,若任务耗时较长,建议使用Webhook或轮询方式获取结果,防止HTTP请求超时。
内容的提问来源于stack exchange,提问作者Aleksz Maalmen
相关产品推荐
相关产品推荐

