跳转至

🌐A2A 协议是什么?它解决了什么问题?

一句话:A2A 是一套让不同厂商、不同框架开发的 AI Agent 能够相互发现、分配任务、交换结果的标准化开放协议。

image.png

它要解决的四个核心痛点:

  1. Agent 碎片化 现在做 Agent 的框架非常多(LangChain、AutoGen、CrewAI、Dify 等),每个框架的内部通信机制都是私有的。你想让 LangChain 写的 Agent 去调用一个 AutoGen 构建的专用 Agent?之前需要单独开发适配层。A2A 提供了一个“最大公约数”的通信模式,跨框架即插即用。

  2. 互操作性缺失 即便都是 HTTP API,A Agent 的任务定义(比如“分析这份财报”)和 B Agent 的任务定义(比如“financial_report_analysis”)大概率字段名、参数结构完全不同。A2A 定义了统一的任务对象(Task)、消息格式(Message)和工件(Artifact),让 Agent 之间能互相“听懂”。

  3. 长任务与异步协作 很多 Agent 任务不是秒级能返回的(比如生成一份几十页的尽调报告、训练一个模型)。HTTP 请求-响应模式撑不了这么久,需要长连接或异步回调。A2A 原生支持异步任务,有完整的 Task 状态机,支持长轮询、流式推送,甚至支持人工介入(Human-in-the-loop)。

  4. 能力发现 传统 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 状态机图示:

image.png

状态详解:

  • 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 角色,它们的协作模式非常像“项目经理”和“领域专家”的关系。

image.png

协作流程四步走:

① 发现(Discovery) Orchestrator 启动时或定期从一个注册中心(或本地配置)拉取所有可用 Worker 的 Agent Card。卡片里写明了 Worker 的 ID、支持的技能(如 financial_analysiscontract_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 能“找帮手”。两者在项目里是典型的互补关系,一个向下连接万物,一个横向连接彼此。

定位差异:

image.png

配合方式:Orchestrator 同时使用 MCP 和 A2A

在一个复杂的企业项目中,Orchestrator Agent 通常同时持有 MCP Client 和 A2A Client:

  1. 用户请求:帮我分析上季度的销售数据,并给每个区域经理发一份个性化总结邮件。

  2. Orchestrator 拆解任务:

  3. 子任务 A:从数据仓库拉取销售数据 → 用 MCP 调用 query_database 工具。
  4. 子任务 B:根据数据生成文本分析 → 用 A2A 派发给“数据分析师 Agent”。
  5. 子任务 C:根据分析结果和区域名单发送邮件 → 用 A2A 派发给“邮件助手 Agent”。

  6. 执行过程:Orchestrator 先 MCP 拿数据,再通过 A2A 把数据和分析要求发给 Worker Agent,Worker Agent 内部可能再用自己的 MCP 工具(比如发送邮件调 SMTP Server)。

架构图:

image.png

关键点:

  • 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 的五个原则:

  1. 描述要“对人友好,对 LLM 更友好” description 不是写给开发者看的注释,它会被 LLM 读到来做路由决策。所以要具体、无歧义。例如,不要只写“分析数据”,要写“分析 CSV/Excel 格式的销售数据,输出同比/环比增长率和异常检测报告”。

  2. 技能定义要精确,包含 Schema 和示例 每个 Skill 必须有清晰的 inputSchemaoutputSchema(JSON Schema 格式),最好附带 1-2 个输入输出示例。这让 Orchestrator 的 LLM 能准确组装调用参数。

  3. 声明非功能性能力

  4. streaming:是否支持流式返回中间结果。
  5. pushNotifications:是否支持主动推送状态更新(如 Webhook)。
  6. maxTaskDuration:该 Agent 处理任务的最长耗时(秒)。Orchestrator 可以据此设置超时。

  7. 版本号是必须的 技能会迭代,version 字段让 Orchestrator 知道当前连接的 Agent 是哪个版本,便于做兼容性检查或技能降级。

  8. 安全边界要清晰

  9. 认证方式(OAuth2, API Key, mTLS)。
  10. 授权范围(这个 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 的技能升级了(输入参数变了、行为逻辑改了),不能一次性中断所有正在运行的任务。要做到平滑过渡、新旧兼容。

设计蓝图:

image.png

四步实现:

① 注册中心存储多版本 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 系统可持续进化的根基。