Skip to content

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 的角色、消息类型和能力描述,是整个协作系统的基础数据结构:

python
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 都收到)。还维护了消息历史,方便调试和审计。

python
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 方法,专注于业务逻辑:

python
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 都忙时,任务进队列等待)。下面的代码完整实现了这套逻辑:

python
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 同时干)、辩论(互相说服)、投票(少数服从多数)。就像工程队里可以串行施工、并行赶工、开会争论方案、投票表决。

python
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 的数据结构定义:

python
@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)  # 当前任务列表

五、常见误区

  1. 误区一:所有任务都用多 Agent。多 Agent 带来了通信开销和协调成本,如果任务本身简单(比如单轮问答),用多 Agent 纯属浪费。判断标准是:任务是否需要不同专业能力、是否能并行拆分。如果不能,就用单 Agent。

  2. 误区二:Manager 成为单点瓶颈。所有任务都经过 Manager 分配,如果 Manager 处理慢或挂了,整个系统就卡住了。解决方案是给 Manager 做主备冗余,或者用分布式锁让多个 Manager 协调分配。

  3. 误区三:消息队列无限增长MessageBusmessage_history 如果不限制大小,长时间运行会吃光内存。代码中用 max_history = 1000 做了截断,生产环境还应考虑持久化到数据库。

  4. 误区四:辩论模式不加轮数限制。辩论如果没有轮数上限,两个 Agent 可能无限争论下去。debate 方法默认 3 轮,这是合理的——再多就是浪费 Token。

  5. 误区五:Worker 失败不重试就放弃。网络抖动、API 限流等临时故障很常见,直接放弃会导致用户体验差。代码中 max_retries = 3 是个合理的默认值,配合指数退避效果更好。

六、本节小结

要点说明
组织架构层级式、对等式、混合式三种模式,根据任务特点选择
消息总线Agent 间通信的核心基础设施,支持点对点和广播
任务分配能力匹配 -> 负载均衡 -> 排队等待,三层策略保证合理分配
协作模式顺序流水线、并行处理、多轮辩论、投票表决,四种经典模式
容错机制任务重试、Worker 状态监控、错误自动重新分配

多 Agent 协作系统是复杂 Agent 应用的必经之路。当单个 Agent 的上下文窗口撑不住、能力边界覆盖不了时,把任务拆开分给多个专业 Agent 是自然的选择。掌握了消息总线、任务分配和协作模式后,你可以搭建代码审查流水线、多模型辩论系统、分布式任务处理系统等各种复杂场景。下一节我们讨论 Agent 系统的稳定性与可靠性——多 Agent 不仅要多,还要稳。