🌐A2A 协议是什么?它解决了什么问题?¶
一句话:A2A 是一套让不同厂商、不同框架开发的 AI Agent 能够相互发现、分配任务、交换结果的标准化开放协议。

它要解决的四个核心痛点:
-
Agent 碎片化 现在做 Agent 的框架非常多(LangChain、AutoGen、CrewAI、Dify 等),每个框架的内部通信机制都是私有的。你想让 LangChain 写的 Agent 去调用一个 AutoGen 构建的专用 Agent?之前需要单独开发适配层。A2A 提供了一个“最大公约数”的通信模式,跨框架即插即用。
-
互操作性缺失 即便都是 HTTP API,A Agent 的任务定义(比如“分析这份财报”)和 B Agent 的任务定义(比如“financial_report_analysis”)大概率字段名、参数结构完全不同。A2A 定义了统一的任务对象(Task)、消息格式(Message)和工件(Artifact),让 Agent 之间能互相“听懂”。
-
长任务与异步协作 很多 Agent 任务不是秒级能返回的(比如生成一份几十页的尽调报告、训练一个模型)。HTTP 请求-响应模式撑不了这么久,需要长连接或异步回调。A2A 原生支持异步任务,有完整的 Task 状态机,支持长轮询、流式推送,甚至支持人工介入(Human-in-the-loop)。
-
能力发现 传统 API 需要人看文档。A2A 让 Orchestrator Agent 能够通过 Agent Card 自动发现远程 Worker Agent 具备哪些技能(Skills)、支持什么输入输出模式,从而实现动态的任务路由。这在多 Agent 协作中至关重要。
A2A 的生态定位:
MCP 解决的是 Agent 与工具/数据源的连接问题(Agent ↔ Tool),而 A2A 解决的是 Agent 与 Agent 之间的协作问题(Agent ↔ Agent)。两者互补,共同构成了 Agent 世界的“USB 协议栈”。
⏳A2A 协议的 Task 生命周期是什么?如何处理长时间运行的任务?¶
A2A 把一切工作都抽象为 Task。一个 Task 从诞生到结束,有严格定义的状态机,这对分布式协同至关重要。
Task 状态机图示:

状态详解:
-
Submitted:任务已被 Orchestrator 创建并提交给 Worker,但 Worker 还没开始处理。可以在此阶段取消。
-
Working:Worker 确认接收,正在执行。这是唯一可以持续输出中间状态的阶段(通过流式消息或 Artifact 更新)。
-
Completed:任务成功,产出物就绪。这是个终态,不可再变。
-
Failed:任务执行出错。终态,但可以附带错误信息。
-
Cancelled:在 Submitted 或 Working 阶段被主动取消。终态。
如何处理长时间运行的任务?
长任务(比如超过 30 秒)不能靠一个简单的 HTTP 请求-响应同步等待。A2A 提供了两种互补机制:
方式一:异步轮询(Polling)
Worker 收到任务后立即返回 taskId 和状态 Working。Orchestrator 端维护一个后台循环,定期用 tasks/get 查询状态,直到变成终态。
# Orchestrator 侧异步等待任务完成
async def wait_for_completion(worker_client, task_id, poll_interval=2.0):
while True:
task = await worker_client.get_task(task_id)
if task.status in ("completed", "failed", "cancelled"):
return task
# 可以在这里处理中间 Artifact,例如打印进度
if task.artifacts:
print(f"进度更新: {task.artifacts[-1].description}")
await asyncio.sleep(poll_interval)
方式二:流式推送(Streaming / SSE)
在支持流式传输的场景下,Worker 可以通过 tasks/send 或类似机制,主动向 Orchestrator 推送状态变更。这需要一个长连接(如 WebSocket、SSE 或 gRPC stream)。Orchestrator 不需要轮询,能实时看到进度条和中间产出。
方式三:Webhook 回调
如果 Worker 和 Orchestrator 彼此网络可达,可以在提交任务时附带一个 callbackUrl。Worker 在状态变为终态时,主动 POST 一次通知过去。这结合了异步的灵活性和较低的网络开销。
实战代码示例:模拟一个执行数据分析长任务的 Worker
# Worker 侧:处理长任务,逐步更新状态
class WorkerAgent:
async def process_task(self, task: Task):
# 1. 更新为 Working
await self.send_status(task.id, "working", "开始分析数据...")
# 2. 阶段一:数据清洗(耗时操作)
await asyncio.sleep(5) # 模拟耗时
await self.send_artifact(task.id, "清洗完成,发现 12 条异常数据")
# 3. 阶段二:建模
await asyncio.sleep(10)
await self.send_artifact(task.id, "模型训练完成,准确率 94.2%")
# 4. 完成
result = {"report_url": "https://storage/report_123.pdf"}
await self.send_status(task.id, "completed", "分析完成", artifacts=result)
对应的 Orchestrator 在接收到状态更新时,就能在 UI 上画出一个真实的进度条,而不是干等着。这种异步、分阶段执行能力,是 A2A 区别于传统 RPC 的一大特征。
🤝A2A 中 Orchestrator Agent 和 Worker Agent 如何协作?¶
A2A 定义了两种典型的 Agent 角色,它们的协作模式非常像“项目经理”和“领域专家”的关系。

