你肯定遇到过这样的场景:一个Agent能帮你写代码,另一个Agent能帮你调API,还有一个Agent能帮你分析日志。它们各自都很能干,但当你需要它们接力完成一个复杂任务时——比如先分析需求、再生成代码、最后部署测试——你就得手动在它们之间“传话”,复制粘贴输出,检查格式,处理错误。整个过程笨拙、低效,还容易出错。
这恰恰是当前AI应用从“单点智能”走向“协同工作流”时最普遍的痛点。我们缺的不是强大的Agent,而是让Agent之间能像团队一样顺畅协作的“沟通协议”和“协作框架”。今天要聊的Agent2Agent(A2A),就是为了解决这个问题而生。它不是一个具体的工具,而是一种设计模式或通信范式,核心目标是让不同的AI Agent能够自主、结构化地交换信息、传递任务和协同工作。
很多人第一次听到A2A,会下意识地把它等同于简单的“API调用”或“消息队列”。这其实是一个常见的误解。A2A的挑战和魅力,远不止于技术上的连通性。它真正要解决的是:如何让拥有不同“技能”和“思维模式”的Agent,理解彼此的“意图”,处理不完整的“上下文”,并在协作失败时能进行有效的“协商”或“回退”。
本文将通过一个具体的、可运行的Demo,带你一步步拆解A2A的核心实现逻辑。我们不会停留在概念层面,而是深入到代码和设计决策中,回答三个关键问题:
- Agent之间到底“聊”什么?(消息协议的设计)
- 它们怎么知道该找谁聊?(路由与发现机制)
- 聊崩了怎么办?(错误处理与状态管理)
你会发现,实现一个基础的A2A通信层并不复杂,但要让这个协作网络稳定、可靠、可扩展,里面充满了值得深思的工程细节。
1. 超越简单的函数调用:A2A要解决的核心问题是什么?
在开始写代码之前,我们必须先厘清一个根本问题:既然我们可以用一个超级Agent(比如GPT-4)通过长上下文处理复杂任务,或者用脚本串联多个API,为什么还需要专门的A2A通信?
答案在于复杂度转移和专业化分工。
一个超级Agent处理长链条任务,如同让一位百科全书式的专家从头到尾负责一个大型项目。他可能行,但效率不高,且任何一个环节的深度需求都可能成为瓶颈。而A2A的思路,是组建一个专家团队:架构师、开发、测试、运维各司其职。这时,团队内部的沟通成本就成了主要矛盾。
A2A通信要解决的就是这个“团队沟通”问题,具体拆解为以下几个层面:
1.1 语义理解而不仅是数据传递
两个Agent交换一个JSON字符串很简单。难的是确保接收方能够正确解析发送方的“意图”。例如,一个“代码生成Agent”发给“代码审查Agent”的消息,不仅包含代码片段,还应包含元数据:生成这段代码的原始需求是什么(original_requirement)、使用的框架和语言(context)、期望审查的重点(focus,如安全性、性能、风格)。没有这些上下文,审查Agent可能给出无关紧要的反馈。
1.2 对话状态与任务上下文管理
一次协作往往涉及多轮对话。Agent A问:“用户想要一个登录页面。” Agent B回复:“需要前端还是后端?” Agent A需要记住这是关于“登录页面”任务的延续,并将新的答案(“前端”)补充到任务上下文中,再传递给负责UI的Agent C。A2A框架需要维护这个共享的、不断演进的“任务会话状态”,而不是让每个消息都是孤立的。
1.3 动态路由与能力发现
在一个多Agent系统中,新的Agent可能随时加入,旧的Agent可能离线。当任务到来时,谁最适合处理?A2A框架需要提供一种机制,让Agent能够“广播”自己的能力(如I can review Python code),或者让一个中央协调器(Orchestrator)根据任务类型动态地将消息路由到最合适的Agent。这比在代码里写死调用关系要灵活得多。
1.4 错误处理与协商逻辑
协作不可能一帆风顺。Agent B可能无法理解Agent A的请求,或者执行失败。一个健壮的A2A框架需要定义标准的错误消息格式,并可能支持简单的协商协议。例如,Agent B可以回复:“无法处理此请求,缺少参数X。建议你补充X,或转而求助Agent D,它擅长处理此类模糊请求。” 这要求Agent之间对“协作协议”有共同的理解。
理解了这些核心问题,我们就能明白,一个A2A Demo的价值不在于实现最复杂的路由算法,而在于清晰地展示如何定义消息、建立连接、处理响应和错误,从而为更复杂的协作打下基础。我们的Demo将聚焦于最本质的通信模式。
2. 搭建最小可行Demo:两个Agent如何“对话”
我们设计一个经典场景:一个“任务规划Agent”(Planner)和一个“代码执行Agent”(Executor)的协作。Planner负责解析用户的自然语言需求,并将其分解为具体的、可执行的步骤。Executor则负责执行这些步骤(这里我们简化为执行系统命令或调用代码解释器)。
这个场景虽然简单,但完整包含了A2A的核心要素:请求、响应、结构化数据交换和简单的错误流。
2.1 第一步:定义通信协议(消息格式)
这是A2A的“宪法”。所有Agent都必须遵循同一套消息格式才能互相理解。我们采用一个扩展性较好的JSON结构:
{ "message_id": "unique-uuid-1234", "from_agent": "planner", "to_agent": "executor", "conversation_id": "conv-uuid-5678", "type": "request", // 或 "response", "error" "payload": { "action": "execute_command", "parameters": { "command": "ls -la", "timeout": 10 } }, "context": { "original_task": "列出当前目录文件", "step": 1, "max_steps": 2 }, "timestamp": "2023-10-27T10:00:00Z" }关键字段解析:
message_id&conversation_id: 实现异步通信和会话追踪的基石。每条消息独立,但属于同一个会话。type: 明确消息意图是请求、成功响应还是错误。payload: 核心数据区。action字段定义了接收方应该做什么(如execute_command,analyze_data),parameters是动作所需的参数。这是Agent“技能”的接口定义。context: 承载任务上下文。它让接收方知道自己正在处理一个更大任务的哪一部分,从而做出更合理的决策。这是避免“对话断层”的关键。
2.2 第二步:实现Agent基础类与通信层
我们不依赖复杂的中件间,先用最简单的进程内消息队列(如Python的queue.Queue)模拟通信总线。每个Agent都是一个独立的线程或异步任务,从自己的接收队列读取消息,处理后再放入目标Agent的发送队列。
import json import uuid import threading import queue import subprocess import time from dataclasses import dataclass, asdict from typing import Any, Dict, Optional @dataclass class A2AMessage: message_id: str from_agent: str to_agent: str conversation_id: str type: str # "request", "response", "error" payload: Dict[str, Any] context: Dict[str, Any] timestamp: str def to_dict(self): return asdict(self) @classmethod def from_dict(cls, data: Dict): return cls(**data) class Agent: def __init__(self, name: str, inbox: queue.Queue, outbox: queue.Queue): self.name = name self.inbox = inbox # 接收消息的队列 self.outbox = outbox # 发送消息的队列 self.running = False def send_message(self, to_agent: str, msg_type: str, payload: Dict, context: Dict, conversation_id: str = None): """发送消息的通用方法""" if conversation_id is None: conversation_id = str(uuid.uuid4()) message = A2AMessage( message_id=str(uuid.uuid4()), from_agent=self.name, to_agent=to_agent, conversation_id=conversation_id, type=msg_type, payload=payload, context=context, timestamp=time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()) ) # 在实际A2A中,这里可能是HTTP请求、WebSocket或真正的消息中间件 # 我们简化处理,直接放入“通信总线”(对方的inbox模拟) # 注意:这里需要全局的agent_registry来查找to_agent的inbox,Demo中我们简化使用共享outbox self.outbox.put(message.to_dict()) print(f"[{self.name}] Sent to {to_agent}: {msg_type} - {payload.get('action', 'N/A')}") def process_message(self, message_dict: Dict): """处理接收到的消息。子类必须重写此方法。""" raise NotImplementedError def start(self): """启动Agent,持续监听inbox""" self.running = True def listen(): while self.running: try: # 非阻塞获取,避免线程卡死 msg = self.inbox.get(timeout=0.1) self.process_message(msg) except queue.Empty: continue except Exception as e: print(f"[{self.name}] Error processing message: {e}") thread = threading.Thread(target=listen, daemon=True) thread.start() print(f"[{self.name}] Started.")2.3 第三步:实现具体的Planner和Executor Agent
现在,我们基于基础类实现两个具有特定能力的Agent。
class PlannerAgent(Agent): """任务规划Agent。接收用户请求,分解步骤,并指挥Executor。""" def process_message(self, message_dict: Dict): msg = A2AMessage.from_dict(message_dict) # Planner通常只处理来自“用户”或“协调器”的初始请求 # 本例中,我们假设第一条消息直接发给了Planner if msg.type == "request" and msg.payload.get("action") == "plan_and_execute": user_task = msg.payload["parameters"]["task"] print(f"[{self.name}] Received task: {user_task}") # 简单的规划逻辑:分解任务步骤 steps = self._plan_task(user_task) conversation_id = msg.conversation_id context = msg.context context["original_task"] = user_task context["total_steps"] = len(steps) # 按步骤发送给Executor for i, step in enumerate(steps): step_context = context.copy() step_context["current_step"] = i + 1 self.send_message( to_agent="executor", msg_type="request", payload={ "action": "execute_command", "parameters": {"command": step["command"], "timeout": step.get("timeout", 30)} }, context=step_context, conversation_id=conversation_id ) def _plan_task(self, task: str) -> list: """极简的任务分解逻辑。实际应用中这里会调用LLM。""" # 这是一个硬编码的示例。真实场景中,这里会是一个LLM调用,进行任务分解。 if "list files" in task.lower(): return [{"command": "ls -la", "description": "List all files in current directory"}] elif "current directory" in task.lower(): return [{"command": "pwd", "description": "Print working directory"}] else: # 默认返回一个echo命令 return [{"command": f"echo 'Executing task: {task}'", "description": "Echo the task"}] class ExecutorAgent(Agent): """代码执行Agent。执行系统命令并返回结果。""" def process_message(self, message_dict: Dict): msg = A2AMessage.from_dict(message_dict) if msg.type == "request" and msg.payload.get("action") == "execute_command": command = msg.payload["parameters"]["command"] timeout = msg.payload["parameters"].get("timeout", 30) print(f"[{self.name}] Executing: {command}") try: # 执行系统命令 result = subprocess.run( command, shell=True, capture_output=True, text=True, timeout=timeout ) if result.returncode == 0: response_payload = { "action": "command_result", "result": { "stdout": result.stdout, "stderr": result.stderr, "returncode": result.returncode } } response_type = "response" else: response_payload = { "action": "command_failed", "error": { "stderr": result.stderr, "returncode": result.returncode } } response_type = "error" # 将执行失败定义为一种错误类型 except subprocess.TimeoutExpired: response_payload = { "action": "command_timeout", "error": f"Command timed out after {timeout} seconds." } response_type = "error" except Exception as e: response_payload = { "action": "execution_error", "error": str(e) } response_type = "error" # 将结果返回给发送者(Planner) self.send_message( to_agent=msg.from_agent, msg_type=response_type, payload=response_payload, context=msg.context, # 携带原上下文返回 conversation_id=msg.conversation_id )2.4 第四步:运行Demo并观察通信流
让我们把上述组件组装起来,并模拟一个用户请求。
def main(): # 创建通信队列。在实际分布式系统中,这些队列会是RabbitMQ、Kafka等消息代理。 planner_inbox = queue.Queue() executor_inbox = queue.Queue() # 使用一个共享的“总线”队列来简化消息路由。实际每个Agent应有自己的地址。 message_bus = queue.Queue() # 创建Agent实例。注意,我们将它们的outbox都指向message_bus。 planner = PlannerAgent("planner", planner_inbox, message_bus) executor = ExecutorAgent("executor", executor_inbox, message_bus) # 启动Agent planner.start() executor.start() # 一个简单的路由器线程:从总线读取消息,根据`to_agent`字段投递到对应Agent的inbox def router(): while True: try: msg_dict = message_bus.get(timeout=0.1) msg = A2AMessage.from_dict(msg_dict) if msg.to_agent == "planner": planner_inbox.put(msg_dict) elif msg.to_agent == "executor": executor_inbox.put(msg_dict) else: print(f"[Router] Unknown destination agent: {msg.to_agent}") except queue.Empty: continue except Exception as e: print(f"[Router] Error: {e}") router_thread = threading.Thread(target=router, daemon=True) router_thread.start() # 模拟用户发起一个任务 print("\n=== 模拟用户请求:'请列出当前目录的文件' ===") user_message = A2AMessage( message_id=str(uuid.uuid4()), from_agent="user", to_agent="planner", conversation_id=str(uuid.uuid4()), type="request", payload={"action": "plan_and_execute", "parameters": {"task": "请列出当前目录的文件"}}, context={"user_id": "demo_user"}, timestamp=time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()) ) # 将用户请求放入Planner的收件箱 planner_inbox.put(user_message.to_dict()) # 等待一段时间,让Agent完成处理 time.sleep(3) print("\n=== Demo 结束 ===") if __name__ == "__main__": main()运行这段代码,你将在控制台看到类似以下的输出:
[planner] Started. [executor] Started. === 模拟用户请求:'请列出当前目录的文件' === [planner] Received task: 请列出当前目录的文件 [planner] Sent to executor: request - execute_command [executor] Executing: ls -la [executor] Sent to planner: response - command_result这个简单的流程清晰地展示了A2A通信的骨架:
- 用户/系统向 Planner 发送一个结构化请求。
- Planner解析请求,进行规划(分解任务),生成一个给 Executor 的标准化请求消息。
- 消息通过“总线”(router)被路由到Executor的收件箱。
- Executor执行命令,并将结果封装成标准响应(或错误)消息,发回给 Planner。
- 消息再次通过总线路由回Planner。
至此,一个最小可运行的A2A通信Demo就完成了。它虽然简陋,但已经包含了消息定义、Agent角色、请求-响应模式、错误反馈和上下文传递这些核心要素。
3. 从Demo到生产:A2A工程化必须考虑的四个维度
Demo跑通了,但如果你认为这就是A2A的全部,那就把问题想简单了。单次成功通信只是起点,要让Agent团队真正可靠地工作,我们必须面对工程化的挑战。以下四个维度,是评估一个A2A框架是否成熟的关键。
3.1 通信模式:不止于请求-响应
我们的Demo使用了最简单的同步请求-响应模式。但在实际场景中,Agent协作可能需要更灵活的模式:
| 模式 | 描述 | 适用场景 |
|---|---|---|
| 请求-响应 | 一对一,发送方等待回复。 | 明确的指令执行、查询。 |
| 发布-订阅 | 一个Agent广播消息,多个感兴趣的Agent接收并处理。 | 事件通知(如“任务完成”、“系统异常”)。 |
| 工作流/管道 | 消息按预定顺序流经多个Agent,每个处理完传递给下一个。 | 有严格顺序的数据处理流水线。 |
| 广播 | 向所有Agent发送消息。 | 系统配置更新、全局状态同步。 |
例如,一个“日志监控Agent”可能以发布-订阅模式广播错误警报,而“告警聚合Agent”和“自动修复Agent”同时订阅并采取不同行动。选择哪种模式,取决于Agent间的耦合度和任务性质。
3.2 状态、上下文与记忆管理
这是A2A中最容易出问题的地方。我们的Demo在消息中携带了context字段,这是一个好的开始,但远远不够。
- 会话状态 vs Agent内部状态:
conversation_id关联的是“任务会话”状态。而Agent自身也可能有需要维护的内部状态(如已使用的API额度、缓存的历史结果)。这两者需要区分管理。 - 上下文窗口与摘要:在多轮复杂协作中,完整的上下文可能非常大。需要设计摘要机制,将冗长的历史对话提炼成关键信息,再传递给下一个Agent,以避免超出LLM的上下文限制。
- 共享记忆体:对于需要多个Agent频繁访问的公共信息(如项目规范、API密钥配置),可以设计一个“共享记忆Agent”或使用外部数据库(如矢量数据库),其他Agent通过查询来获取,而不是在消息中反复传递。
一个进阶的设计是引入**“协调器Agent”**。它不直接处理具体任务,而是专职维护整个工作流的状态机,记录哪个步骤已完成、哪个正在执行、哪个失败了,并负责将适当的上下文传递给下一个执行的Agent。这大大减轻了业务Agent的负担。
3.3 错误处理、重试与降级策略
Demo中,Executor只是将错误封装成消息返回。在生产环境中,这不够。
- 错误分类与处理策略:
- 瞬时错误(如网络超时):应自动重试,并有指数退避策略。
- 逻辑错误(如参数无效):应通知上游Agent,并可能携带修正建议。
- 致命错误(如依赖服务不可用):应触发工作流暂停,并通知人工或更高层级的协调器。
- 重试机制:重试不应无限进行。需要在消息或协调器中定义最大重试次数。重试时,可以考虑微调参数(如增加超时时间)后再次尝试。
- 降级与备选路径:如果某个Agent持续失败,系统是否有一条备选路径?例如,当“图像生成Agent”超时时,协调器是否可以转而请求“文本描述Agent”生成一段详细描述作为替代输出?这需要预先定义好工作流的备选分支。
3.4 安全、权限与监控
当Agent能够自主通信时,安全就成为重中之重。
- 身份认证与授权:每个Agent都应有身份标识。消息传递需要验证发送者是否有权向接收者发送此类消息,以及接收者是否有权执行请求的操作。这通常通过令牌(Token)或双向TLS实现。
- 输入验证与净化:Executor Agent直接执行系统命令是极其危险的。生产环境必须对
command参数进行严格的白名单过滤,或仅允许调用安全的内部API。 - 通信加密:所有跨进程或跨网络的Agent通信必须加密(如使用HTTPS、WSS)。
- 可观测性:必须记录所有A2A消息的流向、耗时和结果。这需要集中的日志、指标(Metrics)和分布式追踪(Tracing)系统。当协作出错时,你可以通过
conversation_id完整回溯整个工作流的执行轨迹,快速定位问题节点。
4. 主流框架的实践与我们的选择
了解了原理和挑战,我们来看看业界是如何实践的。目前,实现A2A通信主要有两种路径:
4.1 基于现有Agent框架的“编排”方案
像LangChain、LlamaIndex、AutoGen、CrewAI这类高阶框架,它们在内核已经抽象了Agent间的协作模式。
- LangChain通过
AgentExecutor和Tool机制,让一个主Agent根据LLM的思考过程,决定调用哪个工具(可视为一个简化Agent)。其多Agent协作更多通过SequentialChain或RouterChain来实现工作流。 - AutoGen则直接以“多Agent对话”为核心范式。你定义多个
AssistantAgent和UserProxyAgent,它们在一个群聊中通过发送消息自动协作。框架底层处理了消息路由和会话管理。 - CrewAI明确引入了
Agent、Task和Crew的概念。Crew(团队)负责协调Agent按顺序或并行执行Task,并管理它们之间的上下文传递。
选择这类框架的好处是“开箱即用”。你无需从零设计消息协议和路由器,可以快速搭建复杂的多Agent工作流。但代价是被框架的设计哲学和复杂度所绑定,定制深度通信逻辑或集成非标准组件可能会比较困难。
4.2 自建轻量级通信总线
这正是我们Demo所演示的路径。你可以基于RabbitMQ、Apache Kafka、Redis Pub/Sub甚至HTTP Webhook来构建自己的消息总线。每个Agent作为一个独立服务,订阅特定的主题或队列。
这种方案的优点是极致灵活和可控。你可以完全自定义消息格式、路由逻辑、持久化策略和监控指标。它适合对性能、可靠性和架构有极高要求的场景,或者当你需要将AI Agent与已有的、非AI的微服务进行深度集成时。
但它的缺点也很明显:复杂度高。你需要自己实现之前讨论的所有工程化特性:服务发现、负载均衡、重试、死信队列、分布式追踪等。这本质上是在构建一个分布式的消息驱动系统。
4.3 如何选择?一个简单的决策框架
面对具体项目时,你可以问自己以下几个问题来做决定:
- 协作复杂度:Agent之间是简单的线性管道,还是复杂的网状对话?线性管道用工作流引擎或简单编排即可;网状对话可能需要更通用的消息总线。
- 集成需求:是否需要与大量现有系统(数据库、API服务、监控告警)通信?是的话,基于标准消息中间件(如Kafka)的自建方案更合适。
- 团队技能:团队是否熟悉分布式系统开发和运维?如果不是,使用成熟的Agent框架(如AutoGen)能大幅降低入门门槛。
- 控制与定制:是否需要绝对控制通信的每个细节(如加密算法、压缩格式、自定义的共识机制)?自建是唯一选择。
- 开发速度 vs 长期维护:原型验证阶段,框架能帮你快速看到效果。但如果预计系统会长期演进、规模扩大,早期在通信层投入设计往往是值得的。
对于大多数从0到1的AI应用项目,我的建议是:先从高阶框架(如AutoGen或CrewAI)开始,快速验证多Agent协作的业务价值。当协作模式稳定,且遇到框架无法满足的特定性能或集成需求时,再考虑将核心的通信层抽离出来,用更底层的工具进行定制化实现。我们的Demo价值就在于,它揭开了这层抽象,让你理解了框架底层可能发生的故事。
回过头看,A2A通信的本质,是为AI能力模块化之后产生的“集成问题”提供标准化的解决方案。它让每个Agent可以专注于自己的核心技能(如编码、分析、执行),而将复杂的协作逻辑交给通信框架来管理。理解了这个本质,无论是选用现成框架还是自建轮子,你都能做出更明智的设计决策。
最终,一个健壮的A2A系统,看起来不像是一群AI在对话,而更像是一个高度自动化、职责清晰、能够自我协调的数字团队在默默工作。而构建这个团队的起点,就是从理解两个Agent之间如何说好第一句话开始。