多 Agent 协作系统设计:从单点智能到群体决策的工程实践
多 Agent 协作系统设计:从单点智能到群体决策的工程实践

一、单 Agent 的能力天花板:复杂任务的分解困境
单个 LLM Agent 在处理简单任务时表现尚可,但面对多步骤、多领域交叉的复杂任务时,暴露出三个结构性问题。第一,上下文窗口溢出:一个 Agent 同时承担需求分析、代码生成、测试验证三重角色,Prompt 长度迅速突破 Token 限制,导致早期指令被截断遗忘。第二,角色混淆:同一个 Agent 在"创意发散"和"严格校验"之间频繁切换,输出质量波动显著,基准测试显示角色切换后的首次输出准确率下降约 22%。第三,单点故障:Agent 崩溃或陷入死循环时,整个任务链断裂,无法局部恢复。
多 Agent 协作的核心思路是"分而治之"——将复杂任务拆解为子任务,分配给专业化 Agent,通过编排层协调执行顺序与数据流转。这不是简单的并行调用,而是需要精心设计的通信协议与冲突消解机制。
二、多 Agent 协作的架构模式与通信机制
多 Agent 系统的架构模式主要有三种:顺序管道(Pipeline)、层级调度(Hierarchical)、事件驱动(Event-Driven)。不同模式适用于不同的任务特征。
graph TB
subgraph 顺序管道模式
P1[规划Agent] -->|任务描述| P2[执行Agent]
P2 -->|执行结果| P3[审查Agent]
P3 -->|修正意见| P2
end
subgraph 层级调度模式
H1[主控Agent] -->|子任务1| H2[搜索Agent]
H1 -->|子任务2| H3[编码Agent]
H1 -->|子任务3| H4[测试Agent]
H2 -->|搜索结果| H1
H3 -->|代码产出| H1
H4 -->|测试报告| H1
end
subgraph 事件驱动模式
E1[需求Agent] -->|需求事件| EB[事件总线]
EB -->|编码触发| E2[编码Agent]
EB -->|审查触发| E3[审查Agent]
E2 -->|代码事件| EB
E3 -->|反馈事件| EB
end
顺序管道适合线性依赖的任务,如"分析→编码→测试"。优点是流程清晰、调试简单;缺点是吞吐量受最慢环节制约。
层级调度适合可并行的子任务,主控 Agent 负责任务分解与结果聚合。优点是并行度高;缺点是主控 Agent 成为瓶颈,且子任务间的隐式依赖可能导致结果冲突。
事件驱动适合松耦合、动态演化的任务,Agent 通过事件总线发布/订阅消息。优点是扩展性强;缺点是调试困难,事件风暴时需要背压控制。
Agent 间的通信协议需要定义三个核心字段:task_id(任务唯一标识,用于追踪与去重)、payload(结构化数据,避免自然语言的歧义)、priority(优先级,用于资源争抢时的调度决策)。
三、生产级多 Agent 编排框架实现
以下基于 Python 实现一个层级调度模式的多 Agent 编排框架,重点展示任务分解、结果聚合与异常恢复机制。
import asyncio
import json
import uuid
from dataclasses import dataclass, field
from enum import Enum
from typing import Any, Callable, Optional
class TaskStatus(Enum):
PENDING = "pending"
RUNNING = "running"
COMPLETED = "completed"
FAILED = "failed"
@dataclass
class AgentMessage:
"""Agent 间通信的标准化消息协议"""
task_id: str
sender: str
payload: dict
priority: int = 0 # 0=低, 1=中, 2=高,调度器按优先级排序
correlation_id: str = "" # 关联ID,用于请求-响应匹配
@dataclass
class TaskResult:
status: TaskStatus
data: Any = None
error: Optional[str] = None
retry_count: int = 0
class Agent:
"""基础 Agent 抽象,所有专业化 Agent 继承此类"""
def __init__(self, name: str, max_retries: int = 2):
self.name = name
self.max_retries = max_retries # 最大重试次数,避免无限重试
async def execute(self, message: AgentMessage) -> TaskResult:
"""执行任务,内置重试与超时机制"""
for attempt in range(self.max_retries + 1):
try:
result = await asyncio.wait_for(
self._process(message),
timeout=30.0, # 单次执行超时30秒,防止Agent陷入死循环
)
return TaskResult(status=TaskStatus.COMPLETED, data=result)
except asyncio.TimeoutError:
if attempt == self.max_retries:
return TaskResult(
status=TaskStatus.FAILED,
error=f"Agent {self.name} 超时,已重试{attempt}次",
retry_count=attempt,
)
except Exception as e:
if attempt == self.max_retries:
return TaskResult(
status=TaskStatus.FAILED,
error=str(e),
retry_count=attempt,
)
return TaskResult(status=TaskStatus.FAILED, error="unreachable")
async def _process(self, message: AgentMessage) -> dict:
"""子类必须实现的具体处理逻辑"""
raise NotImplementedError
class Orchestrator:
"""主控 Agent,负责任务分解、调度与结果聚合"""
def __init__(self):
self.agents: dict[str, Agent] = {}
self.task_registry: dict[str, TaskResult] = {} # 任务注册表,用于追踪与去重
def register(self, agent: Agent):
"""注册专业化 Agent,按名称索引"""
self.agents[agent.name] = agent
async def dispatch(
self,
task_type: str,
message: AgentMessage,
) -> TaskResult:
"""
调度任务到指定 Agent,支持优先级排序与异常隔离
task_type 对应 Agent 名称,实现路由映射
"""
agent = self.agents.get(task_type)
if not agent:
return TaskResult(
status=TaskStatus.FAILED,
error=f"未注册的Agent类型: {task_type}",
)
# 任务去重:相同 task_id 不重复执行
if message.task_id in self.task_registry:
return self.task_registry[message.task_id]
result = await agent.execute(message)
self.task_registry[message.task_id] = result
return result
async def parallel_dispatch(
self,
tasks: list[tuple[str, AgentMessage]],
) -> dict[str, TaskResult]:
"""并行调度多个子任务,收集全部结果后返回"""
coroutines = [
self.dispatch(task_type, msg)
for task_type, msg in tasks
]
results = await asyncio.gather(*coroutines, return_exceptions=True)
output = {}
for (task_type, msg), result in zip(tasks, results):
if isinstance(result, Exception):
output[msg.task_id] = TaskResult(
status=TaskStatus.FAILED, error=str(result)
)
else:
output[msg.task_id] = result
return output
关键设计决策:Orchestrator 不持有 LLM 调用逻辑,只负责调度与聚合,保持单一职责。Agent 的 _process 方法由子类实现,可对接不同 LLM 后端。parallel_dispatch 使用 asyncio.gather 实现并行,但通过 return_exceptions=True 确保单个 Agent 异常不会中断整体流程。
四、多 Agent 系统的工程代价与适用边界
Token 消耗倍增。每个 Agent 独立持有上下文,3 个 Agent 的总 Token 消耗约为单 Agent 的 2.5-3 倍(含通信开销)。通过 A/B 测试对比,在代码生成任务中,多 Agent 方案的单次任务成本约为单 Agent 的 2.8 倍,但首次通过率从 61% 提升到 83%。成本与质量需要根据业务场景权衡。
通信延迟叠加。Agent 间的消息传递引入额外延迟,层级调度模式中,主控 Agent 的串行决策成为瓶颈。实测数据:3 层调度的端到端延迟约为单 Agent 的 1.7 倍。
一致性风险。多个 Agent 并行产出时,结果可能存在逻辑冲突。例如编码 Agent 生成的接口与测试 Agent 假设的接口不一致。需要在聚合阶段增加一致性校验层,但这又引入了额外的复杂度与延迟。
适用边界:任务步骤少于 3 步且无并行需求时,单 Agent 更高效;Agent 数量超过 7 个时,编排复杂度急剧上升,建议引入分层编排(Orchestrator of Orchestrators)。
五、总结
多 Agent 协作系统的核心价值在于"专业化分工"与"并行加速",但代价是通信开销、Token 消耗与一致性管理的复杂度上升。架构选型应基于任务特征:线性依赖选管道模式,可并行选层级调度,松耦合选事件驱动。
落地路线建议:第一步,从单 Agent 出发,明确任务分解点;第二步,实现双 Agent 管道(规划+执行),验证协作可行性;第三步,引入主控 Agent 与并行调度,处理复杂任务;第四步,补齐任务追踪、结果校验与异常恢复机制。每一步都必须有可量化的质量指标(如首次通过率、端到端延迟),用数据驱动架构演进。
更多推荐
所有评论(0)