Agent-to-Agent (A2A) 协议深度实战:构建可互操作的 AI Agent 生态系统

2025 年 4 月,Google 发布了 Agent-to-Agent Protocol(A2A),为 AI Agent 之间的通信与协作定义了一套开放标准。本文将从工程实现角度深度拆解 A2A 协议的核心机制,并通过完整的代码示例展示如何构建生产级的 A2A Agent 系统。

一、为什么需要 A2A:从 Agent 孤岛到协作网络

当前的 AI Agent 生态面临严重的互操作性问题:不同框架(LangChain、AutoGPT、CrewAI、OpenAI Agents SDK)构建的 Agent 彼此无法通信;企业内部部署的 Agent 无法与外部服务商的 Agent 安全交互;每个 Agent 都需要独立实现工具发现、状态同步和协作逻辑。

A2A 协议试图解决的就是这个问题。它的核心设计理念包括:

  • 基于 HTTP/JSON-RPC 的轻量级通信,无需专用中间件
  • Agent Card(服务发现),类似 API 的 OpenAPI 描述
  • 任务生命周期管理,支持长时间运行的异步任务
  • 流式与推送通知,支持 SSE 实时流和 Webhook 回调
  • 多模态内容传输,支持 text、file、data 三种 Part 类型

与 MCP(Model Context Protocol)专注于 Agent-Tool 交互不同,A2A 聚焦于 Agent-Agent 协作,二者是互补关系。

二、A2A 核心概念解析

2.1 Agent Card:Agent 的"能力名片"

每个 A2A Server 在 /.well-known/agent.json 端点暴露一个 Agent Card JSON 文件,声明该 Agent 的能力、认证方式和交互协议版本。

{
  "name": "天气分析助手",
  "description": "提供全球天气数据分析、极端天气预警和气候趋势预测",
  "url": "https://weather-agent.example.com",
  "version": "1.0.0",
  "capabilities": {
    "streaming": true,
    "pushNotifications": true,
    "stateTransitionHistory": true
  },
  "authentication": {
    "schemes": ["bearer"]
  },
  "defaultInputModes": ["text"],
  "defaultOutputModes": ["text", "data"],
  "skills": [
    {
      "id": "weather-query",
      "description": "查询指定地区的历史和实时天气数据",
      "inputModes": ["text"],
      "outputModes": ["text", "data"],
      "examples": ["北京今天天气怎么样", "分析上海过去30天的降雨趋势"]
    }
  ]
}

2.2 Task:有状态协作的核心

A2A 以 Task 为基本通信单位,每个 Task 都有完整的生命周期。Task 的状态流转遵循以下规则:

  • submitted → working → completed
  • submitted → working → input-required(向请求方索要额外输入)
  • 任意状态 → failed(携带 error message)
  • 任意状态 → canceled

Task 的核心字段包括:

字段 说明
id Task 唯一标识(UUID)
contextId 关联上下文 ID,用于多轮对话
status 当前状态对象
history 交互历史(Message 数组)
artifacts 产出物(Agent 生成的最终结果)
metadata 自定义元数据

2.3 Message 与 Part

Message 代表一次通信回合,由多个 Part 组成。Part 有三种类型:

  • TextPart:纯文本内容
  • FilePart:文件(内联 base64 或 URI 引用)
  • DataPart:结构化数据(JSON Schema 定义)

这种设计使得 A2A 天然支持多模态交互,Agent 可以在一次响应中包含文本说明、数据图表和文件附件。

2.4 通信模式对比

A2A 支持四种通信模式:

  1. 请求-响应:标准 HTTP POST,同步等待结果
  2. 流式传输:通过 SSE (Server-Sent Events) 实时推送进度
  3. 推送通知:Agent 通过 Webhook 异步回调
  4. 轮询:客户端定期查询 Task 状态

生产环境中,长时间运行的任务应优先使用流式传输或推送通知模式,避免 HTTP 连接超时。

三、服务端实现:构建 A2A Server

下面使用 Python + FastAPI 构建一个完整的 A2A Server,支持流式传输和任务状态管理。

