AI 探索技术手记

AIGC任务队列与WebSocket实时状态更新系统

2026年7月23日◷ 5 分钟阅读
AIGC任务队列与WebSocket实时状态更新系统

概述

AI 图片生成、视频生成、批量文件处理等耗时任务的典型痛点:用户提交后不知道进度,刷新页面任务丢失,多个任务同时跑导致服务崩溃。本文将构建一套 前端任务队列 + WebSocket 实时状态推送 系统,涵盖:任务状态机、并发控制、进度条实时更新、localStorage 持久化、失败重试/取消、以及与大文件分片上传的联动。

1. 系统架构总览

用户提交任务 → 进入前端队列(并发锁控制)→ 调用后端 API
                                                  ↓
                                            WebSocket 实时推送
                                                  ↓
                                        前端更新任务状态 & 进度
                                                  ↓
                              completed → 展示结果 / failed → 重试

2. 任务状态机与数据结构

// ===== task.types.ts =====

type TaskStatus = 'pending' | 'running' | 'progress' | 'completed' | 'failed' | 'canceled';
type TaskType = 'image_gen' | 'video_gen' | 'file_process' | 'batch';

interface Task {
  id: string;
  type: TaskType;
  status: TaskStatus;
  progress: number;      // 0-100
  message: string;        // 状态文字
  result?: string;        // 图片/视频 URL
  error?: string;
  params: Record<string, any>;
  createdAt: number;
}

// 状态流转
const STATUS_FLOW: Record<TaskStatus, TaskStatus[]> = {
  pending:    ['running', 'canceled'],
  running:    ['progress', 'failed', 'canceled'],
  progress:   ['completed', 'failed', 'canceled'],
  completed:  [],
  failed:     ['pending'],    // 重试 → 回到排队
  canceled:   [],
};

3. 任务队列管理器(单例)

// ===== TaskQueueManager.ts =====
import { v4 as uuid } from 'uuid';

type Listener = (tasks: Task[]) => void;

class TaskQueueManager {
  private tasks: Task[] = [];
  private listeners: Set<Listener> = new Set();
  private maxConcurrency: number = 2;
  private runningCount: number = 0;

  static STORAGE_KEY = 'aigc-tasks';

  constructor() {
    this.restoreFromLocalStorage();
  }

  // ---- 订阅 ----
  subscribe(fn: Listener) {
    this.listeners.add(fn);
    fn([...this.tasks]);
    return () => this.listeners.delete(fn);
  }

  private notify() {
    const snapshot = [...this.tasks];
    this.listeners.forEach((fn) => fn(snapshot));
    this.persist();
  }

  // ---- 创建任务 ----
  addTask(type: TaskType, params: Record<string, any>): Task {
    const task: Task = {
      id: uuid(),
      type,
      status: 'pending',
      progress: 0,
      message: '排队中...',
      params,
      createdAt: Date.now(),
    };
    this.tasks = [task, ...this.tasks];
    this.notify();
    this.tryProcessQueue();
    return task;
  }

  // ---- 并发控制 ----
  private async tryProcessQueue() {
    while (this.runningCount < this.maxConcurrency) {
      const pending = this.tasks.find((t) => t.status === 'pending');
      if (!pending) break;
      this.runningCount++;
      this.updateTask(pending.id, { status: 'running', message: '执行中...' });
      this.executeTask(pending.id).finally(() => {
        this.runningCount--;
        this.tryProcessQueue();
      });
    }
  }

  // ---- 执行任务(模拟) ----
  private async executeTask(taskId: string) {
    const task = this.tasks.find((t) => t.id === taskId);
    if (!task || task.status === 'canceled') return;

    try {
      // 模拟:通过 WebSocket 或轮询接收进度
      // 真实项目中这里会调用 fetch 提交任务,然后监听 WebSocket 消息
      for (let p = 10; p <= 100; p += 10) {
        await this.sleep(500);
        // 检查是否被取消
        const current = this.tasks.find((t) => t.id === taskId);
        if (!current || current.status === 'canceled') return;

        this.updateTask(taskId, {
          status: 'progress',
          progress: p,
          message: `生成中 ${p}%...`,
        });
      }
      this.updateTask(taskId, {
        status: 'completed',
        progress: 100,
        message: '生成完成!',
        result: `https://picsum.photos/seed/${taskId}/800/600`,
      });
    } catch (err: any) {
      this.updateTask(taskId, {
        status: 'failed',
        message: `失败:${err.message}`,
        error: err.message,
      });
    }
  }

