20分钟实战A2A协议:构建可协作AI Agent系统的核心通信框架

📅 2026/8/7 12:02:52
20分钟实战A2A协议:构建可协作AI Agent系统的核心通信框架
最近在尝试将多个AI Agent串联起来完成复杂任务时你是否也遇到了这样的困扰Agent之间如何高效、可靠地通信消息格式五花八门状态难以同步错误处理更是让人头疼。这正是A2AAgent-to-Agent协议要解决的核心问题。本文将带你从零开始通过一个完整的实战项目在20分钟内快速掌握A2A协议的核心思想与实现方法。无论你是想了解多智能体协作原理还是需要在项目中集成Agent通信能力这篇教程都能提供一套可直接复用的代码方案和清晰的工程化思路。1. AI Agent 与 A2A 协议为什么需要它在深入代码之前我们首先要厘清几个核心概念。AI Agent智能体并非一个遥不可及的概念你可以将它理解为一个具备感知、决策和执行能力的软件实体。它接收来自环境可能是用户、传感器或其他系统的输入通过内部逻辑如今常由大语言模型驱动进行推理最终输出行动来影响环境。单个Agent的能力是有限的。当任务变得复杂时比如需要同时处理数据分析、报告生成和结果通知让一个“全能”Agent来完成所有工作不仅效率低下而且难以维护和扩展。这时多智能体系统Multi-Agent System, MAS的优势就显现出来了。我们可以设计多个各司其职的Agent如“数据分析师”、“报告撰写员”、“通知专员”让它们协同工作。然而协同工作的前提是通信。A2A协议就是为了规范Agent之间的通信而诞生的一套“约定”或“标准”。它定义了消息格式Agent之间说什么内容如何组织如JSON结构通信模式怎么说是请求-响应还是发布-订阅会话管理如何关联一次完整的对话或任务流程错误与重试如果通信失败怎么办没有A2A协议Agent之间的协作就像一群人用各自方言开会混乱且低效。有了A2A协议它们就像使用同一种标准语言和流程进行协作可靠性和效率大大提升。2. 环境准备与项目初始化我们将使用Python作为开发语言因为它拥有丰富的AI开发生态。本项目不依赖特定的大模型API重点在于通信逻辑的实现因此你可以轻松替换底层的LLM调用部分。环境要求操作系统Windows 10/11, macOS, 或 Linux (Ubuntu 20.04)Python版本3.8 或更高版本包管理工具pip创建项目并安装依赖首先创建一个新的项目目录并初始化虚拟环境这是保证依赖隔离的最佳实践。# 创建项目目录 mkdir a2a-agent-demo cd a2a-agent-demo # 创建虚拟环境 (Windows用户使用 python -m venv venv) python3 -m venv venv # 激活虚拟环境 # Windows: venv\Scripts\activate # macOS/Linux: source venv/bin/activate # 安装核心依赖 pip install pydantic # 用于数据验证和设置管理 pip install requests # 用于HTTP通信模拟Agent间调用 pip install loguru # 用于更美观、结构化的日志输出项目结构规划一个清晰的项目结构有助于代码管理和扩展。我们采用以下设计a2a-agent-demo/ ├── agents/ # 存放各个Agent的实现 │ ├── __init__.py │ ├── base_agent.py # Agent基类定义通用接口 │ ├── analyst_agent.py # 数据分析Agent │ └── reporter_agent.py # 报告生成Agent ├── protocols/ # 存放通信协议相关定义 │ ├── __init__.py │ └── a2a.py # A2A协议的核心消息定义 ├── utils/ # 工具函数 │ ├── __init__.py │ └── logger.py # 日志配置 ├── main.py # 主程序入口编排Agent工作流 └── requirements.txt # 项目依赖列表你可以使用以下命令快速创建这个结构mkdir -p agents protocols utils touch agents/__init__.py agents/base_agent.py agents/analyst_agent.py agents/reporter_agent.py touch protocols/__init__.py protocols/a2a.py touch utils/__init__.py utils/logger.py touch main.py requirements.txt将已安装的依赖写入requirements.txtpydantic2.0.0 requests2.28.0 loguru0.7.03. 定义A2A协议消息是协作的基石协议层是A2A系统的核心。我们使用pydantic来定义强类型的消息模型这能在开发阶段就捕获许多数据格式错误。在protocols/a2a.py中我们定义通信所需的基本消息结构# protocols/a2a.py from enum import Enum from typing import Any, Dict, Optional from pydantic import BaseModel, Field from datetime import datetime class MessageType(str, Enum): 定义消息类型枚举明确通信意图 TASK_REQUEST task_request # 任务请求 TASK_RESULT task_result # 任务结果 ERROR error # 错误信息 HEARTBEAT heartbeat # 心跳检测用于健康检查 class AgentMessage(BaseModel): A2A协议核心消息体。 所有Agent间通信都必须封装在此消息结构内。 # 消息元数据 msg_id: str Field(..., description全局唯一消息ID用于追踪) msg_type: MessageType Field(..., description消息类型) timestamp: datetime Field(default_factorydatetime.now, description消息创建时间戳) # 会话与路由信息 session_id: str Field(..., description会话ID关联一次完整的任务流程) sender_id: str Field(..., description发送方Agent ID) receiver_id: str Field(..., description接收方Agent ID) # 消息内容 payload: Dict[str, Any] Field(default_factorydict, description消息实际负载任务参数或结果) # 上下文与溯源 parent_msg_id: Optional[str] Field(None, description父消息ID用于消息链溯源) metadata: Dict[str, Any] Field(default_factorydict, description扩展元数据如优先级、超时时间等) class Config: # 允许使用枚举的字符串值进行序列化/反序列化 use_enum_values True关键字段解析msg_id和session_id这是实现可靠通信的关键。msg_id标识单条消息用于去重和确认session_id标识一个完整的业务会话方便日志聚合和问题排查。msg_type使用枚举强制约束消息类型避免拼写错误并使接收方的消息路由逻辑更清晰。payload这是一个灵活的字典用于承载任何业务数据。在实际项目中你可以根据不同的msg_type进一步定义Payload的子类来约束其结构。parent_msg_id实现了消息的“链式”追踪。当Reporter Agent回复Analyst Agent时其消息的parent_msg_id就是请求消息的msg_id。这个简单的消息模型已经具备了生产级通信协议的雏形涵盖了身份、时序、内容和关联关系。4. 实现基础Agent类在agents/base_agent.py中我们创建一个所有具体Agent都将继承的基类。它封装了Agent的通用属性、消息发送/接收的模板方法以及生命周期管理。# agents/base_agent.py import uuid from abc import ABC, abstractmethod from typing import Any, Dict from loguru import logger from protocols.a2a import AgentMessage, MessageType class BaseAgent(ABC): Agent基类。 定义了Agent的通用接口和基础行为如ID管理、消息发送/接收模板。 def __init__(self, agent_id: str, agent_name: str): 初始化Agent。 :param agent_id: Agent的唯一标识符 :param agent_name: Agent的可读名称 self.agent_id agent_id self.agent_name agent_name self._session_cache {} # 简单的会话缓存用于存储上下文 logger.info(fAgent 初始化: {self.agent_name} (ID: {self.agent_id})) def receive_message(self, message: AgentMessage) - AgentMessage: 接收并处理消息的公共入口。这是一个模板方法。 1. 验证消息接收者是否正确。 2. 根据消息类型路由到具体的处理函数。 3. 封装响应消息。 :param message: 接收到的AgentMessage :return: 处理后的响应AgentMessage # 1. 基础验证这条消息是发给我的吗 if message.receiver_id ! self.agent_id: error_msg f消息接收者ID不匹配。预期: {self.agent_id}, 实际: {message.receiver_id} logger.warning(error_msg) return self._create_error_message( original_messagemessage, error_infoerror_msg ) logger.info(f{self.agent_name} 收到消息: {message.msg_type} (会话: {message.session_id})) # 2. 根据消息类型路由处理逻辑 try: if message.msg_type MessageType.TASK_REQUEST: result_payload self.process_task(message.payload, message.session_id) response_type MessageType.TASK_RESULT elif message.msg_type MessageType.HEARTBEAT: result_payload {status: alive, agent_id: self.agent_id} response_type MessageType.HEARTBEAT else: # 暂时不支持其他类型的消息处理 raise ValueError(f不支持的消息类型: {message.msg_type}) # 3. 创建并返回响应消息 response_msg AgentMessage( msg_idstr(uuid.uuid4()), msg_typeresponse_type, session_idmessage.session_id, sender_idself.agent_id, receiver_idmessage.sender_id, payloadresult_payload, parent_msg_idmessage.msg_id # 关键关联请求与响应 ) return response_msg except Exception as e: logger.error(f{self.agent_name} 处理消息时发生异常: {e}) # 返回错误消息 return self._create_error_message( original_messagemessage, error_infostr(e) ) abstractmethod def process_task(self, task_payload: Dict[str, Any], session_id: str) - Dict[str, Any]: 处理任务请求的抽象方法。由子类实现具体业务逻辑。 :param task_payload: 任务负载数据 :param session_id: 当前会话ID :return: 处理结果将被放入响应消息的payload中 pass def _create_error_message(self, original_message: AgentMessage, error_info: str) - AgentMessage: 创建一个标准的错误响应消息。 return AgentMessage( msg_idstr(uuid.uuid4()), msg_typeMessageType.ERROR, session_idoriginal_message.session_id, sender_idself.agent_id, receiver_idoriginal_message.sender_id, payload{error: error_info, original_msg_id: original_message.msg_id}, parent_msg_idoriginal_message.msg_id ) def send_message(self, receiver: BaseAgent, message: AgentMessage) - AgentMessage: 向另一个Agent发送消息并获取响应的简化方法。 在实际分布式系统中这里会是网络调用HTTP/RPC等。 本例中我们直接调用对方Agent的receive_message方法。 logger.debug(f{self.agent_name} 向 {receiver.agent_name} 发送消息: {message.msg_type}) # 模拟网络传输直接调用接收者的处理方法 response receiver.receive_message(message) logger.debug(f{self.agent_name} 收到来自 {receiver.agent_name} 的响应: {response.msg_type}) return response设计要点模板方法模式receive_message方法定义了消息处理的固定流程验证-路由-响应子类只需实现process_task这个可变部分。这保证了所有Agent行为的一致性。错误处理标准化任何在处理过程中抛出的异常都会被捕获并封装成标准格式的ERROR类型消息返回给发送方实现了错误的跨Agent传递。松耦合通信send_message方法目前是本地调用但将其抽象出来意味着未来可以轻松替换为HTTP、gRPC或消息队列等远程通信方式而不需要修改每个Agent的业务逻辑。5. 构建业务Agent数据分析师与报告员现在我们基于上述框架实现两个具有简单业务逻辑的Agent。第一个Agent数据分析师 (AnalystAgent)它的职责是接收原始数据进行“分析”这里简化为计算平均值并返回结果。# agents/analyst_agent.py from typing import Any, Dict from loguru import logger from .base_agent import BaseAgent class AnalystAgent(BaseAgent): 数据分析师Agent负责处理数值数据计算统计指标。 def process_task(self, task_payload: Dict[str, Any], session_id: str) - Dict[str, Any]: 处理数据分析任务。 期望payload格式: {data: [list_of_numbers], operation: mean/max/min} logger.info(fAnalystAgent 开始处理任务会话: {session_id}) # 1. 提取并验证输入 data task_payload.get(data, []) operation task_payload.get(operation, mean) if not isinstance(data, list) or not all(isinstance(x, (int, float)) for x in data): raise ValueError(Payload中data字段必须是一个数字列表) if len(data) 0: raise ValueError(数据列表不能为空) # 2. 执行“分析”逻辑 result None if operation mean: result sum(data) / len(data) analysis_desc f计算了 {len(data)} 个数据的平均值 elif operation max: result max(data) analysis_desc f找出了数据列表中的最大值 elif operation min: result min(data) analysis_desc f找出了数据列表中的最小值 else: raise ValueError(f不支持的操作类型: {operation}) # 3. 组织并返回结果 logger.success(fAnalystAgent 分析完成。操作[{operation}]结果: {result}) return { analysis_result: result, operation_performed: operation, description: analysis_desc, original_data_size: len(data) }第二个Agent报告生成员 (ReporterAgent)它的职责是将分析结果格式化为一份可读的报告。# agents/reporter_agent.py from typing import Any, Dict from datetime import datetime from loguru import logger from .base_agent import BaseAgent class ReporterAgent(BaseAgent): 报告生成员Agent负责将结构化的分析结果转化为文本报告。 def process_task(self, task_payload: Dict[str, Any], session_id: str) - Dict[str, Any]: 处理报告生成任务。 期望payload格式: 包含analysis_result等字段的字典即AnalystAgent的输出。 logger.info(fReporterAgent 开始生成报告会话: {session_id}) # 1. 提取分析结果 analysis_result task_payload.get(analysis_result) operation task_payload.get(operation_performed, 未知操作) desc task_payload.get(description, ) if analysis_result is None: raise ValueError(无法生成报告缺少‘analysis_result’字段) # 2. 生成报告文本这里可以集成LLM调用例如使用OpenAI API或本地模型 # 本例中我们进行简单的字符串格式化 report_text f **数据分析报告** ---------------------------- 生成时间: {datetime.now().strftime(%Y-%m-%d %H:%M:%S)} 会话ID: {session_id} ---------------------------- 执行操作: {operation} 操作描述: {desc} 分析结果: {analysis_result} ---------------------------- 结论: 根据分析数据集的{operation}值为 {analysis_result}。 # 3. 返回报告 logger.success(fReporterAgent 报告生成完成长度: {len(report_text)} 字符) return { report_content: report_text.strip(), format: markdown, generated_at: datetime.now().isoformat() }至此我们拥有了两个具备明确职责、通过标准化A2A消息进行通信的智能体。它们的业务逻辑简单但架构是完整且可扩展的。6. 实战演练编排一个完整的工作流现在让我们在main.py中将这些组件串联起来模拟一个完整的“数据分析并生成报告”的业务流程。# main.py import uuid from loguru import logger from protocols.a2a import AgentMessage, MessageType from agents.analyst_agent import AnalystAgent from agents.reporter_agent import ReporterAgent def main(): 主函数演示A2A协议下的多Agent协作工作流。 logger.info(开始A2A多Agent协作演示...) # 1. 初始化Agent analyst AnalystAgent(agent_idagent_analyst_001, agent_name高级数据分析师) reporter ReporterAgent(agent_idagent_reporter_001, agent_name报告生成专家) logger.info(fAgent初始化完成: {analyst.agent_name}, {reporter.agent_name}) # 2. 创建本次任务的唯一会话ID session_id fsession_{uuid.uuid4().hex[:8]} logger.info(f创建新会话: {session_id}) # 3. 第一步用户或调度器向AnalystAgent发送数据分析任务 # 创建任务请求消息 task_to_analyst AgentMessage( msg_idstr(uuid.uuid4()), msg_typeMessageType.TASK_REQUEST, session_idsession_id, sender_iduser_orchestrator, # 发送者可以是用户或一个编排器Agent receiver_idanalyst.agent_id, payload{ data: [23.5, 45.2, 67.8, 12.1, 89.0, 34.4], operation: mean } ) logger.info(f【步骤1】发送数据分析任务给 {analyst.agent_name}...) # 这里我们模拟“用户”直接调用analyst实际可能通过消息总线 analysis_result_msg analyst.receive_message(task_to_analyst) # 检查第一步结果 if analysis_result_msg.msg_type MessageType.ERROR: logger.error(f数据分析阶段失败: {analysis_result_msg.payload}) return logger.success(f数据分析完成。结果: {analysis_result_msg.payload}) # 4. 第二步AnalystAgent将分析结果发送给ReporterAgent请求生成报告 # 注意这里sender_id变成了analyst体现了Agent间的主动协作 task_to_reporter AgentMessage( msg_idstr(uuid.uuid4()), msg_typeMessageType.TASK_REQUEST, session_idsession_id, # 保持同一会话ID便于追踪 sender_idanalyst.agent_id, receiver_idreporter.agent_id, payloadanalysis_result_msg.payload # 将上一步的结果作为payload传递 ) logger.info(f【步骤2】{analyst.agent_name} 将结果发送给 {reporter.agent_name} 以生成报告...) # 使用Agent基类提供的send_message方法进行通信 final_report_msg analyst.send_message(reporter, task_to_reporter) # 检查最终结果 if final_report_msg.msg_type MessageType.ERROR: logger.error(f报告生成阶段失败: {final_report_msg.payload}) return # 5. 输出最终报告 final_report final_report_msg.payload logger.success(*50) logger.success(✅ 任务执行成功最终报告如下) logger.success(*50) print(final_report.get(report_content, 无报告内容)) # 打印报告内容 logger.success(*50) # 6. 演示消息链追溯通过parent_msg_id可以追溯整个任务流 logger.info(【消息链追溯演示】) logger.info(f最终报告消息ID: {final_report_msg.msg_id}) logger.info(f其父消息ID即给Reporter的请求: {final_report_msg.parent_msg_id}) logger.info(f给Reporter的请求消息ID: {task_to_reporter.msg_id}) logger.info(f其父消息ID即Analyst的结果: {task_to_reporter.parent_msg_id}) logger.info(fAnalyst的结果消息ID: {analysis_result_msg.msg_id}) logger.info(f其父消息ID即最初的任务请求: {analysis_result_msg.parent_msg_id}) logger.info(f最初的任务请求消息ID: {task_to_analyst.msg_id}) logger.info(通过parent_msg_id可以清晰重构出完整的任务执行链路。) if __name__ __main__: # 配置日志使其更美观可选 logger.add(a2a_demo.log, rotation1 MB, levelINFO) main()运行与验证在项目根目录下执行python main.py你应该能看到类似以下的控制台输出清晰地展示了消息在Agent间的流动、处理以及最终的报告结果2024-XX-XX XX:XX:XX.XXX | INFO | __main__:main:16 - 开始A2A多Agent协作演示... 2024-XX-XX XX:XX:XX.XXX | INFO | agents.base_agent:__init__:24 - Agent 初始化: 高级数据分析师 (ID: agent_analyst_001) ... 2024-XX-XX XX:XX:XX.XXX | INFO | __main__:main:30 - 【步骤1】发送数据分析任务给 高级数据分析师... 2024-XX-XX XX:XX:XX.XXX | INFO | agents.base_agent:receive_message:38 - 高级数据分析师 收到消息: task_request (会话: session_xxxxxx) 2024-XX-XX XX:XX:XX.XXX | SUCCESS | agents.analyst_agent:process_task:41 - AnalystAgent 分析完成。操作[mean]结果: 45.333333333333336 2024-XX-XX XX:XX:XX.XXX | SUCCESS | __main__:main:43 - 数据分析完成。结果: {analysis_result: 45.3333..., operation_performed: mean, ...} ... ✅ 任务执行成功最终报告如下 **数据分析报告** ------------------------------------ 生成时间: 2024-XX-XX XX:XX:XX 会话ID: session_xxxxxx ------------------------------------ 执行操作: mean 操作描述: 计算了 6 个数据的平均值 分析结果: 45.333333333333336 ------------------------------------ 结论: 根据分析数据集的mean值为 45.333333333333336。 这个简单的流程演示了任务触发 - Agent A 处理 - A2A 协议通信 - Agent B 处理 - 结果返回的完整闭环。parent_msg_id形成的链条让你可以轻松追踪整个会话中每一跳的消息。7. 常见问题与排查思路FAQ在实际开发和集成A2A系统时你可能会遇到以下典型问题问题现象可能原因排查思路与解决方案Agent 收不到消息1.receiver_id拼写错误或与目标Agent ID不匹配。2. 消息路由机制如消息队列主题配置错误。3. 网络问题或Agent进程未启动。1.日志优先检查发送和接收方的日志确认消息ID、发送/接收者ID。2.验证协议确保消息体符合AgentMessage的Pydantic模型可通过message.model_dump_json()序列化后检查。3.简化测试先使用本文的本地直接调用模式 (send_message) 测试业务逻辑再排查分布式通信问题。消息处理超时或无响应1. 接收方Agent的process_task方法存在死循环或耗时过长。2. 未设置合理的消息超时机制。3. 出现了未捕获的异常导致receive_message方法未能返回。1.超时设置在send_message或网络客户端中增加超时参数。2.异步处理对于长任务应考虑异步模式。接收方立即返回ACCEPTED类型消息再通过回调或另一个消息通道返回结果。3.完善异常处理确保process_task内部有细粒度的try-catch并将异常信息通过_create_error_message返回。消息顺序错乱或丢失1. 在异步或分布式环境下消息可能不按发送顺序到达。2. 消息队列的持久化或确认机制未开启。1.使用session_id和msg_id在应用层通过session_id关联业务通过msg_id去重。对于强顺序需求可在payload中增加序列号。2.选择可靠中间件使用如 RabbitMQ、Kafka 等提供消息持久化、顺序性和确认机制的消息队列。协议扩展性差新增消息类型麻烦1.MessageType枚举和receive_message中的if-elif逻辑硬编码每加一个类型都要改多处。1.策略模式将不同msg_type的处理逻辑抽象为独立的Handler类并在Agent初始化时注册。receive_message只需根据msg_type查找对应的Handler并执行。这符合开闭原则。如何集成真实的LLM本文示例使用了模拟逻辑实际需要调用大模型API。1.抽象LLM客户端在agents/目录下创建llm_client.py封装对 OpenAI、通义千问等API的调用。2.在process_task中调用将任务描述和上下文构造成 Prompt调用LLM客户端获取结果。务必注意处理LLM的异步响应、速率限制和token长度限制。8. 最佳实践与进阶架构建议掌握了基础实现后要将A2A协议用于生产环境还需要考虑以下工程化实践1. 通信层抽象与实现本文的send_message是本地调用。在生产中你需要一个真正的通信层Transport Layer。建议定义一个MessageTransport抽象类然后为其提供不同实现# protocols/transport.py from abc import ABC, abstractmethod from .a2a import AgentMessage class MessageTransport(ABC): abstractmethod def send(self, message: AgentMessage, target_agent_id: str) - AgentMessage: 发送消息到指定Agent并等待响应。 pass abstractmethod def register_agent(self, agent_id: str, callback_function): 注册Agent及其消息回调函数。 pass # 实现类示例HTTPTransport, RabbitMQTransport, RedisPubSubTransport这样Agent基类只需持有MessageTransport的实例无需关心底层是HTTP、gRPC还是消息队列。2. 引入Harness层基础设施层正如网络热词中提到的Harness是包裹在Agent核心逻辑之外的基础设施层。它不替代Agent而是提供通用能力。你可以构建一个AgentHarness类为Agent提供可观测性自动记录所有入站/出站消息的指标和日志。弹性能力自动重试、熔断、降级。安全检查对输入/输出payload进行验证或过滤。上下文管理自动维护和注入会话上下文。 Agent只需关注process_task中的纯业务逻辑其他交叉关切点由Harness统一处理。3. 设计清晰的Skill与Tool调用规范当Agent需要调用外部能力如搜索、数据库查询、工具函数时应定义统一的Skill/Tool接口。这可以与A2A协议结合对内Agent间使用A2A消息通信。对外Agent与工具定义一套Tool接口Agent通过调用tool.execute(params)来使用能力。这有助于能力复用和测试。4. 会话Session与状态管理对于多轮交互的复杂任务简单的session_id可能不够。需要设计一个Session对象存储会话元数据创建时间、状态、所属用户。消息历史列表。共享的上下文数据键值对。 每个Agent在处理消息时可以从Harness或通信层获取当前Session对象实现跨Agent的上下文传递。5. 测试策略单元测试针对每个Agent的process_task方法模拟输入payload验证输出。集成测试启动多个Agent实例模拟完整的A2A消息流验证端到端功能。契约测试利用Pydantic模型确保消息格式在Agent版本迭代中保持兼容。通过以上步骤你不仅实现了一个可运行的A2A协议Demo更掌握了一套构建可维护、可扩展、高可靠的多智能体系统的设计思路。从定义协议、实现基类、开发业务Agent到编排工作流和规划进阶架构每一步都着眼于解决实际协作中的痛点。你可以在此基础上引入网络通信、集成真实LLM、添加更复杂的业务逻辑逐步构建起属于你自己的AI Agent应用。