当一个团队开始用 Agent 替代日常开发任务时,最棘手的问题不是 Agent 能不能写好代码,而是怎么把正确的任务在正确的时间交给正确的 Agent。本文从零设计一套生产级 Agent 任务队列系统,覆盖入队、调度、执行、重试、超时、死信和人工接管全流程。

Agent 任务队列设计:从 Issue 到执行器的调度模型

当一个团队开始用 Agent 替代日常开发任务时,最棘手的问题不是 Agent 能不能写好代码,而是怎么把正确的任务在正确的时间交给正确的 Agent。本文从零设计一套生产级 Agent 任务队列系统,覆盖入队、调度、执行、重试、超时、死信和人工接管全流程。

一、为什么需要任务队列

没有队列的 Agent 平台就像没有调度器的出租车公司——每个司机自己找活干。

1.1 直调模式的痛点

先看一个反面例子。很多团队一开始把 Agent 当"高级函数"调用:

python
# 反例:直调模式
agent.run("帮我修复这个 Bug")
# 问题:
# 1. 当前 Agent 忙怎么办?—— 阻塞等待
# 2. 任务优先级怎么处理?—— 没有优先级
# 3. 失败后怎么办?—— 没有重试
# 4. 谁发起了什么任务?—— 没有审计
# 5. 高峰期怎么限流?—— 无法控制

直调模式在个人场景下可用,但团队场景下 5 个基础需求它一个都满足不了:

需求 直调模式 队列模式
异步处理 ❌ 阻塞等待 ✅ 入队即返回
优先级调度 ❌ 先到先做 ✅ 按 P0/P1/P2 插队
失败重试 ❌ 失败即丢 ✅ 指数退避重试
审计追踪 ❌ 无记录 ✅ 全生命周期事件
弹性伸缩 ❌ 固定实例 ✅ 根据队列深度扩缩

1.2 队列带来的核心收益

引入任务队列后,整个系统从"函数调用"变成了"事件驱动":

text
┌─────────────┐    入队    ┌──────────┐    出队    ┌─────────────┐
│  Issue/PR   │ ────────→ │  Task     │ ────────→ │   Agent     │
│  API 请求   │           │  Queue    │           │   执行器    │
│  定时任务   │           │           │           │  (沙箱内)   │
│  人工指派   │           │  持久化   │           │             │
└─────────────┘           └──────┬───┘           └─────────────┘
                                 │
                          ┌──────▼───────┐
                          │  Audit Log   │
                          │  (事件溯源)  │
                          └──────────────┘
  • 解耦:任务生产者和消费者互相不知道对方的存在
  • 削峰:Issue 洪峰不会压垮 Agent,队列自然缓冲
  • 可观测:每个任务的状态、耗时、重试次数都可追踪
  • 可恢复:队列持久化到数据库,系统重启不丢任务

二、任务数据模型

2.1 队列表设计

无论底层用 Redis、BullMQ、Celery 还是自己写 SQL 队列表,核心字段是通用的。这里给出一个完整的 PostgreSQL 实现:

sql
CREATE TABLE agent_tasks (
    id              UUID PRIMARY KEY DEFAULT gen_random_uuid(),
    -- 任务来源
    source          VARCHAR(32)  NOT NULL,  -- 'issue', 'pr', 'api', 'cron', 'manual'
    source_id       VARCHAR(128),           -- GitHub Issue #ID, PR #ID, etc.
    -- 任务内容
    title           VARCHAR(256) NOT NULL,
    task_type       VARCHAR(64)  NOT NULL,  -- 'bugfix', 'feature', 'review', 'test', 'docs'
    payload         JSONB        NOT NULL,  -- Agent 执行所需的上下文包
    priority        SMALLINT     NOT NULL DEFAULT 5,  -- 0=最高, 9=最低
    -- 调度控制
    status          VARCHAR(16)  NOT NULL DEFAULT 'pending',
                    -- pending → queued → executing → completed
                    --                                  → failed → pending(重试)
                    --              → cancelled
                    --              → pending(超时重试)
                    --              → pending(人工接管后重指派)
    max_retries     SMALLINT     NOT NULL DEFAULT 3,
    retry_count     SMALLINT     NOT NULL DEFAULT 0,
    timeout_secs    INTEGER      NOT NULL DEFAULT 300,  -- 单次执行超时
    lock_token      UUID,                  -- 执行器持有,防重复调度
    locked_at       TIMESTAMPTZ,           -- 锁定时间,用于超时检测
    -- 执行与回调
    assigned_to     VARCHAR(64),            -- 执行器 ID 或 Agent 名称
    callback_url    VARCHAR(512),           -- 完成后回调(可选)
    result          JSONB,                  -- 执行结果摘要
    -- 审计
    requested_by    VARCHAR(64)  NOT NULL,  -- 发起人
    approved_by     VARCHAR(64),            -- 审批人(如果需要)
    created_at      TIMESTAMPTZ  NOT NULL DEFAULT NOW(),
    updated_at      TIMESTAMPTZ  NOT NULL DEFAULT NOW()
);

-- 核心查询索引
CREATE INDEX idx_tasks_status_priority ON agent_tasks (status, priority, created_at)
    WHERE status IN ('pending', 'queued');
CREATE INDEX idx_tasks_locked_at ON agent_tasks (lock_token, locked_at)
    WHERE lock_token IS NOT NULL;
CREATE INDEX idx_tasks_source ON agent_tasks (source, source_id);

-- 去重约束:同一来源同一 ID 只能有一个活跃任务
CREATE UNIQUE INDEX idx_tasks_dedup ON agent_tasks (source, source_id)
    WHERE status NOT IN ('completed', 'cancelled', 'failed');

2.2 状态机设计

任务在它的生命周期里会在以下状态之间流转:

text
                  ┌──────────────────────┐
                  │      pending          │ ← 新创建,未入调度
                  └──────────┬───────────┘
                             │ 调度器扫描
                             ▼
                  ┌──────────────────────┐
                  │       queued          │ ← 已入队,等待执行器领取
                  └──────────┬───────────┘
                             │ 执行器领取
                             ▼
                  ┌──────────────────────┐
                  │     executing         │ ← 正在执行中
                  └──────┬───────────┬───┘
                         │           │
                   成功完成     失败/超时
                         │           │
                         ▼           ▼
                  ┌──────────┐ ┌──────────────┐
                  │completed │ │    failed     │
                  └──────────┘ │              │
                               │  重试次数用完?│
                               ├── 是 → dead   │
                               │  否 → pending │
                               └──────────────┘
                                         │
                                   人工接管
                                         │
                                         ▼
                               ┌──────────────────┐
                               │ manual_intervene  │
                               └──────────────────┘

关键状态变更规则:

  • pending → queued:调度器扫描,只有 status='pending'retry_count < max_retries 的才会入队
  • queued → executing:执行器用 UPDATE ... WHERE lock_token IS NULL LIMIT 1 原子领取
  • executing → failed → pending:失败后自动重试,retry_count +1
  • executing → failed (dead):重试次数用完,标记 dead 等待人工处理
  • 任何时候都可以 cancelled:由用户或上游系统主动取消

三、调度器核心实现

3.1 调度器职责

调度器是整个系统的"心脏",它只做三件事:

  1. 心跳检测:定期扫描 executing 状态中超过 timeout_secs 的任务,标记为失败
  2. 入队扫描:将 pending 状态的任务按优先级排序后标记为 queued
  3. 死信转移:将重试次数超限的任务标记为 dead

调度器不需要知道 Agent 在哪里、在做什么——它只关心任务的状态和时间。下面是一个轻量级调度器实现:

python
import time
import uuid
from datetime import datetime, timedelta, timezone
from typing import Optional