3.1 项目结构

a2a-server/
├── server.py          # FastAPI 应用入口
├── agent.py           # Agent 业务逻辑
├── tasks.py           # Task 存储与管理
├── models.py          # Pydantic 数据模型
└── agent_card.json    # Agent 能力声明

3.2 核心数据模型

from pydantic import BaseModel
from enum import Enum
from typing import List, Optional, Dict, Any
from uuid import uuid4

class TaskState(str, Enum):
    SUBMITTED = "submitted"
    WORKING = "working"
    COMPLETED = "completed"
    FAILED = "failed"
    CANCELED = "canceled"
    INPUT_REQUIRED = "input-required"

class TextPart(BaseModel):
    type: str = "text"
    text: str

class FilePart(BaseModel):
    type: str = "file"
    file: Dict[str, Any]  # {name, mimeType, bytes/buri}
    name: Optional[str] = None

class DataPart(BaseModel):
    type: str = "data"
    data: Dict[str, Any]

class Message(BaseModel):
    role: str  # "user" or "agent"
    parts: List[Dict[str, Any]]
    metadata: Optional[Dict[str, Any]] = None

class TaskStatus(BaseModel):
    state: TaskState
    message: Optional[Message] = None
    timestamp: Optional[str] = None

class Task(BaseModel):
    id: str
    contextId: str
    status: TaskStatus
    history: List[Message] = []
    artifacts: List[Any] = []
    metadata: Optional[Dict[str, Any]] = None

3.3 Task 管理器

from datetime import datetime, timezone
from typing import Optional, Dict
import asyncio
from collections import defaultdict

class TaskManager:
    def __init__(self):
        self._tasks: Dict[str, Task] = {}
        self._subscribers: Dict[str, list] = defaultdict(list)

    def create_task(self, context_id: str = None) -> Task:
        task_id = str(uuid4())
        ctx_id = context_id or str(uuid4())
        task = Task(
            id=task_id,
            contextId=ctx_id,
            status=TaskStatus(
                state=TaskState.SUBMITTED,
                timestamp=datetime.now(timezone.utc).isoformat()
            )
        )
        self._tasks[task_id] = task
        return task

    def get_task(self, task_id: str) -> Optional[Task]:
        return self._tasks.get(task_id)

    def update_status(self, task_id: str, status: TaskStatus):
        if task := self._tasks.get(task_id):
            task.status = status
            # 通知所有订阅者
            for callback in self._subscribers.get(task_id, []):
                callback(task)

    def add_message(self, task_id: str, message: Message):
        if task := self._tasks.get(task_id):
            task.history.append(message)

    def add_artifact(self, task_id: str, artifact: dict):
        if task := self._tasks.get(task_id):
            task.artifacts.append(artifact)

    async def subscribe(self, task_id: str, queue: asyncio.Queue):
        """SSE 订阅:将任务更新推送到异步队列"""
        if task_id in self._tasks:
            await queue.put(self._tasks[task_id])

3.4 Agent 执行引擎

import time
import json

