多 Agent 系统通信的实现原理与最佳实践

多 Agent 系统通信的实现原理与最佳实践 1. 引言随着大语言模型能力的快速提升单 Agent 已经难以覆盖复杂业务场景。多 Agent 系统通过将任务拆解给多个具备不同专长的 Agent 协作完成能够显著提升系统的可扩展性、鲁棒性和任务完成质量。而这一切的基础正是 Agent 之间的通信机制。本文将从通信模型、消息协议、同步与异步、路由与编排、容错与安全等维度系统讲解多 Agent 系统通信的实现原理并给出基于 Python 的完整代码实战。2. 多 Agent 通信的核心模型多 Agent 系统的通信模型决定了 Agent 之间如何发现彼此、如何传递消息、如何协调任务。常见的通信模型有以下四种。通信模型特点适用场景点对点Peer-to-PeerAgent 之间直接通信延迟低耦合高小规模、固定拓扑中心化Hub-and-Spoke所有消息经中心协调器转发易管理有单点风险任务编排、权限控制严格消息总线Message Bus通过发布/订阅解耦生产者和消费者扩展性好事件驱动、大规模系统黑板系统Blackboard共享工作区Agent 读写公共状态适合协作求解复杂问题分解、多专家协作在实际工程中中心化编排配合消息总线是最常见的组合编排器负责任务分解和结果汇总Agent 之间通过总线异步通信。3. 消息协议设计通信协议是 Agent 之间约定的消息格式。一个健壮的消息协议应当包含以下核心字段。{ message_id: msg_8f3a2c1e, sender: agent_planner, receiver: agent_coder, type: task_assign, timestamp: 2026-08-03T09:30:00Z, correlation_id: task_42, payload: { task: 实现用户登录接口, requirements: [支持 JWT, 包含单元测试], deadline: 2026-08-03T12:00:00Z }, metadata: { priority: high, retry_count: 0 } }设计消息协议时应重点关注以下几点消息 ID 与关联 ID用于幂等处理和请求-响应关联。发送方与接收方支持点对点路由也支持广播receiver 为通配符。消息类型区分任务分配、结果回传、状态查询、错误上报等。时间戳用于超时判断和消息排序。版本号协议演进时保证向后兼容。4. 同步通信与异步通信同步通信中调用方阻塞等待被调用方返回结果实现简单但吞吐低异步通信中调用方发送消息后立即返回通过回调、轮询或事件驱动获取结果吞吐高但复杂度上升。下面给出一个基于 Pythonasyncio的异步消息队列实现演示 Agent 之间如何通过队列解耦通信。import asyncio import uuid from dataclasses import dataclass, field from typing import Dict, Optional dataclass class Message: sender: str receiver: str msg_type: str payload: dict message_id: str field(default_factorylambda: uuid.uuid4().hex) correlation_id: Optional[str] None class MessageQueue: 基于 asyncio.Queue 的轻量级消息队列支持点对点和广播。 def __init__(self): self._queues: Dict[str, asyncio.Queue] {} self._lock asyncio.Lock() async def register(self, agent_id: str) - None: async with self._lock: if agent_id not in self._queues: self._queues[agent_id] asyncio.Queue() async def send(self, message: Message) - None: 发送消息receiver 为 时广播否则点对点投递。 if message.receiver : for queue in self._queues.values(): await queue.put(message) else: if message.receiver not in self._queues: raise ValueError(fAgent {message.receiver} 未注册) await self._queues[message.receiver].put(message) async def receive(self, agent_id: str, timeout: float 5.0) - Optional[Message]: queue self._queues.get(agent_id) if queue is None: return None try: return await asyncio.wait_for(queue.get(), timeouttimeout) except asyncio.TimeoutError: return None class Agent: def init(self, agent_id: str, queue: MessageQueue): self.agent_id agent_id self.queue queue async def start(self) - None: await self.queue.register(self.agent_id) while True: message await self.queue.receive(self.agent_id) if message is None: continue await self.handle(message) async def handle(self, message: Message) - None: 子类重写此方法处理消息。 raise NotImplementedError class PlannerAgent(Agent): async def handle(self, message: Message) - None: if message.msg_type task_assign: print(f[Planner] 收到任务: {message.payload[task]}) 模拟任务分解 await asyncio.sleep(0.1) reply Message( senderself.agent_id, receivermessage.sender, msg_typetask_result, payload{status: ok, plan: [step1, step2]}, correlation_idmessage.message_id, ) await self.queue.send(reply) async def main(): queue MessageQueue() planner PlannerAgent(agent_planner, queue) task asyncio.create_task(planner.start()) await queue.register(agent_orchestrator) await queue.send(Message( senderagent_orchestrator, receiveragent_planner, msg_typetask_assign, payload{task: 制定发布计划}, )) 等待 planner 回传结果 result await queue.receive(agent_orchestrator, timeout3.0) if result: print(f[Orchestrator] 收到结果: {result.payload}) task.cancel() if name main: asyncio.run(main())上述代码展示了三个关键设计Agent 启动时注册自己的队列发送方通过receiver字段路由消息通过correlation_id关联请求与响应。5. 消息路由与任务编排在复杂系统中消息需要经过路由层转发到正确的 Agent。路由策略包括基于内容的路由根据消息 payload 中的字段如任务类型决定目标 Agent。基于能力注册的路由Agent 启动时声明自身能力路由层维护能力到 Agent 的映射。基于负载的路由将消息分发给当前负载最低的 Agent实现负载均衡。下面给出一个基于能力注册的路由器实现。from typing import Dict, List, Optional class CapabilityRouter: 根据 Agent 声明的能力进行消息路由。 def __init__(self): self._capabilities: Dict[str, List[str]] {} def register(self, agent_id: str, capabilities: List[str]) - None: self._capabilities[agent_id] capabilities def route(self, required_capability: str) - Optional[str]: 返回具备指定能力的第一个 Agent无匹配时返回 None。 for agent_id, caps in self._capabilities.items(): if required_capability in caps: return agent_id return None 使用示例 router CapabilityRouter() router.register(agent_coder, [python, java]) router.register(agent_reviewer, [code_review, security]) target router.route(python) print(fPython 任务路由到: {target}) # agent_coder在编排层面常见模式包括顺序编排Agent 按固定顺序依次执行前一个的输出作为后一个的输入。并行编排多个独立任务同时分发给多个 Agent最后汇总结果。条件编排根据中间结果动态决定后续执行路径。递归编排Agent 发现任务过大时自行拆解并分发给子 Agent。6. 容错与重试机制分布式环境下Agent 可能崩溃、超时或返回错误结果。健壮的通信层必须提供以下保障。超时控制为每次请求设置超时时间避免无限等待。重试与退避对可重试的失败如网络抖动进行指数退避重试。幂等处理通过消息 ID 去重确保重复投递不会产生副作用。死信队列多次重试仍失败的消息进入死信队列供人工排查。心跳检测定期检测 Agent 存活状态及时摘除失联节点。下面给出一个带超时和重试的请求-响应封装。import asyncio import random async def send_with_retry(queue, message, max_retries3, base_timeout2.0): 带指数退避重试的消息发送。 for attempt in range(max_retries): try: await queue.send(message) result await queue.receive(message.sender, timeoutbase_timeout) if result is not None: return result except asyncio.TimeoutError: pass # 指数退避2s, 4s, 8s wait_time base_timeout * (2 ** attempt) random.uniform(0, 0.5) print(f第 {attempt 1} 次重试等待 {wait_time:.2f}s) await asyncio.sleep(wait_time) raise TimeoutError(f消息 {message.message_id} 重试 {max_retries} 次仍失败)/code/pre 7. 安全与权限控制 多 Agent 系统通信面临身份伪造、消息篡改、越权访问等安全风险。最佳实践包括 身份认证每个 Agent 使用独立的 API Key 或 JWT 进行身份认证。 消息签名对消息体进行 HMAC 签名防止传输过程中被篡改。 最小权限每个 Agent 只授予完成任务所需的最小权限。 敏感信息脱敏日志和消息中避免明文传输密钥、Token 等敏感信息。 审计日志记录所有跨 Agent 通信的关键信息便于追溯。 下面给出一个基于 HMAC 的消息签名示例。 import hashlib import hmac import json def sign_message(payload: dict, secret: str) - str: 对消息 payload 计算 HMAC-SHA256 签名。 body json.dumps(payload, sort_keysTrue, separators(,, :)) return hmac.new(secret.encode(), body.encode(), hashlib.sha256).hexdigest() def verify_message(payload: dict, signature: str, secret: str) - bool: 校验消息签名是否合法。 expected sign_message(payload, secret) return hmac.compare_digest(expected, signature) 使用示例 SECRET my_shared_secret msg_payload {task: deploy, target: prod} sig sign_message(msg_payload, SECRET) print(f签名: {sig}) print(f校验通过: {verify_message(msg_payload, sig, SECRET)}) 8. 实战构建一个完整的多 Agent 协作系统 下面综合前面所有知识点构建一个「需求分析 → 代码生成 → 代码审查」的三 Agent 协作系统。系统使用中心化编排器 异步消息队列并加入超时重试与能力路由。 import asyncio import uuid from dataclasses import dataclass, field from typing import Dict, List, Optional ---------- 消息层 ---------- dataclass class Message: sender: str receiver: str msg_type: str payload: dict message_id: str field(default_factorylambda: uuid.uuid4().hex) correlation_id: Optional[str] None class MessageQueue: def init(self): self._queues: Dict[str, asyncio.Queue] {} self._lock asyncio.Lock() async def register(self, agent_id: str) - None: async with self._lock: self._queues.setdefault(agent_id, asyncio.Queue()) async def send(self, message: Message) - None: if message.receiver *: for q in self._queues.values(): await q.put(message) else: if message.receiver not in self._queues: raise ValueError(fAgent {message.receiver} 未注册) await self._queues[message.receiver].put(message) async def receive(self, agent_id: str, timeout: float 5.0) - Optional[Message]: q self._queues.get(agent_id) if q is None: return None try: return await asyncio.wait_for(q.get(), timeouttimeout) except asyncio.TimeoutError: return None ---------- Agent 基类 ---------- class Agent: def init(self, agent_id: str, queue: MessageQueue, capabilities: List[str]): self.agent_id agent_id self.queue queue self.capabilities capabilities async def start(self) - None: await self.queue.register(self.agent_id) while True: msg await self.queue.receive(self.agent_id) if msg is None: continue await self.handle(msg) async def reply(self, original: Message, payload: dict, msg_type: str task_result) - None: await self.queue.send(Message( senderself.agent_id, receiveroriginal.sender, msg_typemsg_type, payloadpayload, correlation_idoriginal.message_id, )) async def handle(self, message: Message) - None: raise NotImplementedError ---------- 具体 Agent ---------- class AnalystAgent(Agent): 需求分析 Agent async def handle(self, message: Message) - None: if message.msg_type analyze: req message.payload[requirement] print(f[Analyst] 分析需求: {req}) await asyncio.sleep(0.2) await self.reply(message, { status: ok, spec: f需求「{req}」已拆解为 3 个功能点, }) class CoderAgent(Agent): 代码生成 Agent async def handle(self, message: Message) - None: if message.msg_type code: spec message.payload[spec] print(f[Coder] 根据规格生成代码: {spec}) await asyncio.sleep(0.3) await self.reply(message, { status: ok, code: def hello():\n return Hello Multi-Agent, }) class ReviewerAgent(Agent): 代码审查 Agent async def handle(self, message: Message) - None: if message.msg_type review: code message.payload[code] print(f[Reviewer] 审查代码: {code}) await asyncio.sleep(0.2) await self.reply(message, { status: ok, verdict: 通过, suggestions: [建议补充类型注解], }) ---------- 编排器 ---------- class Orchestrator: def init(self, queue: MessageQueue): self.queue queue self._capabilities: Dict[str, List[str]] {} def register_agent(self, agent: Agent) - None: self._capabilities[agent.agent_id] agent.capabilities def route(self, capability: str) - Optional[str]: for agent_id, caps in self._capabilities.items(): if capability in caps: return agent_id return None async def run_pipeline(self, requirement: str) - None: # 1. 路由到分析 Agent analyst self.route(analysis) if not analyst: raise RuntimeError(没有可用的分析 Agent) await self.queue.send(Message( senderorchestrator, receiveranalyst, msg_typeanalyze, payload{requirement: requirement}, )) spec_msg await self.queue.receive(orchestrator, timeout3.0) spec spec_msg.payload[spec] 2. 路由到代码 Agent coder self.route(coding) await self.queue.send(Message( senderorchestrator, receivercoder, msg_typecode, payload{spec: spec}, )) code_msg await self.queue.receive(orchestrator, timeout3.0) code code_msg.payload[code] 3. 路由到审查 Agent reviewer self.route(review) await self.queue.send(Message( senderorchestrator, receiverreviewer, msg_typereview, payload{code: code}, )) review_msg await self.queue.receive(orchestrator, timeout3.0) print(\n 最终结果 ) print(f规格: {spec}) print(f代码: {code}) print(f审查: {review_msg.payload[verdict]} - {review_msg.payload[suggestions]}) async def main(): queue MessageQueue() analyst AnalystAgent(agent_analyst, queue, [analysis]) coder CoderAgent(agent_coder, queue, [coding]) reviewer ReviewerAgent(agent_reviewer, queue, [review]) 启动 Agent 后台任务 tasks [ asyncio.create_task(analyst.start()), asyncio.create_task(coder.start()), asyncio.create_task(reviewer.start()), ] orchestrator Orchestrator(queue) orchestrator.register_agent(analyst) orchestrator.register_agent(coder) orchestrator.register_agent(reviewer) await orchestrator.run_pipeline(实现一个用户注册接口) for t in tasks: t.cancel() if name main: asyncio.run(main()) 运行上述代码输出如下 [Analyst] 分析需求: 实现一个用户注册接口 [Coder] 根据规格生成代码: 需求「实现一个用户注册接口」已拆解为 3 个功能点 [Reviewer] 审查代码: def hello(): return Hello Multi-Agent 最终结果 规格: 需求「实现一个用户注册接口」已拆解为 3 个功能点 代码: def hello(): return Hello Multi-Agent 审查: 通过 - [建议补充类型注解] 9. 最佳实践总结 综合以上原理与实战多 Agent 系统通信的最佳实践可以归纳为以下几点 优先异步通信异步消息队列能有效解耦 Agent提升系统吞吐和可扩展性。 协议先行在开发前定义好消息协议包含消息 ID、关联 ID、类型、时间戳和版本号。 能力注册 路由让 Agent 声明能力由路由层动态分发避免硬编码调用关系。 编排器只做协调编排器负责任务分解、路由和结果汇总不参与具体业务计算。 全面考虑容错超时、重试、幂等、死信队列和心跳检测缺一不可。 安全内建身份认证、消息签名、最小权限和审计日志应在设计阶段就纳入。 可观测性为每条消息链路注入 Trace ID便于全链路追踪和问题定位。 10. 结语 多 Agent 系统的通信层是整个协作体系的骨架。选择合理的通信模型、设计健壮的消息协议、实现可靠的路由与容错机制是构建生产级多 Agent 应用的关键。希望本文的原理讲解和代码实战能帮助你快速上手在实际项目中构建出稳定、高效、可扩展的多 Agent 协作系统。