class TaskScheduler:
    """轻量级任务调度器,适用于中小团队(<50 Agent)"""

    def __init__(self, db_pool, poll_interval: float = 1.0):
        self.db = db_pool
        self.poll_interval = poll_interval  # 秒
        self.running = False

    def tick(self):
        """一次调度心跳:超时检测 → 优先级入队 → 死信转移"""
        now = datetime.now(timezone.utc)

        # 1. 超时检测:执行超时的任务标记为失败
        self.db.execute("""
            UPDATE agent_tasks
            SET status = 'failed',
                retry_count = retry_count + 1,
                updated_at = NOW()
            WHERE status = 'executing'
              AND locked_at < $1 - $2::interval
              AND lock_token IS NOT NULL
        """, now, timedelta(seconds=300))  # 默认超时 300s

        # 2. 优先级入队:把 pending 任务按规定数量入队
        self.db.execute("""
            UPDATE agent_tasks
            SET status = 'queued',
                updated_at = NOW()
            WHERE id IN (
                SELECT id FROM agent_tasks
                WHERE status = 'pending'
                  AND retry_count < max_retries
                ORDER BY priority ASC, created_at ASC
                LIMIT $1   -- 每次最多入队 N 个
                FOR UPDATE SKIP LOCKED
            )
        """, self.max_concurrent)

        # 3. 死信转移:重试次数超限的任务
        self.db.execute("""
            UPDATE agent_tasks
            SET status = 'dead',
                updated_at = NOW()
            WHERE status = 'failed'
              AND retry_count >= max_retries
        """)

    def run_forever(self):
        """启动调度循环"""
        self.running = True
        while self.running:
            self.tick()
            time.sleep(self.poll_interval)

    def stop(self):
        self.running = False

3.2 优先级策略

优先级分两级:静态优先级(任务创建时指定)和 动态优先级(调度器运行时调整)。

静态优先级参考标准:

优先级 适用场景 响应时间
P0 (Critical) 0 生产故障修复、安全漏洞 < 1 min
P1 (High) 1 阻塞性 Bug、紧急需求 < 5 min
P2 (Medium) 3 普通 Bug、Feature < 30 min
P3 (Low) 5 重构、技术债务 < 2 h
P4 (Backlog) 7 文档、实验性任务 空闲时
P5 (Idea) 9 想法、待讨论 永不

动态优先级调整原则:

text
# 等得越久,优先级越高(aging)
有效优先级 = 静态优先级 - floor(等待时间 / aging_interval)
aging_interval = 10 分钟

# 例:一个 P3 任务等了 30 分钟
# 有效优先级 = 5 - floor(30/10) = 5 - 3 = 2 → 相当于 P1

这样设计确保低优先级任务不会永久饿死——它们等得够久后会自动"插队"。

四、执行器设计与领取机制

4.1 原子领取

执行器通过乐观锁原子领取任务,避免两个执行器同时处理同一任务:

python
def claim_task(self) -> Optional[dict]:
    """原子领取一个 queued 任务"""
    lock = uuid.uuid4()
    now = datetime.now(timezone.utc)

    result = self.db.fetchone("""
        UPDATE agent_tasks
        SET status = 'executing',
            lock_token = $1,
            locked_at = $2,
            assigned_to = $3,
            updated_at = $2
        WHERE id = (
            SELECT id FROM agent_tasks
            WHERE status = 'queued'
            ORDER BY priority ASC, created_at ASC
            LIMIT 1
            FOR UPDATE SKIP LOCKED
        )
        RETURNING id, title, task_type, payload, timeout_secs, callback_url
    """, lock, now, self.executor_id)

    return result

关键设计点:

  • FOR UPDATE SKIP LOCKED:PostgreSQL 9.5+ 特性,跳过已被其他事务锁定的行。多个执行器可以同时领任务而不会互相阻塞
  • lock_token:持有锁的凭证,执行器在完成任务更新时必须带上这个 token,防止其他执行器误操作
  • locked_at:记录锁定时间,调度器的超时检测依赖这个字段