class ResearchAgent:
    """研究型 Agent:执行多步研究任务"""

    async def execute(self, task_manager: TaskManager, task_id: str, query: str):
        """执行长时间研究任务"""
        try:
            # 标记为执行中
            task_manager.update_status(task_id, TaskStatus(
                state=TaskState.WORKING,
                message=Message(role="agent", parts=[{
                    "type": "text", 
                    "text": f"开始研究:{query}"
                }])
            ))

            # 步骤 1:信息收集
            await asyncio.sleep(1)  # 模拟耗时操作
            task_manager.add_message(task_id, Message(role="agent", parts=[
                {"type": "text", "text": "正在收集相关信息..."}
            ]))

            # 步骤 2:数据分析
            await asyncio.sleep(1.5)
            analysis_data = {
                "query": query,
                "sources": 42,
                "relevance_score": 0.94,
                "findings": [
                    "趋势分析显示年增长 23%",
                    "主要竞争对手市场份额下降"
                ]
            }

            # 步骤 3:生成结构化报告
            artifact = {
                "type": "data",
                "name": "research_report",
                "data": analysis_data
            }
            task_manager.add_artifact(task_id, artifact)

            # 步骤 4:生成摘要文本
            summary = Message(role="agent", parts=[
                {"type": "text", "text": f"研究完成。基于 {analysis_data['sources']} 个数据源,"
                                          f"相关度 {analysis_data['relevance_score']*100:.0f}%。"
                                          f"\\n\n关键发现:\n" + 
                                          "\n".join(f"- {f}" for f in analysis_data['findings'])}
                },
                {"type": "data", "data": analysis_data}
            ])
            task_manager.add_message(task_id, summary)

            # 标记完成
            task_manager.update_status(task_id, TaskStatus(
                state=TaskState.COMPLETED,
                message=summary
            ))

        except Exception as e:
            task_manager.update_status(task_id, TaskStatus(
                state=TaskState.FAILED,
                message=Message(role="agent", parts=[
                    {"type": "text", "text": f"执行失败:{str(e)}"}
                ])
            ))

3.5 FastAPI 服务端

from fastapi import FastAPI, HTTPException, Response
from fastapi.responses import StreamingResponse
from contextlib import asynccontextmanager
import asyncio
import json

app = FastAPI(title="A2A Research Agent Server")
task_manager = TaskManager()
agent = ResearchAgent()

@app.get("/.well-known/agent.json")
async def agent_card():
    """Agent Card 端点 - 服务发现"""
    return {
        "name": "研究分析助手",
        "description": "提供多维度数据分析、趋势研究和结构化报告生成",
        "url": "http://localhost:9999",
        "version": "1.0.0",
        "capabilities": {"streaming": True, "pushNotifications": True},
        "authentication": {"schemes": ["bearer"]},
        "defaultInputModes": ["text"],
        "defaultOutputModes": ["text", "data"],
        "skills": [{
            "id": "research",
            "description": "执行多步研究分析并生成结构化报告",
            "inputModes": ["text"],
            "outputModes": ["text", "data"]
        }]
    }

@app.post("/tasks")
async def create_and_run_task(request: dict):
    """创建任务并同步返回(非流式)"""
    task = task_manager.create_task()

    # 启动异步执行
    asyncio.create_task(agent.execute(
        task_manager, task.id, request.get("message", "")
    ))

    return task.model_dump()

@app.post("/tasks/{task_id}/send")
async def send_task(task_id: str, request: dict):
    """发送消息到已有 Task(支持 input-required 状态)"""
    task = task_manager.get_task(task_id)
    if not task:
        raise HTTPException(404, "Task not found")

    user_message = Message(role="user", parts=request.get("parts", []))
    task_manager.add_message(task_id, user_message)

    # 如果之前是 input-required,现在恢复执行
    if task.status.state == TaskState.INPUT_REQUIRED:
        asyncio.create_task(agent.execute(
            task_manager, task_id, request.get("message", "")
        ))

    return task.model_dump()

@app.get("/tasks/{task_id}")
async def get_task(task_id: str):
    """查询任务状态"""
    task = task_manager.get_task(task_id)
    if not task:
        raise HTTPException(404, "Task not found")
    return task.model_dump()

@app.post("/tasks/{task_id}/subscribe")
async def subscribe_task(task_id: str):
    """SSE 流式订阅任务状态变更"""
    if not task_manager.get_task(task_id):
        raise HTTPException(404, "Task not found")

    async def event_stream():
        queue = asyncio.Queue()
        await task_manager.subscribe(task_id, queue)

        while True:
            try:
                task = await asyncio.wait_for(queue.get(), timeout=30)
                data = json.dumps(task.model_dump(), ensure_ascii=False)
                yield f"event: task-update\ndata: {data}\n\n"

                if task.status.state in [TaskState.COMPLETED, TaskState.FAILED, TaskState.CANCELED]:
                    break
            except asyncio.TimeoutError:
                # 发送心跳保活
                yield "event: ping\ndata: {}\n\n"

    return StreamingResponse(
        event_stream(),
        media_type="text/event-stream",
        headers={"Cache-Control": "no-cache", "Connection": "keep-alive"}
    )

