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

NestJS控制器如何等待队列任务完成后返回结果?

优化NestJS中Bull队列任务结果的获取方案

我在NestJS中实现了一个控制器,用于将任务写入Bull队列。任务处理时长为100毫秒至5000毫秒,我希望在获取任务结果前不与客户端断开连接。
目前我使用了一段性能极差的代码(如下),通过while循环轮询任务状态,该方案会占用大量事件循环资源且存在安全隐患,现寻求更优解决方案:

@Post('job')
async job(@Body() request): Promise<JobRes>{
   const job = await this.bullQueue.add('hande', request);

   while(true){ //UGLY AND BAD
       const state = await job.getState();
       if(state == 'completed'){
          return await this.db.findByJobId(job.id);
       }
       await sleet(1_000);
   }
}

方案1:利用Bull内置的finished方法

Bull提供了finished工具方法,专门用于等待任务完成/失败,完全替代轮询逻辑,性能更优且代码简洁。

import { finished } from 'bull';
import { BadRequestException } from '@nestjs/common';

@Post('job')
async job(@Body() request): Promise<JobRes> {
  const job = await this.bullQueue.add('handle', request);

  try {
    // 设置超时时间(建议略长于任务最长处理时长,比如6000ms)
    await finished(job, { timeout: 6000 });
    return await this.db.findByJobId(job.id);
  } catch (error) {
    // 捕获任务失败、超时等异常
    throw new BadRequestException(`任务处理异常: ${error.message}`);
  }
}

方案2:Server-Sent Events(SSE)实时推送

如果需要给客户端实时反馈任务状态,或者任务时长波动较大,SSE是更合适的选择——客户端建立长连接后,服务器在任务完成时主动推送结果。

控制器代码

import { Observable } from 'rxjs';
import { Sse, Post, Body, Res } from '@nestjs/common';
import { Response } from 'express';

@Post('job')
@Sse('job-result')
async job(@Body() request, @Res() res: Response): Promise<Observable<any>> {
  const job = await this.bullQueue.add('handle', request);
  
  return new Observable((observer) => {
    // 监听任务完成事件
    const onCompleted = async () => {
      const result = await this.db.findByJobId(job.id);
      observer.next({ data: result });
      observer.complete();
    };

    // 监听任务失败事件
    const onFailed = (error: Error) => {
      observer.error({ error: error.message });
      observer.complete();
    };

    job.on('completed', onCompleted);
    job.on('failed', onFailed);

    // 客户端断开连接时清理监听,避免内存泄漏
    res.on('close', () => {
      job.off('completed', onCompleted);
      job.off('failed', onFailed);
      observer.complete();
    });
  });
}

浏览器客户端示例

const eventSource = new EventSource('/job-result');
eventSource.onmessage = (event) => {
  const result = JSON.parse(event.data);
  console.log('任务完成:', result);
  eventSource.close();
};
eventSource.onerror = (error) => {
  console.error('任务处理失败:', error);
  eventSource.close();
};

方案3:Websocket双向通信

如果需要支持客户端主动取消任务等双向交互场景,用Websocket(NestJS集成Socket.IO)是最佳选择。

网关代码

import { WebSocketGateway, SubscribeMessage, WebSocketServer } from '@nestjs/websockets';
import { Socket } from 'socket.io';

@WebSocketGateway()
export class JobGateway {
  constructor(private bullQueue: Queue) {}

  @SubscribeMessage('submitJob')
  async handleJobSubmission(client: Socket, payload: any): Promise<void> {
    const job = await this.bullQueue.add('handle', payload);

    job.on('completed', async () => {
      const result = await this.db.findByJobId(job.id);
      client.emit('jobResult', { success: true, data: result });
    });

    job.on('failed', (error) => {
      client.emit('jobResult', { success: false, error: error.message });
    });

    // 可选:支持客户端取消任务
    client.on('cancelJob', () => {
      job.remove();
      client.emit('jobCancelled');
    });
  }
}

客户端示例

const socket = io();
socket.emit('submitJob', { /* 任务数据 */ });

socket.on('jobResult', (response) => {
  if (response.success) {
    console.log('任务完成:', response.data);
  } else {
    console.error('任务失败:', response.error);
  }
  socket.disconnect();
});

// 可选:取消任务
// socket.emit('cancelJob');
// socket.on('jobCancelled', () => console.log('任务已取消'));

内容的提问来源于stack exchange,提问作者SexyMF

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 20:40:31