NestJS长任务:SSE进度推送与Redis进度持久化方案咨询
解决方案:基于Redis+Bull队列+SSE的长任务进度追踪方案
针对你的需求,我们可以通过异步任务队列+Redis进度存储+可恢复的SSE推送实现,既保证API不阻塞,又能在用户导航返回后恢复进度查看,同时避免频繁写库的性能问题。
核心思路
- 用Bull队列(基于Redis)管理长耗时任务,实现异步执行,避免阻塞API
- 任务执行过程中,将进度数据实时写入Redis(高效读写)
- SSE接口支持重连恢复:前端重新连接时,先从Redis拉取历史进度,再接收后续实时推送
- 任务完成后自动触发后续流程(通过队列事件或Redis监听)
具体实现步骤
1. 依赖安装
先安装必要的NestJS包:
npm install @nestjs/bull bull @nestjs/redis ioredis uuid
2. 配置Redis与任务队列
在app.module.ts中配置Redis和Bull模块:
import { Module } from '@nestjs/common'; import { BullModule } from '@nestjs/bull'; import { RedisModule } from '@nestjs/redis'; @Module({ imports: [ RedisModule.forRoot({ url: 'redis://localhost:6379', }), BullModule.forRoot({ redis: { host: 'localhost', port: 6379, }, }), // 注册长任务队列 BullModule.registerQueue({ name: 'long-task-queue', }), ], }) export class AppModule {}
3. 创建任务处理器
编写长任务的处理器,执行任务并更新Redis进度:
import { Processor, Process } from '@nestjs/bull'; import { Job } from 'bull'; import { RedisService } from '@nestjs/redis'; @Processor('long-task-queue') export class LongTaskProcessor { constructor(private readonly redisService: RedisService) {} @Process('document-processing') async handleDocumentProcessing(job: Job) { const { taskId, documentId } = job.data; const redisClient = this.redisService.getClient(); const progressKey = `task:progress:${taskId}`; // 初始化进度 await redisClient.hSet(progressKey, { status: 'running', currentStep: 0, totalSteps: 3, steps: JSON.stringify([ { name: '生成文档摘要', completed: false }, { name: '执行相关处理', completed: false }, { name: '生成最终结果', completed: false }, ]), updateTime: Date.now(), }); // 步骤1:生成文档摘要 await this.simulateLongOperation(5000); await redisClient.hSet(progressKey, { currentStep: 1, steps: JSON.stringify([ { name: '生成文档摘要', completed: true }, { name: '执行相关处理', completed: false }, { name: '生成最终结果', completed: false }, ]), updateTime: Date.now(), }); // 步骤2:执行相关处理 await this.simulateLongOperation(8000); await redisClient.hSet(progressKey, { currentStep: 2, steps: JSON.stringify([ { name: '生成文档摘要', completed: true }, { name: '执行相关处理', completed: true }, { name: '生成最终结果', completed: false }, ]), updateTime: Date.now(), }); // 步骤3:生成最终结果 await this.simulateLongOperation(6000); await redisClient.hSet(progressKey, { status: 'completed', currentStep: 3, steps: JSON.stringify([ { name: '生成文档摘要', completed: true }, { name: '执行相关处理', completed: true }, { name: '生成最终结果', completed: true }, ]), updateTime: Date.now(), }); // 任务完成后自动推进流程 await this.triggerPostProcess(taskId); // 设置进度数据过期(7天) await redisClient.expire(progressKey, 60 * 60 * 24 * 7); } private async simulateLongOperation(delay: number) { return new Promise(resolve => setTimeout(resolve, delay)); } private async triggerPostProcess(taskId: string) { // 实现任务完成后的逻辑,比如更新数据库、发送通知等 console.log(`任务${taskId}完成,触发后续流程`); } }
4. 编写触发任务的API
提供POST接口启动任务,立即返回taskId,不阻塞:
import { Controller, Post, Body } from '@nestjs/common'; import { Queue } from 'bull'; import { InjectQueue } from '@nestjs/bull'; import { v4 as uuidv4 } from 'uuid'; @Controller('tasks') export class TasksController { constructor(@InjectQueue('long-task-queue') private readonly longTaskQueue: Queue) {} @Post('start-document-processing') async startDocumentProcessing(@Body() body: { documentId: string }) { const taskId = uuidv4(); await this.longTaskQueue.add('document-processing', { taskId, documentId: body.documentId, }); return { taskId, message: '任务已启动' }; } }
5. 实现可恢复的SSE进度推送接口
SSE接口在前端连接时,先从Redis拉取当前进度,再监听Redis键变化实时推送更新:
import { Controller, Sse, Param, Get } from '@nestjs/common'; import { Observable, interval, from, merge } from 'rxjs'; import { map, switchMap } from 'rxjs/operators'; import { RedisService } from '@nestjs/redis'; @Controller('tasks') export class TaskProgressController { constructor(private readonly redisService: RedisService) {} @Sse(':taskId/progress') async getProgressStream(@Param('taskId') taskId: string): Promise<Observable<any>> { const redisClient = this.redisService.getClient(); const progressKey = `task:progress:${taskId}`; // 1. 推送初始进度 const initialProgress = await redisClient.hGetAll(progressKey); if (initialProgress) { initialProgress.steps = JSON.parse(initialProgress.steps); } // 2. 监听Redis键变化,实时拉取进度 await redisClient.configSet('notify-keyspace-events', 'Kh'); const progressUpdates = from(redisClient.subscribe(`__keyspace@0__:${progressKey}`)).pipe( switchMap(() => interval(1000)), switchMap(async () => { const progress = await redisClient.hGetAll(progressKey); progress.steps = JSON.parse(progress.steps); return progress; }), map(progress => ({ data: progress })), ); // 合并初始进度和后续更新 return merge( from([{ data: initialProgress }]), progressUpdates, ); } // 供前端拉取历史进度的GET接口 @Get(':taskId/progress') async getCurrentProgress(@Param('taskId') taskId: string) { const redisClient = this.redisService.getClient(); const progress = await redisClient.hGetAll(`task:progress:${taskId}`); if (progress) { progress.steps = JSON.parse(progress.steps); } return progress; } }
6. 前端适配逻辑
前端处理SSE重连,页面加载时先拉取历史进度:
let sse; const taskId = '用户的taskId'; // 初始化进度 async function initProgress() { const response = await fetch(`/tasks/${taskId}/progress`); const progress = await response.json(); renderProgress(progress); // 启动SSE连接 connectSSE(); } // 连接SSE function connectSSE() { sse = new EventSource(`/tasks/${taskId}/progress`); sse.onmessage = (event) => { const progress = JSON.parse(event.data); renderProgress(progress); // 任务完成后关闭SSE并推进流程 if (progress.status === 'completed') { sse.close(); window.location.href = `/tasks/${taskId}/result`; } }; sse.onerror = () => { console.log('SSE断开,尝试重连...'); sse.close(); setTimeout(connectSSE, 3000); }; } // 渲染进度UI function renderProgress(progress) { console.log('当前进度:', progress); // 这里实现步骤列表、进度条等UI渲染逻辑 } // 页面加载时初始化 window.onload = initProgress;
关键细节说明
- Redis进度存储:用哈希结构存储任务进度,字段清晰,读写高效;任务完成后设置过期时间避免内存浪费
- SSE重连恢复:前端断开后自动重连,重连时先拉取最新进度,再接收实时更新,保证用户返回后能看到完整进度
- 任务队列管理:Bull队列自带任务重试、失败处理、状态追踪功能,适合长耗时任务的生命周期管理
- 性能优化:Redis读写性能远高于关系型数据库,避免频繁写库的性能瓶颈;SSE单向推送比轮询更高效
内容的提问来源于stack exchange,提问作者AnAverageDev
相关产品推荐
相关产品推荐