四、客户端实现:发现与调用 Agent

4.1 Agent 发现流程

import httpx

class A2ADiscovery:
    """Agent 发现客户端"""

    @staticmethod
    def discover(base_url: str) -> dict:
        """从 /.well-known/agent.json 获取 Agent Card"""
        with httpx.Client() as client:
            resp = client.get(f"{base_url}/.well-known/agent.json")
            resp.raise_for_status()
            return resp.json()

    @staticmethod
    def filter_skills(agent_card: dict, keyword: str) -> list:
        """按关键词筛选 Agent 技能集"""
        skills = agent_card.get("skills", [])
        return [s for s in skills if keyword.lower() in s.get("description", "").lower()]

# 使用示例
discovery = A2ADiscovery()
card = discovery.discover("http://localhost:9999")
print(f"发现 Agent: {card['name']}")
print(f"支持能力: {list(card['capabilities'].keys())}")

4.2 流式任务调用

import httpx
import json
import sseclient  # pip install sseclient-py

class A2AClient:
    """A2A 协议客户端 - 支持流式调用"""

    def __init__(self, base_url: str, token: str = None):
        self.base_url = base_url
        self.headers = {"Content-Type": "application/json"}
        if token:
            self.headers["Authorization"] = f"Bearer {token}"

    def send_task(self, message: str, context_id: str = None) -> dict:
        """发送任务并同步等待"""
        payload = {
            "message": message,
            "contextId": context_id
        }
        with httpx.Client() as client:
            resp = client.post(
                f"{self.base_url}/tasks",
                json=payload,
                headers=self.headers,
                timeout=120
            )
            resp.raise_for_status()
            return resp.json()

    def stream_task(self, message: str, context_id: str = None):
        """流式接收任务执行过程"""
        payload = {
            "message": message,
            "contextId": context_id
        }
        with httpx.Client() as client:
            # 1. 创建任务
            resp = client.post(
                f"{self.base_url}/tasks",
                json=payload,
                headers=self.headers,
                timeout=30
            )
            task = resp.json()
            task_id = task["id"]
            yield ("task_created", task)

            # 2. 订阅 SSE 流
            with client.stream(
                "POST",
                f"{self.base_url}/tasks/{task_id}/subscribe",
                headers=self.headers,
                timeout=120
            ) as resp:
                for line in resp.iter_lines():
                    if line.startswith("data:"):
                        data = json.loads(line[5:].strip())
                        yield ("update", data)

                        # 任务完成/失败时退出
                        state = data.get("status", {}).get("state", "")
                        if state in ("completed", "failed", "canceled"):
                            break

# 使用示例
client = A2AClient("http://localhost:9999")

for event_type, data in client.stream_task("分析 2026 年 AI Agent 市场趋势"):
    if event_type == "task_created":
        print(f"[Task Created] ID: {data['id']}")
    elif event_type == "update":
        state = data["status"]["state"]
        # 输出最新的 agent 消息
        if history := data.get("history"):
            latest = history[-1]
            if latest["role"] == "agent":
                parts = latest["parts"]
                for part in parts:
                    if part["type"] == "text":
                        print(f"[{state}] {part['text'][:100]}...")

4.3 多 Agent 协作编排

A2A 真正的威力在于多 Agent 协作。以下示例展示如何将研究 Agent 和报告 Agent 串联:

