You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在@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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 18:05:28