2026-05-23 18:49:50 +08:00
|
|
|
|
# 异步任务
|
|
|
|
|
|
|
2026-05-24 10:38:31 +08:00
|
|
|
|
## 概述
|
|
|
|
|
|
|
|
|
|
|
|
生成任务采用「提交-异步执行-轮询/推送」模式:API 同步返回 taskId,后端异步执行四阶段管线,前端通过轮询或 WebSocket 获取进度与结果。
|
|
|
|
|
|
|
|
|
|
|
|
## 任务生命周期
|
|
|
|
|
|
|
|
|
|
|
|
```mermaid
|
|
|
|
|
|
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'` |
|
2026-05-24 10:56:40 +08:00
|
|
|
|
| `completed` | 管线完成,素材已保存到本地 | `task.status = 'completed'`, 写入 `asset` 表 |
|
2026-05-24 10:38:31 +08:00
|
|
|
|
| `failed` | 执行失败或超过重试上限 | `task.status = 'failed'`, `task.error` 记录原因 |
|
|
|
|
|
|
|
|
|
|
|
|
## 任务队列
|
|
|
|
|
|
|
|
|
|
|
|
采用内存 Channel 队列(非 Redis),适合单实例部署场景:
|
|
|
|
|
|
|
|
|
|
|
|
```go
|
|
|
|
|
|
// 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 更新
|
2026-05-24 10:56:40 +08:00
|
|
|
|
├── 5. 保存素材到本地,写入 asset 表
|
2026-05-24 10:38:31 +08:00
|
|
|
|
├── 6. 更新 task.status = completed
|
|
|
|
|
|
└── 7. 异常时更新 task.status = failed, task.error = 错误信息
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
|
|
### 并发控制
|
|
|
|
|
|
|
|
|
|
|
|
通过固定数量的 Worker goroutine 控制并发:
|
|
|
|
|
|
|
|
|
|
|
|
```go
|
|
|
|
|
|
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 | 并发执行任务数 |
|
|
|
|
|
|
|
2026-05-24 10:47:00 +08:00
|
|
|
|
### 超时控制
|
|
|
|
|
|
|
|
|
|
|
|
每个任务的 `processTask` 执行设置超时上下文(默认 10 分钟),超时后:
|
|
|
|
|
|
|
|
|
|
|
|
1. 取消正在执行的管线节点(如 AI 推理调用)
|
|
|
|
|
|
2. 更新 `task.status = failed`,`task.error = "任务执行超时"`
|
|
|
|
|
|
3. Worker 释放,继续处理下一个任务
|
|
|
|
|
|
|
|
|
|
|
|
避免长时间卡住的任务永久占用 Worker。
|
|
|
|
|
|
|
2026-05-24 10:38:31 +08:00
|
|
|
|
## 管线阶段与进度
|
|
|
|
|
|
|
|
|
|
|
|
每个阶段对应 Eino Graph 的一个节点,执行过程中更新 `task.stage` 和 `task.progress`:
|
|
|
|
|
|
|
|
|
|
|
|
| 阶段 | stage 值 | 进度范围 | 说明 |
|
|
|
|
|
|
|------|----------|---------|------|
|
|
|
|
|
|
| 提示词构建 | `prompt_builder` | 0-20% | 合并风格,生成三段式提示词 |
|
|
|
|
|
|
| 素材生成 | `asset_generator` | 20-60% | 调用 AI 推理 API 出图 |
|
|
|
|
|
|
| 质量检查 | `quality_supervisor` | 60-80% | 视觉模型质检 |
|
2026-05-24 10:56:40 +08:00
|
|
|
|
| 格式适配 | `format_adapter` | 80-100% | 格式转换、spritesheet 打包、保存素材 |
|
2026-05-24 10:38:31 +08:00
|
|
|
|
|
|
|
|
|
|
进度更新通过 WebSocket 实时推送给前端(参见 [API 设计 - WebSocket 消息格式](api.md))。
|
|
|
|
|
|
|
|
|
|
|
|
## 重试策略
|
|
|
|
|
|
|
|
|
|
|
|
### 质检重试(管线内重试)
|
|
|
|
|
|
|
|
|
|
|
|
QualitySupervisor 质检不通过时,Eino Graph 分支回到 PromptBuilder 重新生成:
|
|
|
|
|
|
|
|
|
|
|
|
```
|
|
|
|
|
|
QualitySupervisor -- fail --> PromptBuilder --> AssetGenerator --> QualitySupervisor
|
|
|
|
|
|
QualitySupervisor -- pass --> FormatAdapter
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
|
|
| 参数 | 默认值 | 说明 |
|
|
|
|
|
|
|------|--------|------|
|
|
|
|
|
|
| 最大重试次数 | 3 | `task.retry_count` 达到上限后降级输出 |
|
|
|
|
|
|
| 重试触发条件 | 质检不通过 | 视觉模型判定风格不一致 |
|
|
|
|
|
|
|
|
|
|
|
|
超过重试次数后,跳过质检直接进入 FormatAdapter 输出(降级策略,保证任务不会无限循环)。
|
|
|
|
|
|
|
|
|
|
|
|
### 任务级重试(管线外重试)
|
|
|
|
|
|
|
2026-05-24 10:56:40 +08:00
|
|
|
|
管线执行过程中发生不可恢复的错误(如 AI API 超时、文件写入失败):
|
2026-05-24 10:38:31 +08:00
|
|
|
|
|
|
|
|
|
|
| 参数 | 默认值 | 说明 |
|
|
|
|
|
|
|------|--------|------|
|
|
|
|
|
|
| 最大重试次数 | 1 | 仅重试一次 |
|
|
|
|
|
|
| 重试间隔 | 5 秒 | 固定间隔 |
|
|
|
|
|
|
|
|
|
|
|
|
重试时重新执行完整管线,不保留上次中间状态。
|
|
|
|
|
|
|
|
|
|
|
|
## 进度推送
|
|
|
|
|
|
|
|
|
|
|
|
前端可通过两种方式获取任务进度:
|
|
|
|
|
|
|
|
|
|
|
|
### 轮询
|
|
|
|
|
|
|
|
|
|
|
|
```
|
|
|
|
|
|
GET /api/v1/tasks/:taskId
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
|
|
前端定时调用(建议间隔 2-3 秒),简单可靠,适合不需要实时性的场景。
|
|
|
|
|
|
|
|
|
|
|
|
### WebSocket
|
|
|
|
|
|
|
|
|
|
|
|
```
|
2026-05-24 10:56:40 +08:00
|
|
|
|
ws://host/api/v1/tasks/:taskId/ws
|
2026-05-24 10:38:31 +08:00
|
|
|
|
```
|
|
|
|
|
|
|
2026-05-24 10:56:40 +08:00
|
|
|
|
长连接实时推送,通过 httpOnly Cookie 自动认证,每个阶段的状态变更立即通知前端。消息格式参见 [API 设计](api.md)。
|
2026-05-24 10:38:31 +08:00
|
|
|
|
|
|
|
|
|
|
## 去重
|
|
|
|
|
|
|
2026-05-24 10:47:00 +08:00
|
|
|
|
去重策略详见 [数据存储 — 去重策略](database.md#去重策略)。相同输入的并发请求直接返回已有 taskId,任务完成后从内存清除。
|
2026-05-24 10:38:31 +08:00
|
|
|
|
|
|
|
|
|
|
## 文件存储
|
|
|
|
|
|
|
2026-05-24 10:56:40 +08:00
|
|
|
|
任务完成后,FormatAdapter 将素材保存到本地文件系统。素材文件的存储结构与访问方式详见 [数据存储 — 文件存储](database.md#文件存储)。
|