class AgentOrchestrator:
    """多 Agent 编排器 - 串联多个 A2A Agent"""

    def __init__(self, agent_registry: dict[str, str]):
        """
        agent_registry: {"researcher": "http://localhost:9999", 
                         "reporter": "http://localhost:9998"}
        """
        self.clients = {
            name: A2AClient(url) for name, url in agent_registry.items()
        }

    async def execute_pipeline(self, query: str):
        """执行研究 → 报告 生成流水线"""
        # 步骤 1:调用研究 Agent
        print("=== 调用研究 Agent ===")
        research_client = self.clients["researcher"]
        research_result = research_client.send_task(
            f"深度研究:{query}"
        )

        # 提取研究发现
        data_artifact = next(
            (a for a in research_result.get("artifacts", []) 
             if a.get("name") == "research_report"),
            None
        )
        if not data_artifact:
            raise ValueError("研究 Agent 未返回有效数据")

        findings = data_artifact["data"]["findings"]

        # 步骤 2:调用报告 Agent(输入来自研究 Agent 的输出)
        print("=== 调用报告 Agent ===")
        reporter_client = self.clients["reporter"]
        report_result = reporter_client.send_task(
            json.dumps({"findings": findings, "format": "markdown"})
        )

        return report_result

# 编排调用
orchestrator = AgentOrchestrator({
    "researcher": "http://localhost:9999",
    "reporter": "http://localhost:9998"
})
result = await orchestrator.execute_pipeline(
    "2026 年量子计算商业化进展分析"
)

五、安全与生产部署

5.1 认证与授权

A2A 支持多种认证方案:

from fastapi import Depends, HTTPException, status
from fastapi.security import HTTPBearer, HTTPAuthorizationCredentials

security = HTTPBearer()

async def verify_agent_token(
    credentials: HTTPAuthorizationCredentials = Depends(security)
):
    """验证 Agent 间调用令牌"""
    token = credentials.credentials

    # 生产环境:使用 JWT 验证
    # payload = jwt.decode(token, PUBLIC_KEY, algorithms=["RS256"])

    # 简单场景:校验预共享密钥
    if not validate_agent_token(token):
        raise HTTPException(
            status_code=status.HTTP_401_UNAUTHORIZED,
            detail="Agent authentication failed"
        )
    return token

@app.post("/tasks")
async def create_task_secure(
    request: dict,
    token: str = Depends(verify_agent_token)
):
    # 业务逻辑...
    pass

5.2 生产环境 Checklist

维度 要点
认证 mTLS + JWT 双因素认证,Agent 身份白名单
限流 按 Agent 调用方身份设置速率限制
超时 HTTP 连接 30s,SSE 心跳 15s,任务最大 10min
重试 幂等设计(contextId 去重),指数退避重试
监控 Task 状态分布、延迟 P99、错误率按 Agent 维度聚合
沙箱 Agent 间文件交换需过内容安全扫描
日志 全链路 contextId 追踪,保留 30 天

六、A2A vs MCP:互补而非竞争

很多工程师会困惑于 A2A 和 MCP 的关系。二者的定位不同:

  • MCP(Model Context Protocol):解决 Agent 如何与外部工具、数据源交互的问题。它是 Client-Server 架构,Agent 作为 MCP Client 调用 MCP Server 提供的工具(文件读取、数据库查询、API 调用等)。
  • A2A(Agent-to-Agent Protocol):解决 Agent 之间如何协作的问题。它是 Peer-to-Peer 架构,多个 Agent 通过 Task 机制进行有状态的异步协作。

实际生产系统中,二者经常配合使用:Agent A 通过 A2A 协议收到协作请求后,内部通过 MCP 协议调用工具完成子任务,再通过 A2A 将结果返回给请求方 Agent。这种分层架构实现了清晰的职责划分。

七、总结与展望

A2A 协议为构建可互操作的 Agent 生态系统提供了基础协议框架。其核心价值在于:

  1. 标准化:统一的任务生命周期和通信格式,降低集成成本
  2. 可扩展:通过 skills 机制支持任意业务逻辑
  3. 异步优先:原生支持长时间运行的任务和流式传输
  4. 安全可控:内置认证、Agent Card 发现和内容安全规范

随着 2026 年 AI Agent 在企业的规模化部署,Agent 间互操作将成为刚需。A2A 作为一种开放协议,正在被越来越多的 AI 平台采纳。建议工程师尽早熟悉其核心概念和实现模式,为多 Agent 系统的构建做好准备。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部