  // ---- 更新任务(统一入口) ----
  updateTask(taskId: string, patch: Partial<Task>) {
    this.tasks = this.tasks.map((t) =>
      t.id === taskId ? { ...t, ...patch } : t
    );
    this.notify();
  }

  // ---- 重试 ----
  retryTask(taskId: string) {
    this.updateTask(taskId, { status: 'pending', progress: 0, error: '', message: '重新排队...' });
    this.tryProcessQueue();
  }

  // ---- 取消 ----
  cancelTask(taskId: string) {
    this.updateTask(taskId, { status: 'canceled', message: '已取消' });
    // 如果正在执行,这个标记会被 executeTask 轮询捕获
  }

  // ---- 清除已完成/已取消的任务 ----
  clearDone() {
    this.tasks = this.tasks.filter(
      (t) => t.status !== 'completed' && t.status !== 'canceled'
    );
    this.notify();
  }

  // ---- localStorage 持久化 ----
  private persist() {
    const savable = this.tasks.map((t) => {
      // 不持久化运行时变量
      const { ...rest } = t;
      // 排队/执行中任务:刷新页后重置为 pending
      if (rest.status === 'running' || rest.status === 'progress') {
        rest.status = 'pending';
        rest.progress = 0;
        rest.message = '排队中(页面刷新后恢复)';
      }
      return rest;
    });
    localStorage.setItem(TaskQueueManager.STORAGE_KEY, JSON.stringify(savable));
  }

  private restoreFromLocalStorage() {
    try {
      const raw = localStorage.getItem(TaskQueueManager.STORAGE_KEY);
      if (raw) this.tasks = JSON.parse(raw);
    } catch { /* ignore */ }
  }

  private sleep(ms: number) {
    return new Promise((r) => setTimeout(r, ms));
  }

  // 单例
  static instance = new TaskQueueManager();
}

export const taskQueue = TaskQueueManager.instance;

4. WebSocket 实时进度监听

// ===== useTaskWebSocket.ts =====
import { useEffect, useRef, useCallback } from 'react';
import { taskQueue } from './TaskQueueManager';

interface WsProgressMsg {
  type: 'progress';
  task_id: string;
  progress: number;   // 0-100
  message: string;
  result?: string;
}

interface WsErrorMsg {
  type: 'error';
  task_id: string;
  message: string;
}

type WsMessage = WsProgressMsg | WsErrorMsg;

export function useTaskWebSocket(url: string = '/ws/tasks') {
  const wsRef = useRef<WebSocket | null>(null);
  const reconnectTimerRef = useRef<number>();

  const connect = useCallback(() => {
    const ws = new WebSocket(url);
    wsRef.current = ws;

    ws.onopen = () => console.log('[WS] 任务通道已连接');

    ws.onmessage = (event) => {
      try {
        const data: WsMessage = JSON.parse(event.data);

        if (data.type === 'progress') {
          taskQueue.updateTask(data.task_id, {
            status: data.result ? 'completed' : 'progress',
            progress: data.progress,
            message: data.message,
            result: data.result,
          });
        }

        if (data.type === 'error') {
          taskQueue.updateTask(data.task_id, {
            status: 'failed',
            message: data.message,
            error: data.message,
          });
        }
      } catch (e) {
        console.error('[WS] 消息解析错误', e);
      }
    };

    ws.onclose = () => {
      console.log('[WS] 断开,3 秒后重连');
      reconnectTimerRef.current = window.setTimeout(connect, 3000);
    };

    ws.onerror = () => ws.close();
  }, [url]);

  useEffect(() => {
    connect();
    return () => {
      wsRef.current?.close();
      clearTimeout(reconnectTimerRef.current);
    };
  }, [connect]);

  return { sendMessage: (data: any) => wsRef.current?.send(JSON.stringify(data)) };
}

5. 任务 UI 组件