4.2 结果回传

任务完成后,执行器更新任务状态并释放锁:

python
def complete_task(self, task_id: str, lock_token: str, result: dict):
    """完成任务并回传结果"""
    self.db.execute("""
        UPDATE agent_tasks
        SET status = 'completed',
            result = $1::jsonb,
            lock_token = NULL,
            locked_at = NULL,
            updated_at = NOW()
        WHERE id = $2
          AND lock_token = $3   -- 验证锁持有者
    """, json.dumps(result), task_id, lock_token)


def fail_task(self, task_id: str, lock_token: str, error: str):
    """标记任务失败(如果还有重试次数,调度器会自动重试)"""
    self.db.execute("""
        UPDATE agent_tasks
        SET status = 'failed',
            lock_token = NULL,
            locked_at = NULL,
            result = jsonb_build_object('error', $1),
            updated_at = NOW()
        WHERE id = $2
          AND lock_token = $3
    """, error, task_id, lock_token)

4.3 超时与取消机制

执行器必须支持两种中断信号:

python
import signal

class AgentExecutor:
    def execute_with_timeout(self, task: dict) -> dict:
        """带超时的 Agent 执行"""
        timeout_secs = task.get('timeout_secs', 300)

        # 方案一:SIGALRM 信号(仅 Unix)
        signal.signal(signal.SIGALRM, self._timeout_handler)
        signal.alarm(timeout_secs)
        try:
            result = self._run_agent(task)   # 实际调用 Agent
            return {'status': 'ok', 'data': result}
        except TimeoutError:
            # 超时后杀掉 Agent 子进程
            self._kill_agent_process()
            return {'status': 'timeout', 'error': f'超过 {timeout_secs}s 限制'}
        finally:
            signal.alarm(0)  # 取消闹钟


    def cancel_task(self, task_id: str):
        """由调度器或用户发起的取消"""
        # 1. 数据库标记取消
        self.db.execute("""
            UPDATE agent_tasks
            SET status = 'cancelled',
                lock_token = NULL,
                locked_at = NULL,
                updated_at = NOW()
            WHERE id = $1 AND status = 'executing'
        """, task_id)
        # 2. 杀死执行进程
        self._kill_agent_process()

五、参数说明

参数 类型 默认值 说明
priority SMALLINT 5 0=P0 ~ 9=P5,调度器按此排序出队
timeout_secs INTEGER 300 单次 Agent 执行的最大秒数,超时自动标记失败
max_retries SMALLINT 3 任务最大重试次数,重试用完后进入死信队列
retry_count SMALLINT 0 当前已重试次数,调度器自动累加
lock_token UUID NULL 执行器领取任务时生成的唯一凭证,用于防并发
locked_at TIMESTAMPTZ NULL 领取时间,调度器根据此字段判断超时
poll_interval FLOAT 1.0 调度器心跳间隔(秒),影响任务响应延迟
max_concurrent INTEGER 10 单次入队扫描最多放入 waiting 的任务数
callback_url VARCHAR NULL 任务完成/失败后,执行器向此 URL 发送结果
aging_interval INTEGER 600 动态优先级 aging 间隔(秒)
source VARCHAR 任务来源标识,用于去重和审计
source_id VARCHAR 上游系统 ID,例如 Issue #ID,用于同源去重

六、落地检查清单

  • [✅] 任务入队:通过 Issue Webhook、API 请求、定时触发器、人工指派 4 种途径都能成功入队
  • [✅] 去重:同一个 Issue/PR 不会产生重复任务(唯一索引 + status 过滤)
  • [✅] 优先级排序:P0 任务总能比 P3 任务先被执行器领取
  • [✅] 动态 aging:等待超过 aging_interval 的任务优先级自动提升
  • [✅] 原子领取:两个执行器不会领到同一个任务
  • [✅] 超时检测:超过 timeout_secs 未完成的任务自动标记失败并进入重试
  • [✅] 指数退避重试:第 N 次重试前等待 min(2^N * 10, 300) 秒
  • [✅] 死信队列:重试次数用完后进入 dead 状态,人工介入
  • [✅] 人工接管:dead 任务支持人工认领,完成后标记 completed
  • [✅] 取消:执行中的任务可以被取消(DB 标记 + 进程杀 kill)
  • [✅] 审计事件:每次状态变更记录 event_type、actor、old_status、new_status、timestamp
  • [✅] 结果回调:任务完成时向 callback_url POST 结果摘要

