Async Task Queue

Async Task Queue 是将耗时操作从请求-响应链路中解耦出来的模式。在 AI Agent 应用中,它用于后台处理消息持久化、Embedding 生成、评估任务等。

#type / concept #status / evergreen #tech / backend #tech / architecture

[!info] related notes

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. 重试队列

失败任务自动重试,带指数退避。

常见坑

  1. 任务丢失: Worker 执行中崩溃,任务没有重新入队
  2. 不做去重: 同一个任务被多次入队
  3. 不做超时: Worker 卡住导致队列阻塞
  4. 不做死信队列: 多次失败的任务无法排查

参考资料

创建于 2026/6/30 更新于 2026/7/15