Agent 任务队列设计:从 Issue 到执行器的调度模型
当一个团队开始用 Agent 替代日常开发任务时,最棘手的问题不是 Agent 能不能写好代码,而是怎么把正确的任务在正确的时间交给正确的 Agent。本文从零设计一套生产级 Agent 任务队列系统,覆盖入队、调度、执行、重试、超时、死信和人工接管全流程。
一、为什么需要任务队列
没有队列的 Agent 平台就像没有调度器的出租车公司——每个司机自己找活干。
1.1 直调模式的痛点
先看一个反面例子。很多团队一开始把 Agent 当"高级函数"调用:
# 反例:直调模式
agent.run("帮我修复这个 Bug")
# 问题:
# 1. 当前 Agent 忙怎么办?—— 阻塞等待
# 2. 任务优先级怎么处理?—— 没有优先级
# 3. 失败后怎么办?—— 没有重试
# 4. 谁发起了什么任务?—— 没有审计
# 5. 高峰期怎么限流?—— 无法控制直调模式在个人场景下可用,但团队场景下 5 个基础需求它一个都满足不了:
| 需求 | 直调模式 | 队列模式 |
|---|---|---|
| 异步处理 | ❌ 阻塞等待 | ✅ 入队即返回 |
| 优先级调度 | ❌ 先到先做 | ✅ 按 P0/P1/P2 插队 |
| 失败重试 | ❌ 失败即丢 | ✅ 指数退避重试 |
| 审计追踪 | ❌ 无记录 | ✅ 全生命周期事件 |
| 弹性伸缩 | ❌ 固定实例 | ✅ 根据队列深度扩缩 |
1.2 队列带来的核心收益
引入任务队列后,整个系统从"函数调用"变成了"事件驱动":
┌─────────────┐ 入队 ┌──────────┐ 出队 ┌─────────────┐
│ Issue/PR │ ────────→ │ Task │ ────────→ │ Agent │
│ API 请求 │ │ Queue │ │ 执行器 │
│ 定时任务 │ │ │ │ (沙箱内) │
│ 人工指派 │ │ 持久化 │ │ │
└─────────────┘ └──────┬───┘ └─────────────┘
│
┌──────▼───────┐
│ Audit Log │
│ (事件溯源) │
└──────────────┘- 解耦:任务生产者和消费者互相不知道对方的存在
- 削峰:Issue 洪峰不会压垮 Agent,队列自然缓冲
- 可观测:每个任务的状态、耗时、重试次数都可追踪
- 可恢复:队列持久化到数据库,系统重启不丢任务
二、任务数据模型
2.1 队列表设计
无论底层用 Redis、BullMQ、Celery 还是自己写 SQL 队列表,核心字段是通用的。这里给出一个完整的 PostgreSQL 实现:
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 状态机设计
任务在它的生命周期里会在以下状态之间流转:
┌──────────────────────┐
│ 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 调度器职责
调度器是整个系统的"心脏",它只做三件事:
- 心跳检测:定期扫描
executing状态中超过timeout_secs的任务,标记为失败 - 入队扫描:将
pending状态的任务按优先级排序后标记为queued - 死信转移:将重试次数超限的任务标记为
dead
调度器不需要知道 Agent 在哪里、在做什么——它只关心任务的状态和时间。下面是一个轻量级调度器实现:
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 = False3.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 | 想法、待讨论 | 永不 |
动态优先级调整原则:
# 等得越久,优先级越高(aging)
有效优先级 = 静态优先级 - floor(等待时间 / aging_interval)
aging_interval = 10 分钟
# 例:一个 P3 任务等了 30 分钟
# 有效优先级 = 5 - floor(30/10) = 5 - 3 = 2 → 相当于 P1这样设计确保低优先级任务不会永久饿死——它们等得够久后会自动"插队"。
四、执行器设计与领取机制
4.1 原子领取
执行器通过乐观锁原子领取任务,避免两个执行器同时处理同一任务:
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 结果回传
任务完成后,执行器更新任务状态并释放锁:
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 超时与取消机制
执行器必须支持两种中断信号:
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 设置不同的默认超时:
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 指标来观测队列健康度:
# 队列深度(按状态)
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 数量增长时,调度器可能成为瓶颈。解决方案:
- 分区调度:按
task_type或source分到不同的调度器实例 - 读写分离:调度器只读状态和写状态变更,执行结果通过异步队列回写
- 分片队列:按 priority 范围分片,P0-P1 一个队列,P2-P4 一个队列,各由独立的调度器处理
这样设计,一个团队从 5 人到 200 人,队列系统只需要调整配置,不需要改架构。