协作流程四步走:
① 发现(Discovery)
Orchestrator 启动时或定期从一个注册中心(或本地配置)拉取所有可用 Worker 的 Agent Card。卡片里写明了 Worker 的 ID、支持的技能(如 financial_analysis、contract_review)、输入输出规范、端点地址等。这是动态路由的基础。
② 分配(Assignment)
用户说“帮我研究一下投资 AI 芯片的可行性,并拟一份合作框架”。Orchestrator 用自己的大脑分析:
-
子任务 1:市场研究与技术分析 → 路由给擅长
market_research的 Worker A。 -
子任务 2:法律框架起草 → 路由给擅长
legal_draft的 Worker B。 每个子任务被封装为一个 A2A Task,通过tasks/send发送给对应 Worker。
③ 协调与监控(Orchestration)
Orchestrator 不会发完任务就撒手不管。它维护一个任务依赖图:Task 2 可能依赖 Task 1 的数据(市场分析是拟定合同的依据)。它会:
-
等待 Task 1 完成,提取关键结论。
-
将结论作为 Task 2 的输入参数,再下发给 Worker B。
-
如果 Worker A 失败,它可以重新分配给另一个同技能 Worker,或者要求 Worker A 重试。
④ 聚合(Aggregation)
所有子任务终态后,Orchestrator 把各 Worker 返回的 Artifacts 组合成一份完整答案,推送给用户。比如:市场分析报告 PDF + 合同草稿 DOCX + 风险提示摘要。
示例代码:Orchestrator 如何动态分配任务
class OrchestratorAgent:
def __init__(self, registry):
self.registry = registry # 可用 Worker 列表及其 Agent Card
async def handle_request(self, user_query):
# 1. 用自己的 LLM 做任务分解
subtasks = await self.decompose(user_query)
# 分解结果如:[("market_research", "AI芯片市场前景"), ("legal_draft", "拟定合作框架")]
results = []
for skill, description in subtasks:
# 2. 从注册中心找最适合的 Worker
worker = self.registry.find_best(skill)
if not worker:
results.append(f"未找到支持 {skill} 的 Worker")
continue
# 3. 创建并发送任务
task = Task(
description=description,
context={"previous_results": results} # 携带前置结果
)
task_id = await worker.send_task(task)
# 4. 异步等待完成(或轮询)
final_task = await self.wait_for_completion(worker, task_id)
if final_task.status == "completed":
results.append(final_task.artifacts)
else:
# 5. 失败处理:换 Worker 或重试
results.append(await self.handle_failure(skill, description, worker))
# 6. 聚合所有结果
return self.aggregate(results, user_query)
协作中三个关键工程点:
-
输入/输出对齐:Orchestrator 需要确保前一个 Task 的输出是下一个 Task 的合法输入。这通常依靠对 Artifacts 的内容校验和 schema 匹配。
-
容错与重试:如果 Worker 返回
Failed,Orchestrator 可以决定是用同样参数重试、换一个 Worker,还是降级处理(例如人工介入)。 -
安全与信任边界:A2A 强烈建议使用 TLS + 服务账号认证。Orchestrator 应该验证 Worker 的身份,Worker 也应该验证 Orchestrator 是否有权提交这种任务。这种双向信任是生产化的前提。
总结一句:
Orchestrator 是“用什么、怎么做、找谁做”的决策者,Worker 是“做得好、做得快、做得专”的执行者。A2A 把这两者之间的接口固定下来,让所有 Agent 都能在一个规则下自由组队,就像乐高的卡扣一样,任意拼接,牢固协作。
🤝 A2A 和 MCP 在实际项目中如何配合使用?¶
一句话:MCP 让 Agent 能“用工具”,A2A 让 Agent 能“找帮手”。两者在项目里是典型的互补关系,一个向下连接万物,一个横向连接彼此。
定位差异:

