多 Agent 协作编排:从任务分解到执行监控的工程实践
多 Agent 协作编排:从任务分解到执行监控的工程实践

一、单 Agent 的能力天花板:复杂任务的协作困境
当 AI Agent 从概念验证走向生产环境时,一个核心瓶颈迅速暴露:单 Agent 处理复杂任务时,上下文窗口溢出、工具调用链过长、错误恢复困难。以一个典型的"技术调研报告生成"任务为例,单 Agent 需要同时负责信息检索、内容筛选、结构编排、格式校验和引用核查——每个子任务的认知负载叠加后,输出质量急剧下降。
多 Agent 协作的核心思路是将复杂任务拆解为职责单一的子任务,每个 Agent 只负责一个明确的工作域。但多 Agent 系统引入了新的工程挑战:Agent 间的通信协议如何设计?任务依赖关系如何编排?部分 Agent 执行失败时如何降级?这些问题决定了系统从 Demo 到生产的距离。
二、多 Agent 编排的核心架构与通信机制
多 Agent 系统的架构设计需要解决三个核心问题:任务分解、通信协议和执行监控。
flowchart TB
subgraph 任务分解层
A[用户请求] --> B[Planner Agent]
B --> C[任务 DAG 生成]
C --> D[子任务 1: 信息检索]
C --> E[子任务 2: 数据分析]
C --> F[子任务 3: 内容生成]
end
subgraph 通信与执行层
D --> G[消息总线]
E --> G
F --> G
G --> H[共享状态存储]
H --> I[结果聚合器]
end
subgraph 监控与容错层
G --> J[执行追踪器]
J --> K{超时/失败?}
K -->|是| L[降级策略]
K -->|否| M[继续执行]
L --> N[重试/跳过/人工介入]
I --> O[最终输出]
end
任务分解由 Planner Agent 完成,它将用户请求解析为有向无环图(DAG),每个节点代表一个子任务,边代表依赖关系。关键设计决策:DAG 的粒度不宜过细——每个子任务应是一个 Agent 可以独立完成的完整工作单元,而非单步操作。
通信协议采用消息总线模式,Agent 之间不直接调用,而是通过消息总线传递结构化消息。这解耦了 Agent 的执行时序,同时支持异步并行。消息格式必须包含:发送者 ID、接收者 ID、消息类型(请求/响应/事件)、负载内容和关联的任务 ID。
执行监控通过执行追踪器实现,记录每个子任务的开始时间、结束时间、状态和输出摘要。当检测到超时或失败时,触发降级策略:重试(可恢复的临时错误)、跳过(非关键路径的子任务)或人工介入(关键路径的不可恢复错误)。
三、生产级多 Agent 编排框架的实现
import asyncio
import uuid
from dataclasses import dataclass, field
from enum import Enum
from typing import Any, Callable, Optional
class TaskStatus(Enum):
PENDING = "pending"
RUNNING = "running"
SUCCESS = "success"
FAILED = "failed"
SKIPPED = "skipped"
@dataclass
class AgentMessage:
"""Agent 间通信的结构化消息"""
sender: str
receiver: str
msg_type: str # request / response / event
payload: dict
task_id: str
correlation_id: str = field(default_factory=lambda: uuid.uuid4().hex[:8])
@dataclass
class SubTask:
"""子任务定义"""
task_id: str
agent_name: str
handler: Callable
dependencies: list[str] = field(default_factory=list)
timeout: float = 60.0
max_retries: int = 1
# 非关键路径的子任务失败后可跳过
critical: bool = True
class MessageBus:
"""异步消息总线,解耦 Agent 间的直接依赖"""
def __init__(self):
self._subscribers: dict[str, list[asyncio.Queue]] = {}
def subscribe(self, agent_name: str) -> asyncio.Queue:
if agent_name not in self._subscribers:
self._subscribers[agent_name] = []
queue: asyncio.Queue = asyncio.Queue()
self._subscribers[agent_name].append(queue)
return queue
async def publish(self, message: AgentMessage):
subscribers = self._subscribers.get(message.receiver, [])
for queue in subscribers:
await queue.put(message)
class MultiAgentOrchestrator:
"""多 Agent 编排引擎"""
def __init__(self):
self.bus = MessageBus()
self.shared_state: dict[str, Any] = {}
self.task_status: dict[str, TaskStatus] = {}
self._execution_log: list[dict] = []
def _log(self, task_id: str, event: str, detail: str = ""):
self._execution_log.append({
"task_id": task_id,
"event": event,
"detail": detail,
})
async def _execute_task(self, task: SubTask) -> Any:
"""执行单个子任务,含超时与重试"""
self.task_status[task.task_id] = TaskStatus.RUNNING
self._log(task.task_id, "start")
for attempt in range(task.max_retries + 1):
try:
result = await asyncio.wait_for(
task.handler(self.shared_state, self.bus),
timeout=task.timeout,
)
self.task_status[task.task_id] = TaskStatus.SUCCESS
self.shared_state[task.task_id] = result
self._log(task.task_id, "success")
return result
except asyncio.TimeoutError:
self._log(task.task_id, "timeout", f"attempt={attempt}")
except Exception as e:
self._log(task.task_id, "error", f"attempt={attempt}, err={e}")
if attempt < task.max_retries:
await asyncio.sleep(1.0 * (2 ** attempt))
# 重试耗尽
if task.critical:
self.task_status[task.task_id] = TaskStatus.FAILED
raise RuntimeError(f"关键任务 {task.task_id} 执行失败")
else:
self.task_status[task.task_id] = TaskStatus.SKIPPED
self._log(task.task_id, "skipped", "non-critical task failed")
return None
async def run(self, tasks: list[SubTask]) -> dict:
"""按 DAG 依赖关系编排执行"""
# 初始化所有任务状态
for t in tasks:
self.task_status[t.task_id] = TaskStatus.PENDING
completed: set[str] = set()
max_iterations = len(tasks) * 2 # 防止死循环
for _ in range(max_iterations):
# 找出所有依赖已满足的待执行任务
ready = [
t for t in tasks
if self.task_status[t.task_id] == TaskStatus.PENDING
and all(d in completed for d in t.dependencies)
]
if not ready:
if len(completed) == len(tasks):
break
# 存在无法满足依赖的任务,标记为失败
for t in tasks:
if self.task_status[t.task_id] == TaskStatus.PENDING:
self.task_status[t.task_id] = TaskStatus.FAILED
break
# 并行执行所有就绪任务
coros = [self._execute_task(t) for t in ready]
results = await asyncio.gather(*coros, return_exceptions=True)
for t, r in zip(ready, results):
if isinstance(r, Exception):
if t.critical:
# 关键任务失败,终止整个编排
return self._build_result()
completed.add(t.task_id)
return self._build_result()
def _build_result(self) -> dict:
return {
"status": self.task_status,
"shared_state": self.shared_state,
"execution_log": self._execution_log,
}
关键设计点:
-
消息总线解耦:Agent 之间通过
MessageBus通信,而非直接函数调用。这使得 Agent 可以独立部署和扩展,也为后续引入远程 Agent(如通过 HTTP 调用的第三方服务)预留了接口。 -
DAG 依赖调度:编排引擎按拓扑顺序执行任务,同一层级的任务并行执行。依赖检测通过
completed集合实现,避免轮询开销。 -
关键路径保护:
critical标志区分关键和非关键任务。非关键任务失败后跳过,不影响整体流程;关键任务失败则终止编排。
四、多 Agent 系统的工程代价与适用边界
通信开销不可忽视。Agent 间的消息传递引入了序列化/反序列化成本和异步等待延迟。在子任务执行时间小于 100ms 的场景下,通信开销可能占总耗时的 50% 以上。此时单 Agent 方案反而更高效。
状态一致性是隐性地雷。共享状态存储在并发写入时可能出现竞态条件。上述实现使用字典存储,在单进程内是安全的,但一旦扩展到多进程或多节点,就需要引入分布式锁或 CRDT(无冲突复制数据类型)来解决一致性问题。
调试难度指数级增长。当 5 个 Agent 协作完成一个任务时,一次失败可能涉及 3 个 Agent 的交互链路。执行日志(_execution_log)是排障的唯一线索,必须包含完整的消息流转记录。
适用边界:多 Agent 编排适合"子任务之间有明确边界、需要不同工具或模型、执行时间在秒级以上"的场景。对于"子任务之间高度耦合、执行时间在毫秒级"的场景,单 Agent 方案更优。
五、总结
多 Agent 协作编排的核心价值在于将复杂任务的认知负载分散到职责单一的 Agent 上,通过 DAG 依赖调度实现并行执行和容错降级。落地时需把握三个要点:第一,任务分解粒度以"一个 Agent 可独立完成的工作单元"为标准,过细的粒度反而增加通信开销;第二,Agent 间通过消息总线通信而非直接调用,为后续扩展和独立部署预留空间;第三,区分关键路径和非关键路径,非关键任务失败后跳过而非终止,提升系统的韧性。多 Agent 不是万能架构,在子任务高度耦合或执行时间极短的场景下,单 Agent 方案仍然是更务实的选择。
更多推荐


所有评论(0)