f4c1a10403
- 修复 api.md 修改密码接口混入注册场景的 409 错误 - 补充 api.md WebSocket 示例缺失的 quality_supervisor running 消息 - 补充 api.md 注册接口格式校验错误响应、缓存接口详细说明 - 补充 frontend.md Task 接口缺失的 error 字段 - 补充 frontend.md 错误处理章节(API 错误、WebSocket 断连、加载状态) - frontend.md 风格键表改为引用 style-keys.md,消除重复 - async-tasks.md 去重/文件存储改为引用 database.md,消除重复 - backend.md ER 图改为引用 database.md,消除重复 - 补充 backend.md 目录结构中缺失的 auth/project handler 和 user model - 补充 backend.md 缓存层设计说明 - 补充 database.md GEN2D_SERVER_PORT 环境变量和迁移文件归属说明 - 补充 multi-agent-pipeline.md 到 backend.md 和 async-tasks.md 的交叉引用 - 补充 async-tasks.md 超时控制说明
5.7 KiB
5.7 KiB
异步任务
概述
生成任务采用「提交-异步执行-轮询/推送」模式: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 | 并发执行任务数 |
超时控制
每个任务的 processTask 执行设置超时上下文(默认 10 分钟),超时后:
- 取消正在执行的管线节点(如 AI 推理调用)
- 更新
task.status = failed,task.error = "任务执行超时" - Worker 释放,继续处理下一个任务
避免长时间卡住的任务永久占用 Worker。
管线阶段与进度
每个阶段对应 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 设计。
去重
去重策略详见 数据存储 — 去重策略。相同输入的并发请求直接返回已有 taskId,任务完成后从内存清除。
文件存储
任务完成后,FormatAdapter 将素材上传至 OSS。素材文件的存储结构与访问方式详见 数据存储 — 文件存储。