最近在尝试将多个AI Agent串联起来完成复杂任务时,你是否也遇到了这样的困扰:Agent之间如何高效、可靠地通信?消息格式五花八门,状态难以同步,错误处理更是让人头疼。这正是A2A(Agent-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.txt

pydantic>=2.0.0
requests>=2.28.0
loguru>=0.7.0

3. 定义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_factory=datetime.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_factory=dict, description="消息实际负载(任务参数或结果)")
    
    # 上下文与溯源
    parent_msg_id: Optional[str] = Field(None, description="父消息ID,用于消息链溯源")
    metadata: Dict[str, Any] = Field(default_factory=dict, 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(f"Agent 初始化: {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_message=message,
                error_info=error_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_id=str(uuid.uuid4()),
                msg_type=response_type,
                session_id=message.session_id,
                sender_id=self.agent_id,
                receiver_id=message.sender_id,
                payload=result_payload,
                parent_msg_id=message.msg_id  # 关键:关联请求与响应
            )
            return response_msg
            
        except Exception as e:
            logger.error(f"{self.agent_name} 处理消息时发生异常: {e}")
            # 返回错误消息
            return self._create_error_message(
                original_message=message,
                error_info=str(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_id=str(uuid.uuid4()),
            msg_type=MessageType.ERROR,
            session_id=original_message.session_id,
            sender_id=self.agent_id,
            receiver_id=original_message.sender_id,
            payload={"error": error_info, "original_msg_id": original_message.msg_id},
            parent_msg_id=original_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

设计要点:

  1. 模板方法模式 receive_message 方法定义了消息处理的固定流程(验证->路由->响应),子类只需实现 process_task 这个可变部分。这保证了所有Agent行为的一致性。
  2. 错误处理标准化 :任何在处理过程中抛出的异常都会被捕获,并封装成标准格式的 ERROR 类型消息返回给发送方,实现了错误的跨Agent传递。
  3. 松耦合通信 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(f"AnalystAgent 开始处理任务,会话: {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(f"AnalystAgent 分析完成。操作[{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(f"ReporterAgent 开始生成报告,会话: {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(f"ReporterAgent 报告生成完成,长度: {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_id="agent_analyst_001", agent_name="高级数据分析师")
    reporter = ReporterAgent(agent_id="agent_reporter_001", agent_name="报告生成专家")
    
    logger.info(f"Agent初始化完成: {analyst.agent_name}, {reporter.agent_name}")
    
    # 2. 创建本次任务的唯一会话ID
    session_id = f"session_{uuid.uuid4().hex[:8]}"
    logger.info(f"创建新会话: {session_id}")
    
    # 3. 第一步:用户(或调度器)向AnalystAgent发送数据分析任务
    # 创建任务请求消息
    task_to_analyst = AgentMessage(
        msg_id=str(uuid.uuid4()),
        msg_type=MessageType.TASK_REQUEST,
        session_id=session_id,
        sender_id="user_orchestrator",  # 发送者可以是用户或一个编排器Agent
        receiver_id=analyst.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_id=str(uuid.uuid4()),
        msg_type=MessageType.TASK_REQUEST,
        session_id=session_id,  # 保持同一会话ID,便于追踪
        sender_id=analyst.agent_id,
        receiver_id=reporter.agent_id,
        payload=analysis_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(f"Analyst的结果消息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", rotation="1 MB", level="INFO")
    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应用。

更多推荐