12.4 多 Agent 协作系统
前几节我们构建的 Agent 都是"单兵作战"——一个 Agent 包揽意图识别、知识检索、回复生成。当任务简单时这没问题,但当任务变复杂——比如要同时做代码生成、代码审查、测试编写、文档输出——单个 Agent 的上下文会爆炸,能力边界会被撑破。这时候就需要多个 Agent 分工协作。
多 Agent 协作系统的核心思想是分而治之:把大任务拆成小任务,每个 Agent 专精一件事,通过消息传递协调配合。就像一家餐厅不会让一个人同时当厨师、服务员、收银员,而是各司其职、流水线作业。本节我们就来拆解多 Agent 协作系统的架构设计、通信协议、任务分配和协作模式。
生活类比:建筑工程队
想象一栋大楼的施工过程。总包经理(Manager Agent)拿到图纸后,把任务拆成水电、木工、泥瓦、装修等子任务,分给不同的施工队(Worker Agent)。各施工队各干各的,但需要通过 walkie-talkie 对讲机(MessageBus 消息总线)协调——水电队要告诉泥瓦队"管道走完了你可以铺砖了",泥瓦队要告诉装修队"墙面找平了你可以刷漆了"。如果某个施工队出了问题,总包会重新分配任务或叫备用队顶上。
多 Agent 协作系统就是这支"数字工程队"。它的核心组件有四个:组织架构(谁指挥谁)、通信协议(怎么传消息)、任务分配(活儿怎么派)、协作模式(怎么配合)。下面我们逐个拆解。
一、多 Agent 协作的组织架构
多 Agent 系统有三种经典的组织架构模式,选择哪种取决于任务特点:
1. 层级式 (Hierarchical) 2. 对等式 (Peer-to-Peer) 3. 混合式 (Hybrid)
Manager Agent Agent A ←-> Agent B Manager Agent
/ | \ ↕ ↕ / \
Worker Worker Worker Agent C -> Agent D Worker A ←-> Worker B层级式:有一个"经理"Agent 负责分配任务,多个"工人"Agent 执行具体工作。适合任务可以清晰拆分的场景,比如代码审查链(生成 -> 审查 -> 修复 -> 测试)。优点是结构清晰、易于管理,缺点是 Manager 是单点瓶颈。
对等式:所有 Agent 地位平等,可以直接互相通信。适合需要自由协作的场景,如头脑风暴、多轮辩论。优点是灵活、无单点瓶颈,缺点是协调成本高、容易混乱。
混合式:Manager 负责任务分配,但 Worker 之间也可以直接通信。兼顾了层级式的清晰和对等式的灵活,是生产环境中最常用的模式。
下面的代码定义了 Agent 的角色、消息类型和能力描述,是整个协作系统的基础数据结构:
from abc import ABC, abstractmethod # 抽象基类
from typing import List, Dict, Any, Optional, Callable # 类型注解
from dataclasses import dataclass, field # 数据类
from enum import Enum # 枚举类型
import asyncio # 异步IO
import uuid # 唯一ID生成
import json # JSON处理
class AgentRole(Enum): # Agent角色枚举
MANAGER = "manager" # 管理者:分配协调
WORKER = "worker" # 工作者:执行任务
OBSERVER = "observer" # 观察者:监控不参与
COORDINATOR = "coordinator" # 协调者:居中调度
class MessageType(Enum): # 消息类型枚举
TASK_ASSIGNMENT = "task_assignment" # 任务分配
TASK_RESULT = "task_result" # 任务结果
QUERY = "query" # 查询
RESPONSE = "response" # 响应
BROADCAST = "broadcast" # 广播
HEARTBEAT = "heartbeat" # 心跳
ERROR = "error" # 错误
@dataclass
class AgentMessage: # Agent间通信消息
"""Agent间通信消息""" # 封装一次通信的内容
msg_id: str = field(default_factory=lambda: str(uuid.uuid4())) # 唯一消息ID
sender_id: str = "" # 发送者ID
receiver_id: str = "" # 接收者ID,空串表示广播
msg_type: MessageType = MessageType.QUERY # 消息类型
content: Any = None # 消息内容
metadata: Dict = field(default_factory=dict) # 元数据
timestamp: float = field(default_factory=lambda: __import__('time').time()) # 时间戳
reply_to: Optional[str] = None # 回复的消息ID
@dataclass
class AgentCapability: # Agent能力描述
"""Agent能力描述""" # 声明Agent能做什么
name: str # 能力名称
description: str # 能力描述
skills: List[str] # 技能列表
tools: List[str] # 可用工具
max_concurrent_tasks: int = 3 # 最大并发任务数
priority: int = 1 # 优先级二、多 Agent 通信协议
组织架构定义了"谁和谁说话",通信协议定义了"怎么说话"。MessageBus 消息总线是整个通信系统的核心——它就像工程队的对讲机网络,每个 Agent 注册一个频道(Queue),发消息时指定接收者或广播。
消息总线支持两种通信方式:点对点(指定 receiver_id,只有目标 Agent 收到)和广播(receiver_id 为空,所有 Agent 都收到)。还维护了消息历史,方便调试和审计。
class MessageBus: # 消息总线
"""消息总线 - Agent间通信的核心""" # 所有Agent通过它收发消息
def __init__(self): # 初始化
self.subscribers: Dict[str, asyncio.Queue] = {} # 订阅者:AgentID->消息队列
self.message_history: List[AgentMessage] = [] # 消息历史
self.max_history = 1000 # 最多保留1000条历史
def register(self, agent_id: str) -> asyncio.Queue:
"""注册Agent到消息总线""" # Agent启动时调用
queue = asyncio.Queue() # 创建专属消息队列
self.subscribers[agent_id] = queue # 存入订阅表
return queue # 返回队列供Agent轮询
def unregister(self, agent_id: str):
"""注销Agent""" # Agent停止时调用
self.subscribers.pop(agent_id, None) # 从订阅表移除
async def send(self, message: AgentMessage):
"""发送消息""" # 核心方法:投递消息
self.message_history.append(message) # 记入历史
if len(self.message_history) > self.max_history: # 历史超限
self.message_history = self.message_history[-self.max_history:] # 只保留最近N条
if message.receiver_id: # 点对点消息
if message.receiver_id in self.subscribers: # 接收者存在
await self.subscribers[message.receiver_id].put(message) # 投递到接收者队列
else: # 广播消息
for agent_id, queue in self.subscribers.items(): # 遍历所有订阅者
if agent_id != message.sender_id: # 不发给自己
await queue.put(message) # 投递到每个队列
async def receive(self, agent_id: str, timeout: float = None) -> Optional[AgentMessage]:
"""接收消息""" # Agent从自己的队列取消息
queue = self.subscribers.get(agent_id) # 取出Agent的队列
if not queue: # 没注册过
return None # 返回None
try:
if timeout: # 设置了超时
return await asyncio.wait_for(queue.get(), timeout=timeout) # 带超时等待
return await queue.get() # 无限等待
except asyncio.TimeoutError: # 超时了
return None # 返回None有了消息总线,还需要一个 BaseAgent 基类来封装"注册 -> 收消息 -> 处理消息"的通用逻辑。每个具体的 Agent 只需实现 handle_message 方法,专注于业务逻辑:
class BaseAgent(ABC): # Agent抽象基类
"""Agent基类""" # 所有Agent的公共逻辑
def __init__(self, agent_id: str, role: AgentRole, message_bus: MessageBus): # 初始化
self.agent_id = agent_id # Agent唯一ID
self.role = role # 角色类型
self.bus = message_bus # 消息总线引用
self.capabilities: List[AgentCapability] = [] # 能力列表
self.state: Dict = {} # 状态存储
self.running = False # 运行标志
self.inbox: asyncio.Queue = None # 消息收件箱
async def start(self): # 启动Agent
"""启动Agent""" # 注册到总线并启动消息循环
self.inbox = self.bus.register(self.agent_id) # 注册获取消息队列
self.running = True # 标记为运行中
asyncio.create_task(self._message_loop()) # 启动消息处理协程
async def stop(self): # 停止Agent
"""停止Agent""" # 关闭并注销
self.running = False # 停止循环
self.bus.unregister(self.agent_id) # 从总线注销
async def _message_loop(self): # 消息处理循环
"""消息处理循环""" # 不断取消息并处理
while self.running: # 运行中
message = await self.bus.receive(self.agent_id, timeout=1.0) # 取消息
if message: # 有消息
await self.handle_message(message) # 交给子类处理
@abstractmethod
async def handle_message(self, message: AgentMessage): # 抽象方法
"""处理收到的消息""" # 子类必须实现
pass
async def send_message(self, receiver_id: str, msg_type: MessageType,
content: Any, reply_to: str = None) -> str:
"""发送消息""" # 封装并发送
message = AgentMessage( # 构造消息
sender_id=self.agent_id, # 发送者是自己
receiver_id=receiver_id, # 接收者
msg_type=msg_type, # 消息类型
content=content, # 内容
reply_to=reply_to # 回复的目标消息ID
)
await self.bus.send(message) # 通过总线发送
return message.msg_id # 返回消息ID
async def broadcast(self, msg_type: MessageType, content: Any):
"""广播消息""" # 发给所有Agent
message = AgentMessage( # 构造广播消息
sender_id=self.agent_id, # 发送者
receiver_id="", # 空串表示广播
msg_type=msg_type, # 类型
content=content # 内容
)
await self.bus.send(message) # 发送三、任务分配与协调
Manager Agent 是层级式协作的大脑。它维护所有 Worker 的信息(能力、状态、当前负载),根据任务需求匹配最合适的 Worker。就像总包经理知道哪个施工队擅长水电、哪个擅长木工,并且知道谁现在有空。
任务分配有三层策略:能力匹配(Worker 的技能要覆盖任务需求)、负载均衡(优先选当前任务最少的 Worker)、排队等待(所有合适的 Worker 都忙时,任务进队列等待)。下面的代码完整实现了这套逻辑:
class ManagerAgent(BaseAgent): # 管理者Agent
"""管理者Agent - 负责任务分配与协调""" # 层级式协作的核心
def __init__(self, agent_id: str, message_bus: MessageBus): # 初始化
super().__init__(agent_id, AgentRole.MANAGER, message_bus) # 调用父类
self.workers: Dict[str, WorkerInfo] = {} # Worker信息表
self.tasks: Dict[str, Task] = {} # 任务表
self.pending_tasks: asyncio.Queue = asyncio.Queue() # 待分配任务队列
self.completed_tasks: Dict[str, Any] = {} # 已完成任务
async def register_worker(self, worker_id: str, capabilities: List[AgentCapability]):
"""注册Worker""" # Worker上线时注册
self.workers[worker_id] = WorkerInfo( # 创建Worker信息
worker_id=worker_id, # Worker ID
capabilities=capabilities, # 能力列表
status="idle", # 初始空闲
current_tasks=[] # 当前无任务
)
async def assign_task(self, task: 'Task') -> str: # 分配任务
"""分配任务给最合适的Worker""" # 三层策略
suitable_workers = [] # 收集符合条件的Worker
for worker_id, worker_info in self.workers.items(): # 遍历所有Worker
if worker_info.status == "idle" and self._capability_match(task, worker_info): # 空闲且能力匹配
suitable_workers.append(worker_info) # 加入候选列表
if not suitable_workers: # 没有合适的Worker
await self.pending_tasks.put(task) # 任务进等待队列
return None # 返回None表示暂未分配
best_worker = min(suitable_workers, key=lambda w: len(w.current_tasks)) # 选负载最低的
task.assigned_to = best_worker.worker_id # 记录分配给谁
task.status = TaskStatus.ASSIGNED # 更新任务状态
self.tasks[task.task_id] = task # 存入任务表
await self.send_message( # 发送任务分配消息
receiver_id=best_worker.worker_id, # 接收者是选中的Worker
msg_type=MessageType.TASK_ASSIGNMENT, # 类型:任务分配
content={ # 任务详情
"task_id": task.task_id, # 任务ID
"description": task.description, # 任务描述
"input_data": task.input_data, # 输入数据
"deadline": task.deadline # 截止时间
}
)
best_worker.current_tasks.append(task.task_id) # Worker任务列表+1
best_worker.status = "busy" # 标记为忙碌
return task.task_id # 返回任务ID
def _capability_match(self, task: 'Task', worker: 'WorkerInfo') -> bool:
"""检查Worker能力是否匹配任务""" # 技能交集检查
for cap in worker.capabilities: # 遍历Worker的每项能力
if any(skill in task.required_skills for skill in cap.skills): # 有技能在需求列表中
return True # 匹配成功
return False # 没有任何匹配
async def handle_message(self, message: AgentMessage): # 处理消息
"""处理消息""" # Manager主要处理结果和错误
if message.msg_type == MessageType.TASK_RESULT: # 收到任务结果
content = message.content # 取出内容
task_id = content["task_id"] # 取任务ID
self.completed_tasks[task_id] = content["result"] # 存结果
worker = self.workers.get(message.sender_id) # 取Worker信息
if worker: # Worker存在
worker.current_tasks.remove(task_id) # 从任务列表移除
worker.status = "idle" if not worker.current_tasks else "busy" # 更新状态
await self._process_pending_tasks() # 处理等待队列中的任务
if task_id in self.tasks and self.tasks[task_id].callback: # 有回调
await self.tasks[task_id].callback(content["result"]) # 执行回调
elif message.msg_type == MessageType.ERROR: # 收到错误
task_id = message.content.get("task_id") # 取任务ID
if task_id and task_id in self.tasks: # 任务存在
task = self.tasks[task_id] # 取任务
task.retry_count += 1 # 重试计数+1
if task.retry_count < task.max_retries: # 未超最大重试
await self.assign_task(task) # 重新分配
class WorkerAgent(BaseAgent): # 工作者Agent
"""工作者Agent - 执行具体任务""" # 接收任务并用ReAct执行
def __init__(self, agent_id: str, message_bus: MessageBus,
llm_client, tools: List[Callable] = None): # 初始化
super().__init__(agent_id, AgentRole.WORKER, message_bus) # 调用父类
self.llm = llm_client # LLM客户端
self.tools = tools or [] # 可用工具列表
async def handle_message(self, message: AgentMessage): # 处理消息
"""处理消息""" # Worker处理任务分配和查询
if message.msg_type == MessageType.TASK_ASSIGNMENT: # 收到任务
await self._execute_task(message) # 执行任务
elif message.msg_type == MessageType.QUERY: # 收到查询
await self._handle_query(message) # 处理查询
async def _execute_task(self, message: AgentMessage): # 执行任务
"""执行任务""" # 用ReAct循环执行并返回结果
content = message.content # 取任务内容
task_id = content["task_id"] # 取任务ID
try:
result = await self._react_loop( # ReAct循环执行
task=content["description"], # 任务描述
input_data=content["input_data"] # 输入数据
)
await self.send_message( # 发送成功结果
receiver_id=message.sender_id, # 回给发送者
msg_type=MessageType.TASK_RESULT, # 类型:任务结果
content={"task_id": task_id, "result": result}, # 内容
reply_to=message.msg_id # 回复的消息ID
)
except Exception as e: # 执行出错
await self.send_message( # 发送错误
receiver_id=message.sender_id,
msg_type=MessageType.ERROR,
content={"task_id": task_id, "error": str(e)},
reply_to=message.msg_id
)
async def _react_loop(self, task: str, input_data: Any, max_steps: int = 10) -> Any:
"""ReAct循环执行任务""" # Thought->Action->Observation循环
context = { # 上下文
"task": task, # 任务描述
"input": input_data, # 输入数据
"observations": [], # 观察记录
"steps_completed": 0 # 已完成步数
}
for step in range(max_steps): # 最多max_steps步
thought = await self._think(context) # 思考下一步
if thought.get("action") == "FINISH": # 判断是否完成
return thought.get("result") # 返回最终结果
action_result = await self._execute_action(thought, context) # 执行动作
context["observations"].append(action_result) # 记录观察
context["steps_completed"] = step + 1 # 更新步数
return {"error": "达到最大步数限制", "context": context} # 超限返回错误四、多 Agent 协作模式
有了通信和任务分配的基础设施,接下来看 Agent 之间怎么配合干活。常见的协作模式有四种:顺序流水线(A -> B -> C 串行)、并行处理(多个 Agent 同时干)、辩论(互相说服)、投票(少数服从多数)。就像工程队里可以串行施工、并行赶工、开会争论方案、投票表决。
class CollaborationPattern: # 协作模式定义
"""协作模式定义""" # 提供多种协作策略
@staticmethod
async def sequential(agents: List[WorkerAgent], task: str, input_data: Any) -> Any:
"""顺序协作:A->B->C 流水线""" # 前一个的输出是后一个的输入
result = input_data # 初始输入
for agent in agents: # 依次调用每个Agent
result = await agent._react_loop(task, result) # 上一个的输出作为下一个的输入
return result # 返回最终结果
@staticmethod
async def parallel(agents: List[WorkerAgent], task: str, input_data: Any) -> List[Any]:
"""并行协作:多个Agent同时处理""" # 所有Agent同时工作,各自独立
tasks = [agent._react_loop(task, input_data) for agent in agents] # 构造协程列表
return await asyncio.gather(*tasks) # 并发执行并收集结果
@staticmethod
async def debate(agents: List[WorkerAgent], topic: str, rounds: int = 3) -> Dict:
"""辩论模式:Agent之间互相辩论""" # 多轮交替发表观点
opinions = [] # 收集所有观点
for round_num in range(rounds): # 进行rounds轮辩论
round_opinions = [] # 本轮观点
for agent in agents: # 每个Agent轮流发言
context = { # 辩论上下文
"topic": topic, # 辩题
"round": round_num + 1, # 当前轮次
"previous_opinions": opinions # 之前的所有观点
}
opinion = await agent._react_loop(f"就'{topic}'发表观点并回应其他观点", context) # 让Agent发言
round_opinions.append(opinion) # 存入本轮观点
opinions.extend(round_opinions) # 合并到总观点列表
return {"topic": topic, "opinions": opinions} # 返回辩论记录
@staticmethod
async def vote_and_execute(agents: List[WorkerAgent], task: str, input_data: Any) -> Any:
"""投票执行:多个方案投票选择最佳方案""" # 先各自生成方案再投票
proposals = await CollaborationPattern.parallel(agents, f"为'{task}'生成解决方案", input_data) # 各自生成方案
votes = [0] * len(proposals) # 初始化投票计数
for agent in agents: # 每个Agent投票
vote_prompt = f"请从以下方案中选择最佳方案(返回编号0-{len(proposals)-1}):\n"
for i, p in enumerate(proposals): # 列出所有方案
vote_prompt += f"\n方案{i}:{p}\n"
vote = await agent.llm.chat(vote_prompt) # Agent投票
try:
voted_idx = int(vote.strip()) # 解析投票编号
votes[voted_idx] += 1 # 计票
except: # 解析失败
pass # 票作废
best_idx = votes.index(max(votes)) # 得票最高的方案编号
return proposals[best_idx] # 返回最佳方案debate 模式特别适合需要多角度思考的场景——比如技术方案评估,激进派推动方案 A、保守派倾向方案 B,多轮辩论后往往能收敛出更全面的决策。vote_and_execute 模式适合需要减少单一 Agent 偏见的场景——多个 Agent 独立生成方案再投票,类似"群体智慧"。
最后是任务和 Worker 的数据结构定义:
@dataclass
class Task: # 任务定义
"""任务""" # 一个待执行的单元
task_id: str = field(default_factory=lambda: str(uuid.uuid4())) # 唯一ID
description: str = "" # 任务描述
required_skills: List[str] = field(default_factory=list) # 所需技能
input_data: Any = None # 输入数据
assigned_to: Optional[str] = None # 分配给哪个Worker
status: str = "pending" # 状态
deadline: Optional[float] = None # 截止时间
max_retries: int = 3 # 最大重试次数
retry_count: int = 0 # 当前重试次数
callback: Optional[Callable] = None # 完成回调
class TaskStatus: # 任务状态常量
PENDING = "pending" # 等待分配
ASSIGNED = "assigned" # 已分配
RUNNING = "running" # 执行中
COMPLETED = "completed" # 已完成
FAILED = "failed" # 失败
@dataclass
class WorkerInfo: # Worker信息
"""Worker信息""" # Manager维护的Worker状态
worker_id: str # Worker ID
capabilities: List[AgentCapability] # 能力列表
status: str = "idle" # 当前状态
current_tasks: List[str] = field(default_factory=list) # 当前任务列表五、常见误区
误区一:所有任务都用多 Agent。多 Agent 带来了通信开销和协调成本,如果任务本身简单(比如单轮问答),用多 Agent 纯属浪费。判断标准是:任务是否需要不同专业能力、是否能并行拆分。如果不能,就用单 Agent。
误区二:Manager 成为单点瓶颈。所有任务都经过 Manager 分配,如果 Manager 处理慢或挂了,整个系统就卡住了。解决方案是给 Manager 做主备冗余,或者用分布式锁让多个 Manager 协调分配。
误区三:消息队列无限增长。
MessageBus的message_history如果不限制大小,长时间运行会吃光内存。代码中用max_history = 1000做了截断,生产环境还应考虑持久化到数据库。误区四:辩论模式不加轮数限制。辩论如果没有轮数上限,两个 Agent 可能无限争论下去。
debate方法默认 3 轮,这是合理的——再多就是浪费 Token。误区五:Worker 失败不重试就放弃。网络抖动、API 限流等临时故障很常见,直接放弃会导致用户体验差。代码中
max_retries = 3是个合理的默认值,配合指数退避效果更好。
六、本节小结
| 要点 | 说明 |
|---|---|
| 组织架构 | 层级式、对等式、混合式三种模式,根据任务特点选择 |
| 消息总线 | Agent 间通信的核心基础设施,支持点对点和广播 |
| 任务分配 | 能力匹配 -> 负载均衡 -> 排队等待,三层策略保证合理分配 |
| 协作模式 | 顺序流水线、并行处理、多轮辩论、投票表决,四种经典模式 |
| 容错机制 | 任务重试、Worker 状态监控、错误自动重新分配 |
多 Agent 协作系统是复杂 Agent 应用的必经之路。当单个 Agent 的上下文窗口撑不住、能力边界覆盖不了时,把任务拆开分给多个专业 Agent 是自然的选择。掌握了消息总线、任务分配和协作模式后,你可以搭建代码审查流水线、多模型辩论系统、分布式任务处理系统等各种复杂场景。下一节我们讨论 Agent 系统的稳定性与可靠性——多 Agent 不仅要多,还要稳。