Agent SubAgent 任务调度与编排:子代理生命周期管理实战

📅 2026/7/31 10:23:24 👁️ 阅读次数 📝 编程学习
Agent SubAgent 任务调度与编排:子代理生命周期管理实战

一、为什么需要任务调度

在 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

六、最佳实践清单

  1. 永远设超时:每个 SubAgent 调用必须有硬超时,推荐 60-120s
  2. 重试要退避:指数退避 (2^n) 避免雪崩
  3. 依赖显式化:不要在主控代码里隐式 await,用 DAG 声明依赖
  4. 资源隔离:每个 SubAgent 使用独立的 LLM 会话/连接池
  5. 可观测性:每个 SubAgent 上报状态、耗时、Token 消耗

总结:调度层是 SubAgent 架构的"操作系统内核",做好生命周期管理、编排策略和容错机制,才能让子代理真正高效协作。下篇我们将深入 SubAgent 通信协议设计。