配合方式:Orchestrator 同时使用 MCP 和 A2A
在一个复杂的企业项目中,Orchestrator Agent 通常同时持有 MCP Client 和 A2A Client:
-
用户请求:帮我分析上季度的销售数据,并给每个区域经理发一份个性化总结邮件。
-
Orchestrator 拆解任务:
- 子任务 A:从数据仓库拉取销售数据 → 用 MCP 调用
query_database工具。 - 子任务 B:根据数据生成文本分析 → 用 A2A 派发给“数据分析师 Agent”。
-
子任务 C:根据分析结果和区域名单发送邮件 → 用 A2A 派发给“邮件助手 Agent”。
-
执行过程:Orchestrator 先 MCP 拿数据,再通过 A2A 把数据和分析要求发给 Worker Agent,Worker Agent 内部可能再用自己的 MCP 工具(比如发送邮件调 SMTP Server)。
架构图:

关键点:
-
MCP 负责“原子操作”:一次数据库查询、一次文件读取、一次 API 调用。这些是确定性的、无状态的工具。
-
A2A 负责“复合任务”:把一个需要多步推理、有独立上下文的任务交给另一个 Agent。这个 Agent 可能自己也有 MCP 工具。
-
边界:如果一个操作是简单的、一次性的,用 MCP 就够了。如果一个操作需要另一个 Agent 的“大脑”参与(比如需要理解复杂文档并做出判断),用 A2A。
代码骨架:Orchestrator 同时用 MCP 和 A2A
class Orchestrator:
def __init__(self, mcp_client, a2a_registry):
self.mcp = mcp_client
self.a2a = a2a_registry
async def handle_request(self, user_query):
# 1. 用 MCP 取数据
sales_data = await self.mcp.call_tool("query_sales", {"quarter": "Q2"})
# 2. 拆解任务,发现可用 Agent
analyst = self.a2a.find_agent(skill="data_analysis")
mailer = self.a2a.find_agent(skill="email_campaign")
# 3. 用 A2A 派发子任务
analysis_task = Task(description="分析销售数据", context={"data": sales_data})
analysis_result = await self.a2a.send_and_wait(analyst, analysis_task)
# 4. 继续派发下一个 A2A 任务
email_task = Task(description="发送个性化邮件", context={"analysis": analysis_result})
final = await self.a2a.send_and_wait(mailer, email_task)
return final
配合的本质:MCP 是 Agent 的“手和脚”,A2A 是 Agent 的“同事”。 没有手脚,大脑空转;没有同事,复杂项目一个人扛不住。两者结合,才构成完整的 Agent 企业级架构。
🪪Agent Card 是什么?如何设计一个好的 Agent Card?¶
Agent Card 是每个 Worker Agent 对外公开的“名片 + 简历”,用结构化 JSON 描述自己的身份、技能、接口和约束。 Orchestrator 通过读取它来决定“这个 Agent 能帮我做什么,怎么跟它通信”。
一个标准的 Agent Card 包含六个核心字段:
{
"name": "Financial Analyst Agent",
"description": "专精于财务报表分析、趋势预测和异常检测。",
"url": "https://agents.company.com/financial-analyst/a2a",
"version": "2.1.0",
"capabilities": {
"skills": [
{
"id": "financial_statement_analysis",
"description": "分析资产负债表、利润表和现金流量表,提取关键指标并给出评分。",
"inputSchema": { ... },
"outputSchema": { ... },
"examples": [ ... ]
}
],
"streaming": true,
"pushNotifications": false,
"maxTaskDuration": 300
},
"provider": {
"name": "Finance Team",
"contact": "fin-ai@company.com"
},
"security": {
"authentication": "oauth2",
"authorization": "role:finance-readonly"
}
}
设计一个好 Agent Card 的五个原则:
-
描述要“对人友好,对 LLM 更友好”
description不是写给开发者看的注释,它会被 LLM 读到来做路由决策。所以要具体、无歧义。例如,不要只写“分析数据”,要写“分析 CSV/Excel 格式的销售数据,输出同比/环比增长率和异常检测报告”。 -
技能定义要精确,包含 Schema 和示例 每个 Skill 必须有清晰的
inputSchema和outputSchema(JSON Schema 格式),最好附带 1-2 个输入输出示例。这让 Orchestrator 的 LLM 能准确组装调用参数。 -
声明非功能性能力
streaming:是否支持流式返回中间结果。pushNotifications:是否支持主动推送状态更新(如 Webhook)。-
maxTaskDuration:该 Agent 处理任务的最长耗时(秒)。Orchestrator 可以据此设置超时。 -
版本号是必须的 技能会迭代,
version字段让 Orchestrator 知道当前连接的 Agent 是哪个版本,便于做兼容性检查或技能降级。 -
安全边界要清晰
- 认证方式(OAuth2, API Key, mTLS)。
- 授权范围(这个 Agent 能处理哪些数据等级、能调用哪些下游服务)。这决定了它是否适合处理包含敏感信息的任务。
好的 Agent Card 应该让 Orchestrator 在没有文档的情况下,只通过这张卡片就能正确路由任务。 就像你招人,看完简历就知道他能不能胜任这个岗位。在工程上,通常由一个中心注册服务(如 Consul + 自定义元数据)来存储和管理所有 Agent Card,Orchestrator 启动时或定期拉取。
🔄多 Agent 系统中如何处理 Agent 之间的循环依赖和死锁?¶
在多 Agent 协作中,循环依赖指的是:Agent A 等 Agent B 完成,Agent B 又等 Agent A 完成,两者永远等下去。
典型场景:
Orchestrator 把任务 T1 分配给 Agent A(需要知识库 X),T2 分配给 Agent B(负责更新知识库 X)。结果 Agent A 需要等 Agent B 更新完知识库才能开始,Agent B 又被安排等 Agent A 给出一些初步分析后才能决定更新什么。双方互相等待,形成死锁。
三个层次的解决方案:
① 设计层:任务 DAG(有向无环图)强制非循环
Orchestrator 在生成任务计划时,必须构建一个有向无环图(DAG),并通过拓扑排序检测环。如果存在环,拒绝该计划,要求重新分解。
def check_cycle(tasks: List[Task]) -> bool:
"""检查任务依赖图中是否有环"""
from collections import defaultdict
graph = defaultdict(list)
for t in tasks:
for dep in t.depends_on:
graph[t.id].append(dep)
# 拓扑排序
indegree = {t.id: 0 for t in tasks}
for u in graph:
for v in graph[u]:
indegree[v] = indegree.get(v, 0) + 1
q = [n for n, deg in indegree.items() if deg == 0]
visited = 0
while q:
node = q.pop(0)
visited += 1
for neighbor in graph[node]:
indegree[neighbor] -= 1
if indegree[neighbor] == 0:
q.append(neighbor)
return visited != len(indegree)
# 如果返回 True,说明有环,需要人工介入或自动 re-plan
② 运行层:超时打破 + 降级策略
即使计划是无环的,也可能因为执行时的消息乱序或资源争用而出现实际上的相互等待(比如两方都在等对方先释放资源)。因此必须设置任务级超时和资源等待超时。
-
如果一个 Agent 等待某个依赖资源超过 30 秒,主动触发超时异常,回退到 Orchestrator。
-
Orchestrator 收到超时后,执行降级:例如将任务重新分配给另一个 Agent(如果资源无冲突),或者使用本地缓存/默认值,或者直接向用户报告“当前无法完成”。
async def wait_for_dependency(dep_future, timeout=30):
try:
return await asyncio.wait_for(dep_future, timeout=timeout)
except asyncio.TimeoutError:
raise DeadlockException("依赖任务超时,可能存在循环等待")
③ 架构层:集中式调度 + 资源锁
在多 Agent 系统中,最不容易死锁的调度方式是星型拓扑:所有任务都由中心 Orchestrator 发起和协调,Worker 之间不直接互调。如果必须让 Worker 之间 A2A 直接通信,那么必须有一个全局的资源管理器来分配“锁”。
-
资源锁:如果任务需要独占某个知识库的写权限,Orchestrator 为它分配一个带有超时释放的分布式锁。
-
死锁检测线程:定期扫描所有未完成的任务和它们持有的锁/等待的依赖,构建等待图,如果发现环,强制终止其中优先级最低的那个任务,释放锁并通知相关方。
核心思想:绝不假设分布式任务会按理想顺序完成。 通过 DAG 设计避免显式环,通过超时和降级打破隐式环,通过集中调度减少不必要的点对点耦合。这是多 Agent 系统从 Demo 到 Production 的关键一步。
🔄 技能版本管理和热更新如何实现?¶
技能版本管理要解决的问题:Worker Agent 的技能升级了(输入参数变了、行为逻辑改了),不能一次性中断所有正在运行的任务。要做到平滑过渡、新旧兼容。
设计蓝图:

