Files
gen2d/docs/async-tasks.md
T
wonder 5a6bcaffaf docs: 新增用户认证、OSS 存储、异步任务设计
- database.md: 新增 user 表,project 关联 user_id,恢复 OSS 对象存储,
  移除 60 分钟文件过期,新增 JWT/OSS 环境变量配置
- api.md: 新增认证章节和错误码表,补充用户注册/登录/修改密码接口,
  工程管理新增列表/详情/删除接口,各接口补充错误响应说明
- async-tasks.md: 完善任务队列设计、Worker 流程、并发控制、管线进度、
  质检重试与任务级重试策略、进度推送方式
2026-05-24 10:38:31 +08:00

5.8 KiB
Raw Blame History

异步任务

概述

生成任务采用「提交-异步执行-轮询/推送」模式:API 同步返回 taskId,后端异步执行四阶段管线,前端通过轮询或 WebSocket 获取进度与结果。

任务生命周期

stateDiagram-v2
    [*] --> pending : POST /api/v1/generate
    pending --> running : Worker 领取任务
    running --> completed : 四阶段全部通过
    running --> failed : 节点执行失败 / 超过重试次数
    completed --> [*]
    failed --> [*]
状态 说明 持久化
pending 已入队,等待 Worker 领取 task.status = 'pending'
running 管线执行中,task.stage 记录当前阶段 task.status = 'running'
completed 管线完成,素材已上传至 OSS task.status = 'completed', 写入 asset 表
failed 执行失败或超过重试上限 task.status = 'failed', task.error 记录原因

任务队列

采用内存 Channel 队列(非 Redis),适合单实例部署场景:

// jobqueue.go
type JobQueue struct {
    ch   chan string      // taskId channel
    done chan struct{}
}

func NewJobQueue(size int) *JobQueue {
    return &JobQueue{
        ch:   make(chan string, size),
        done: make(chan struct{}),
    }
}

func (q *JobQueue) Push(taskId string) {
    q.ch <- taskId
}

func (q *JobQueue) Pop() (string, bool) {
    select {
    case taskId := <-q.ch:
        return taskId, true
    case <-q.done:
        return "", false
    }
}

提交流程

POST /api/v1/generate
    │
    ├── 1. 去重检查(sync.Map,hash(prompt + assetType + params))
    ├── 2. 写入 task 表(status=pending)
    ├── 3. 推入 JobQueue
    └── 4. 同步返回 taskId

Worker 流程

JobQueue.Pop() 取出 taskId
    │
    ├── 1. 更新 task.status = running
    ├── 2. 查询 task + project_style
    ├── 3. 构造 PipelineInput
    ├── 4. RunPipeline(ctx, input)
    │       ├── PromptBuilder → stage 更新
    │       ├── AssetGenerator → stage 更新
    │       ├── QualitySupervisor → stage 更新(可能触发重试分支)
    │       └── FormatAdapter → stage 更新
    ├── 5. 上传素材至 OSS,写入 asset 表
    ├── 6. 更新 task.status = completed
    └── 7. 异常时更新 task.status = failed, task.error = 错误信息

并发控制

通过固定数量的 Worker goroutine 控制并发:

func (s *TaskService) StartWorkers(n int) {
    for i := 0; i < n; i++ {
        go func(workerID int) {
            for {
                taskId, ok := s.queue.Pop()
                if !ok {
                    return
                }
                s.processTask(context.Background(), taskId)
            }
        }(i)
    }
}
参数 默认值 说明
队列容量 100 内存 channel 缓冲大小
Worker 数 3 并发执行任务数

管线阶段与进度

每个阶段对应 Eino Graph 的一个节点,执行过程中更新 task.stage 和 task.progress:

阶段 stage 值 进度范围 说明
提示词构建 prompt_builder 0-20% 合并风格,生成三段式提示词
素材生成 asset_generator 20-60% 调用 AI 推理 API 出图
质量检查 quality_supervisor 60-80% 视觉模型质检
格式适配 format_adapter 80-100% 格式转换、spritesheet 打包、上传 OSS

进度更新通过 WebSocket 实时推送给前端(参见 API 设计 - WebSocket 消息格式)。

重试策略

质检重试(管线内重试)

QualitySupervisor 质检不通过时,Eino Graph 分支回到 PromptBuilder 重新生成:

QualitySupervisor -- fail --> PromptBuilder --> AssetGenerator --> QualitySupervisor
QualitySupervisor -- pass --> FormatAdapter
参数 默认值 说明
最大重试次数 3 task.retry_count 达到上限后降级输出
重试触发条件 质检不通过 视觉模型判定风格不一致

超过重试次数后,跳过质检直接进入 FormatAdapter 输出(降级策略,保证任务不会无限循环)。

任务级重试(管线外重试)

管线执行过程中发生不可恢复的错误(如 AI API 超时、OSS 上传失败):

参数 默认值 说明
最大重试次数 1 仅重试一次
重试间隔 5 秒 固定间隔

重试时重新执行完整管线,不保留上次中间状态。

进度推送

前端可通过两种方式获取任务进度:

轮询

GET /api/v1/tasks/:taskId

前端定时调用(建议间隔 2-3 秒),简单可靠,适合不需要实时性的场景。

WebSocket

ws://host/api/v1/tasks/:taskId/ws?token=<jwt>

长连接实时推送,每个阶段的状态变更立即通知前端。消息格式参见 API 设计。

去重

相同 prompt + assetType + params 的并发请求,通过应用层去重避免重复生成:

  1. 计算 hash(prompt + assetType + params) 作为去重键
  2. sync.Map 中查找:若存在且任务仍在运行中,直接返回已有 taskId
  3. 任务完成或失败后从内存中清除

去重范围为当前实例内存,不跨实例共享。相同输入在不同时间点允许重新生成。

文件存储

任务完成后,FormatAdapter 将素材上传至 OSS:

  • 开发环境:写入本地 data/ 目录,通过 Gin 静态文件服务访问
  • 生产环境:上传至 OSS 存储桶,asset.url 存储完整 CDN URL

OSS Key 结构:users/{userId}/projects/{projectId}/tasks/{taskId}/output/

素材持久化保存,随工程生命周期管理,删除工程时级联删除 OSS 文件。