// ===== TaskQueuePanel.tsx =====
import { useEffect, useState } from 'react';
import { taskQueue, Task } from './TaskQueueManager';

export default function TaskQueuePanel() {
  const [tasks, setTasks] = useState<Task[]>([]);

  useEffect(() => taskQueue.subscribe(setTasks), []);

  const statusIcon: Record<string, string> = {
    pending: '⏳',
    running: '🏃',
    progress: '🔄',
    completed: '✅',
    failed: '❌',
    canceled: '⏹',
  };

  return (
    <div className="task-queue-panel">
      <h3>任务队列 ({tasks.filter(t => t.status !== 'completed' && t.status !== 'canceled').length} 进行中)</h3>

      {tasks.map((task) => (
        <div key={task.id} className={`task-item task-${task.status}`}>
          <div className="task-header">
            <span>{statusIcon[task.status]} {task.type}</span>
            <span>{task.message}</span>
          </div>

          {/* 进度条 */}
          {(task.status === 'running' || task.status === 'progress') && (
            <div className="progress-bar-bg">
              <div
                className="progress-bar-fill"
                style={{ width: `${task.progress}%` }}
              />
              <span>{task.progress}%</span>
            </div>
          )}

          {/* 结果预览 */}
          {task.status === 'completed' && task.result && (
            <img
              src={task.result}
              alt="生成结果"
              className="task-result-preview"
              style={{ maxWidth: 200, borderRadius: 8 }}
            />
          )}

          {/* 操作按钮 */}
          <div className="task-actions">
            {task.status === 'failed' && (
              <button onClick={() => taskQueue.retryTask(task.id)}>重试</button>
            )}
            {(task.status === 'pending' || task.status === 'running' || task.status === 'progress') && (
              <button onClick={() => taskQueue.cancelTask(task.id)}>取消</button>
            )}
          </div>
        </div>
      ))}

      {tasks.some((t) => t.status === 'completed' || t.status === 'canceled') && (
        <button onClick={() => taskQueue.clearDone()}>清除已结束任务</button>
      )}
    </div>
  );
}

6. 触发创建任务

// 在业务页面中
import { taskQueue } from './TaskQueueManager';
import { useTaskWebSocket } from './useTaskWebSocket';

function ImageGenerator() {
  // 连接 WebSocket
  useTaskWebSocket();

  const handleGenerate = () => {
    taskQueue.addTask('image_gen', {
      prompt: 'a beautiful sunset over the ocean',
      model: 'sd-xl',
      steps: 25,
    });
  };

  return <button onClick={handleGenerate}>生成图片</button>;
}

7. 面试核心要点

问:实时进度为什么用 WebSocket 而不是轮询?

轮询有固定延迟(如 2 秒),浪费请求;WebSocket 是服务端主动推送,毫秒级实时更新,能精准反映 AI 推理每一阶段的进度(如图片生成从噪声到清晰的过程)。降级方案:浏览器不支持时 fallback 到 SSE。

问:刷新页面任务数据怎么恢复?

  • 已完成任务结果持久化到 localStorage
  • 排队中/执行中任务:刷新后标记为 pending 重新入队(前端无法恢复正在运行的进程)
  • 不可持久化内容:定时器、WebSocket 实例、DOM 引用等运行时变量

问:并发控制怎么做的?

纯前端控制并发数(默认 2),通过计数器 runningCount 限制同时执行的任务数,一个任务结束后自动从队列取下一个。后端也需要做并发控制作为保底。

总结

  • 状态机优先:6 种状态 + 流转规则,保证任务生命周期可追溯
  • 单例队列:全局唯一的 TaskQueueManager,任何组件都能 push 任务
  • WebSocket 实时推送:服务端主动推送进度,前端只做展示更新
  • 并发锁runningCount < maxConcurrency 控制并发,完成后自动出队
  • localStorage 持久化:刷新页不丢任务,运行时变量不持久化

Comments 留言讨论

还没有评论,来抢个沙发,聊聊你的看法~

Michael.Meng

michaelnews@126.com
用 AI 记录,用文字沉淀

© 2026 Michael Meng · 保留所有权利 · Powered by FastAPI + Nuxt