七、真实经验与踩坑

7.1 不要直接在生产用 `SKIP LOCKED`

FOR UPDATE SKIP LOCKED 是 PostgreSQL 的利器,但它有一个陷阱:跳过的任务不会重新触发调度器扫描。如果你的执行器一直满负荷,被跳过的任务可能永远不被领取。解决方案是在调度器 tick 中加上一个补偿逻辑——每次扫描时也检查一下 queued 状态的任务是否已经有 executing 的任务在领了,如果没有,说明上一个调度周期被跳过了,手动触发一次入队。

7.2 超时时间不能一刀切

Bugfix 任务和大型重构任务的耗时差 10 倍。我们的实践是按 task_type 设置不同的默认超时:

python
TIMEOUT_MAP = {
    'bugfix':     120,   # 2 min
    'review':     180,   # 3 min
    'feature':    600,   # 10 min
    'test':       120,   # 2 min
    'docs':        60,   # 1 min
    'refactor':   900,   # 15 min
}

7.3 死信队列一定要有人盯着

死信任务不会自己消失。我们的做法是:每天早 9 点向团队 Slack 发一条死信任务汇总,包含每个死信任务的 ID、标题、失败原因和首次创建时间。排查发现,超过 60% 的死信任务是权限不足导致的——执行器没有目标仓库的写权限,而不是 Agent 本身能力不够。

7.4 重试要有上限和止损

不要无限重试。三次重试后基本可以断定:要么是任务描述有问题,要么是环境有问题,要么是 Agent 不适合这个任务类型。继续重试除了浪费 Token 没有意义。我们的经验是:重试 3 次后标记 dead 并通知人工——人工检查后可能补全信息重新入队,也可能直接关闭。

八、生产级部署建议

8.1 推荐技术选型

场景 推荐方案 说明
轻量团队(<10 人) SQLite + 单进程调度器 部署简单,零依赖
中型团队(10-50 人) PostgreSQL + Python 调度器 本文代码可直接用
大型团队(>50 人) BullMQ (Redis) + Node.js 调度器 天然支持分布式、监控、延迟队列
已有基础设施 Celery + Redis/RabbitMQ 与 Python 生态深度集成

8.2 监控指标

配置以下 Prometheus 指标来观测队列健康度:

text
# 队列深度(按状态)
agent_queue_depth{status="pending"}    # 待入队
agent_queue_depth{status="queued"}     # 等待执行
agent_queue_depth{status="executing"}  # 正在执行
agent_queue_depth{status="dead"}       # 死信(重点关注)

# 任务延迟
agent_queue_latency_seconds   # 任务从创建到开始执行的延迟
agent_queue_completion_time   # 任务从开始到完成的耗时

# 任务结果
agent_tasks_total{status="completed"}  # 成功
agent_tasks_total{status="failed"}     # 失败
agent_tasks_total{status="timeout"}    # 超时
agent_retries_total                    # 重试次数

8.3 水平扩展

当 Agent 数量增长时,调度器可能成为瓶颈。解决方案:

  1. 分区调度:按 task_typesource 分到不同的调度器实例
  2. 读写分离:调度器只读状态和写状态变更,执行结果通过异步队列回写
  3. 分片队列:按 priority 范围分片,P0-P1 一个队列,P2-P4 一个队列,各由独立的调度器处理

这样设计,一个团队从 5 人到 200 人,队列系统只需要调整配置,不需要改架构。