ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

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

2026/8/3 11:14:19 拓冰建站 浏览量
多 Agent 系统通信的实现原理与最佳实践

1. 引言

随着大语言模型能力的快速提升,单 Agent 已经难以覆盖复杂业务场景。多 Agent 系统通过将任务拆解给多个具备不同专长的 Agent 协作完成,能够显著提升系统的可扩展性、鲁棒性和任务完成质量。而这一切的基础,正是 Agent 之间的通信机制。

本文将从通信模型、消息协议、同步与异步、路由与编排、容错与安全等维度,系统讲解多 Agent 系统通信的实现原理,并给出基于 Python 的完整代码实战。

2. 多 Agent 通信的核心模型

多 Agent 系统的通信模型决定了 Agent 之间如何发现彼此、如何传递消息、如何协调任务。常见的通信模型有以下四种。

通信模型特点适用场景
点对点(Peer-to-Peer)Agent 之间直接通信,延迟低,耦合高小规模、固定拓扑
中心化(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_factory=lambda: 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(f"Agent {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(), timeout=timeout) 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( sender=self.agent_id, receiver=message.sender, msg_type="task_result", payload={"status": "ok", "plan": ["step1", "step2"]}, correlation_id=message.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( sender="agent_orchestrator", receiver="agent_planner", msg_type="task_assign", payload={"task": "制定发布计划"}, )) 等待 planner 回传结果 result = await queue.receive("agent_orchestrator", timeout=3.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(f"Python 任务路由到: {target}") # agent_coder

在编排层面,常见模式包括:

  • 顺序编排:Agent 按固定顺序依次执行,前一个的输出作为后一个的输入。
  • 并行编排:多个独立任务同时分发给多个 Agent,最后汇总结果。
  • 条件编排:根据中间结果动态决定后续执行路径。
  • 递归编排:Agent 发现任务过大时,自行拆解并分发给子 Agent。

6. 容错与重试机制

分布式环境下,Agent 可能崩溃、超时或返回错误结果。健壮的通信层必须提供以下保障。

  • 超时控制:为每次请求设置超时时间,避免无限等待。
  • 重试与退避:对可重试的失败(如网络抖动)进行指数退避重试。
  • 幂等处理:通过消息 ID 去重,确保重复投递不会产生副作用。
  • 死信队列:多次重试仍失败的消息进入死信队列,供人工排查。
  • 心跳检测:定期检测 Agent 存活状态,及时摘除失联节点。

下面给出一个带超时和重试的请求-响应封装。

import asyncio import random async def send_with_retry(queue, message, max_retries=3, base_timeout=2.0): """带指数退避重试的消息发送。""" for attempt in range(max_retries): try: await queue.send(message) result = await queue.receive(message.sender, timeout=base_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_keys=True, 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_factory=lambda: 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(f"Agent {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(), timeout=timeout) 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( sender=self.agent_id, receiver=original.sender, msg_type=msg_type, payload=payload, correlation_id=original.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( sender="orchestrator", receiver=analyst, msg_type="analyze", payload={"requirement": requirement}, )) spec_msg = await self.queue.receive("orchestrator", timeout=3.0) spec = spec_msg.payload["spec"] 2. 路由到代码 Agent coder = self.route("coding") await self.queue.send(Message( sender="orchestrator", receiver=coder, msg_type="code", payload={"spec": spec}, )) code_msg = await self.queue.receive("orchestrator", timeout=3.0) code = code_msg.payload["code"] 3. 路由到审查 Agent reviewer = self.route("review") await self.queue.send(Message( sender="orchestrator", receiver=reviewer, msg_type="review", payload={"code": code}, )) review_msg = await self.queue.receive("orchestrator", timeout=3.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 协作系统。