如何在@nestjs/bull中移除队列任务?并发任务阻塞问题排查
问题分析与修复方案
核心问题1:send任务挂起导致并发阻塞
你的send处理器声明为async函数,同时还接收了回调参数cb,但全程未调用该回调。在Bull/BullMQ中,async处理器依赖Promise的resolve/reject标记任务完成;若用回调模式,则必须调用cb()结束任务。两种模式混用且未调用回调,会导致任务一直卡在active状态,占用并发槽位,新任务无法被调度,直到任务超时。
核心问题2:receive处理器调用moveToCompleted后的队列停滞
手动调用sendJob.moveToCompleted()完成任务时,需确保任务处于可修改状态(如active),同时BullMQ内部需要时间同步任务状态。手动完成任务后,建议触发队列调度刷新,避免延迟。
修复后的代码
1. 修复send处理器(解决并发阻塞问题)
去掉多余的回调参数,利用async函数的Promise自动标记任务完成:
@Process({ name: 'send', concurrency: 2 }) async sendIso(job: Job<any>) { this.logger.debug(`processing job ${job.id}`); // 消息发送完成后,Promise自动标记任务为completed await this.connectorService.sendMessage(job.data.iso); }
若需手动控制任务状态(如捕获错误标记失败),可添加try/catch:
@Process({ name: 'send', concurrency: 2 }) async sendIso(job: Job<any>) { this.logger.debug(`processing job ${job.id}`); try { await this.connectorService.sendMessage(job.data.iso); await job.moveToCompleted('消息发送成功'); } catch (error) { await job.moveToFailed({ message: error.message }, true); } }
2. 优化receive处理器(解决队列停滞问题)
手动完成任务后,调用队列clean()方法强制触发任务调度:
@Process({ name: 'receive', concurrency: 10 }) async recieve(job: Job<any>) { const unpacked = this.unpackerService.unpack(job.data.message); const id = `${unpacked.FORCEDID}`; const sendJob = await this.sendQueue.getJob(id); this.logger.debug(`received response for ${id}`); if (sendJob) { this.logger.debug(`job ${id} found, moving to complete`); const currentState = await sendJob.getState(); // 仅在任务处于active/waiting状态时执行完成操作 if (['active', 'waiting'].includes(currentState)) { await sendJob.moveToCompleted(JSON.stringify(unpacked)); // 强制刷新队列,触发新任务调度 await this.sendQueue.clean(0, 'active'); this.logger.debug(`job ${id} marked as completed`); } return; } this.logger.debug(`job ${id} not found, saving trace`); }
关键注意事项
- 并发控制生效前提:确保send队列的
concurrency设置正确,且任务不会因未正确结束而占用槽位。 - 任务状态同步:手动修改任务状态后,调用
clean()可加速队列状态同步,避免调度延迟。 - 处理器模式统一:不要混用async/await和回调模式,async处理器依赖Promise状态,回调模式必须调用
cb()。
内容的提问来源于stack exchange,提问作者Carlos Zequeira
相关产品推荐
相关产品推荐