四步实现:
① 注册中心存储多版本 Agent Card
Worker 启动时向注册中心注册自己支持的多个版本。每个版本对应一组完整的技能定义。
{
"agentId": "financial-analyst",
"versions": [
{
"version": "1.0.0",
"url": "https://worker:9090/v1/a2a",
"skills": [...]
},
{
"version": "1.1.0",
"url": "https://worker:9090/v2/a2a",
"skills": [...] // 可能技能名相同但 inputSchema 不同
}
],
"defaultVersion": "1.1.0"
}
② 技能请求携带版本约束
Orchestrator 在分配任务时,可以指定所需技能的最低版本、精确版本或版本范围。
# 任务定义时带版本要求
task = Task(
description="分析财报",
required_skills=["financial_statement_analysis"],
version_constraint=">=1.0.0" # 最低版本
)
③ 路由时选择合适版本
Orchestrator 从注册中心拉取该 Agent 的所有版本,根据约束和当前灰度策略选择最合适的。
def select_agent_version(agent_card, constraint, rollout_policy):
valid_versions = [v for v in agent_card.versions
if version_satisfies(v.version, constraint)]
if rollout_policy.get("canary"):
# 灰度:10% 流量走新版本
if random.random() < 0.1:
return find_version(valid_versions, rollout_policy["canary_version"])
# 否则选最新的稳定版
return valid_versions[-1]
④ 热更新与旧任务排空
-
不中断正在运行的任务:更新版本时,旧的 Worker 实例继续保持服务,直到其上的所有长任务完成。
-
新任务引流:新启动的 Worker 实例注册新版本,Orchestrator 的流量逐渐切过去。
-
优雅关闭:旧 Worker 收到 SIGTERM 后,不再接受新任务,等待现有任务完成(或超时强制结束),然后退出。
-
客户端无感:Orchestrator 连接的是注册中心,它根据版本策略动态选择 URL,Worker 的启停对它透明。
代码骨架:Worker 侧热更新
class Worker:
def __init__(self, version):
self.version = version
self.active_tasks = set()
async def start(self):
await self.register_to_registry(self.version) # 上报新版本
# 启动旧版本排空等待...
async def shutdown(self):
# 1. 取消注册,新任务不再路由过来
await self.deregister()
# 2. 等待所有进行中任务完成,最多等 60 秒
await asyncio.wait_for(
asyncio.gather(*self.active_tasks, return_exceptions=True),
timeout=60
)
# 3. 强制终止剩下未完成的任务并告警
for task in self.active_tasks:
task.cancel()
灰度发布策略:
-
Canary:先让 5% 的新任务路由到新版本,观察 30 分钟无报错,逐步扩大到 50%、100%。
-
按租户/场景灰度:更精细的做法,指定特定测试租户使用新版本,生产租户继续用旧版本。
总结: 技能版本化让 Agent 的迭代从“心惊胆战的停机更新”变成了“从容不迫的平滑演进”。注册中心是多版本的目录,Orchestrator 是灵活的路由器,Worker 的优雅关闭则是保证任务零丢失的最后一道门。这一套机制,正是多 Agent 系统可持续进化的根基。