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
相关产品推荐
相关产品推荐

