Async Task Queue
Async Task Queue 是将耗时操作从请求-响应链路中解耦出来的模式。在 AI Agent 应用中,它用于后台处理消息持久化、Embedding 生成、评估任务等。
#type / concept
#status / evergreen
#tech / backend
#tech / architecture
[!info] related notes
- 所属 MOC: AI Agent Application MOC
- 相关: Job Management
Async Task Queue
一句话定义
Async Task Queue 是将耗时操作从请求-响应链路中解耦出来的模式。把任务放入队列,由后台 worker 异步执行,不阻塞主请求。
它解决什么问题
某些操作耗时但不需要立即返回结果:
- Embedding 生成(文档摄入后)
- 知识库索引更新
- 评估任务批量运行
- 消息摘要生成
- 数据导出
如果同步执行,用户要等很久。
核心原理
架构
API 请求 → 入队 → 立即返回 task_id
│
▼
Task Queue (Redis / RabbitMQ / SQS)
│
▼
Worker 消费 → 执行任务 → 更新状态
Go 实现
// 入队
func (s *TaskService) Enqueue(ctx context.Context, task Task) (string, error) {
task.ID = generateID()
task.Status = "pending"
task.CreatedAt = time.Now()
data, _ := json.Marshal(task)
err := s.redis.LPush(ctx, "tasks:pending", data)
return task.ID, err
}
// Worker 消费
func (w *Worker) Start(ctx context.Context) {
for {
data, _ := w.redis.BRPop(ctx, 0, "tasks:pending")
var task Task
json.Unmarshal([]byte(data[1]), &task)
task.Status = "running"
w.updateStatus(task)
err := w.execute(ctx, task)
if err != nil {
task.Status = "failed"
task.Error = err.Error()
} else {
task.Status = "completed"
}
w.updateStatus(task)
}
}
常见设计模式
1. 延迟队列
任务在指定时间后才执行(如会话过期清理)。
2. 优先级队列
高优先级任务先执行(如实时评估 vs 后台索引)。
3. 重试队列
失败任务自动重试,带指数退避。
常见坑
- 任务丢失: Worker 执行中崩溃,任务没有重新入队
- 不做去重: 同一个任务被多次入队
- 不做超时: Worker 卡住导致队列阻塞
- 不做死信队列: 多次失败的任务无法排查