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→completedsubmitted→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 支持四种通信模式:
- 请求-响应:标准 HTTP POST,同步等待结果
- 流式传输:通过 SSE (Server-Sent Events) 实时推送进度
- 推送通知:Agent 通过 Webhook 异步回调
- 轮询:客户端定期查询 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 生态系统提供了基础协议框架。其核心价值在于:
- 标准化:统一的任务生命周期和通信格式,降低集成成本
- 可扩展:通过 skills 机制支持任意业务逻辑
- 异步优先:原生支持长时间运行的任务和流式传输
- 安全可控:内置认证、Agent Card 发现和内容安全规范
随着 2026 年 AI Agent 在企业的规模化部署,Agent 间互操作将成为刚需。A2A 作为一种开放协议,正在被越来越多的 AI 平台采纳。建议工程师尽早熟悉其核心概念和实现模式,为多 Agent 系统的构建做好准备。

发表评论 取消回复