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

NestJS长任务:SSE进度推送与Redis进度持久化方案咨询

解决方案:基于Redis+Bull队列+SSE的长任务进度追踪方案

针对你的需求,我们可以通过异步任务队列+Redis进度存储+可恢复的SSE推送实现,既保证API不阻塞,又能在用户导航返回后恢复进度查看,同时避免频繁写库的性能问题。

核心思路

  1. 用Bull队列(基于Redis)管理长耗时任务,实现异步执行,避免阻塞API
  2. 任务执行过程中,将进度数据实时写入Redis(高效读写)
  3. SSE接口支持重连恢复:前端重新连接时,先从Redis拉取历史进度,再接收后续实时推送
  4. 任务完成后自动触发后续流程(通过队列事件或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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 21:24:54