一、为什么需要任务调度
在 SubAgent 架构中,主控 Agent 的核心职责不是"做事",而是调度。当系统从 1 个 SubAgent 扩展到 10 个、100 个时,以下问题会立刻暴露:
- 并发控制:同时启动多少个 SubAgent 不会撑爆 Token 配额和 API 限流?
- 依赖编排:SubAgent B 依赖 SubAgent A 的输出,如何保证执行顺序?
- 超时与重试:一个 SubAgent 卡住了,是等还是杀?重试几次?
- 资源回收:SubAgent 用完后的上下文、临时文件、网络连接谁来清理?
结论:没有调度层的 SubAgent 架构,本质上还是手写 if-else,只不过换了个名字。
二、SubAgent 生命周期模型
每个 SubAgent 的完整生命周期分为 5 个阶段:
CREATED → SCHEDULED → RUNNING → COMPLETED↓FAILED → RETRYING → SCHEDULED (循环)
| 阶段 | 状态 | 含义 |
|---|---|---|
| CREATED | 已创建 | SubAgent 实例化完成,资源已分配,等待调度 |
| SCHEDULED | 已调度 | 进入执行队列,等待资源配额 |
| RUNNING | 执行中 | LLM 调用 + 工具调用进行中 |
| COMPLETED | 已完成 | 执行成功,结果已写入共享存储 |
| FAILED | 失败 | 执行出错,触发重试策略或熔断 |
| RETRYING | 重试中 | 等待退避时间后重新 SCHEDULED |
3.1 调度器接口
from enum import Enum
from dataclasses import dataclass
from typing import Optionalclass SubAgentStatus(Enum):CREATED = "created"SCHEDULED = "scheduled"RUNNING = "running"COMPLETED = "completed"FAILED = "failed"RETRYING = "retrying"@dataclass
class SubAgentTask:id: stragent_type: strinput: dictstatus: SubAgentStatusretry_count: int = 0max_retries: int = 3timeout: int = 120depends_on: list[str] = Noneresult: Optional[dict] = Noneerror: Optional[str] = None
3.2 调度引擎
class SubAgentScheduler:def __init__(self, max_concurrency: int = 5):self.tasks: dict[str, SubAgentTask] = {}self.queue = asyncio.Queue()self.active = 0self.max_concurrency = max_concurrencyself.semaphore = asyncio.Semaphore(max_concurrency)async def submit(self, task: SubAgentTask):if task.depends_on:for dep_id in task.depends_on:dep = self.tasks.get(dep_id)if not dep or dep.status != SubAgentStatus.COMPLETED:await asyncio.sleep(1)return await self.submit(task)task.status = SubAgentStatus.SCHEDULEDawait self.queue.put(task)self._schedule_worker()async def _execute(self, task: SubAgentTask):async with self.semaphore:task.status = SubAgentStatus.RUNNINGself.active += 1try:result = await asyncio.wait_for(self._run_agent(task), timeout=task.timeout)task.status = SubAgentStatus.COMPLETEDtask.result = resultreturn resultexcept asyncio.TimeoutError:task.status = SubAgentStatus.FAILEDtask.error = "Timeout"await self._handle_retry(task)except Exception as e:task.status = SubAgentStatus.FAILEDtask.error = str(e)await self._handle_retry(task)finally:self.active -= 1
```## 四、编排策略模式### 4.1 顺序编排(Chain)
```python
async def chain(scheduler, tasks: list[SubAgentTask]):prev_result = Nonefor task in tasks:if prev_result:task.input["prev_output"] = prev_resultprev_result = await scheduler.submit(task)return prev_result
4.2 并行编排(Fan-out)
async def fan_out(scheduler, tasks: list[SubAgentTask]):futures = [scheduler.submit(t) for t in tasks]results = await asyncio.gather(*futures, return_exceptions=True)for r in results:if isinstance(r, Exception):raise rreturn results
4.3 DAG 编排
async def dag_execute(scheduler, dag, tasks):from collections import dequein_degree = {tid: 0 for tid in tasks}for tid, deps in dag.items():in_degree[tid] = len(deps)ready = deque([tid for tid, deg in in_degree.items() if deg == 0])results = {}while ready:batch = []while ready:batch.append(ready.popleft())futures = {tid: scheduler.submit(tasks[tid]) for tid in batch}for tid, future in futures.items():results[tid] = await futurefor tid, task in tasks.items():if task.depends_on and all(d in results for d in task.depends_on):if in_degree[tid] > 0:in_degree[tid] = 0ready.append(tid)return results
五、超时与熔断实战
class CircuitBreaker:def __init__(self, failure_threshold: int = 5, recovery_timeout: int = 30):self.failure_count = 0self.failure_threshold = failure_thresholdself.recovery_timeout = recovery_timeoutself.last_failure_time = 0self.state = "closed"async def call(self, subagent_func, *args, **kwargs):if self.state == "open":if time.time() - self.last_failure_time > self.recovery_timeout:self.state = "half-open"else:raise CircuitBreakerOpen("SubAgent 熔断中")try:result = await subagent_func(*args, **kwargs)if self.state == "half-open":self.state = "closed"self.failure_count = 0return resultexcept Exception as e:self.failure_count += 1self.last_failure_time = time.time()if self.failure_count >= self.failure_threshold:self.state = "open"raise
六、最佳实践清单
- 永远设超时:每个 SubAgent 调用必须有硬超时,推荐 60-120s
- 重试要退避:指数退避 (2^n) 避免雪崩
- 依赖显式化:不要在主控代码里隐式 await,用 DAG 声明依赖
- 资源隔离:每个 SubAgent 使用独立的 LLM 会话/连接池
- 可观测性:每个 SubAgent 上报状态、耗时、Token 消耗
总结:调度层是 SubAgent 架构的"操作系统内核",做好生命周期管理、编排策略和容错机制,才能让子代理真正高效协作。下篇我们将深入 SubAgent 通信协议设计。