Clipping 微信公众号

A2A中的异步任务是如何进行的呢?

by 柠遇AI纪元 原文 ↗
Created: 2026-06-15

公众号名称:柠遇AI纪元

作者名称:柠遇AI纪元

发布时间:2026-06-15 00:47

A2A(Agent-to-Agent)协议系统性详解

1. 什么是 A2A 协议

这个A2A协议的定义来自google的官网,通俗来说是定义Agent与Agent之间的通信协议。用官方的话来说,A2A(Agent-to-Agent)协议是一个开放标准,旨在促进独立、可能不透明的 AI Agent 系统之间的通信和互操作性。在一个 Agent 可能由不同框架、语言构建或由不同供应商提供的生态系统中,他们彼此之间不了解彼此内部的实现,而A2A 提供了一种通用的语言和交互模型。

A2A 协议的主要目标如下:

目标说明
互操作性补充了不同 Agentic 系统之间的通信障碍
协作使 Agent 能够委派任务、交换上下文并协同工作
发现允许 Agent 动态发现和理解其他 Agent 的能力
灵活性支持同步、流式、异步推送等多种交互模式
安全性基于标准 Web 安全实践
异步优先原生支持长时间运行异步任务场景

2. 核心设计原则

A2A 协议的设计遵循以下指导原则:

简单性:复用现有的、被广泛理解的标准(HTTP、JSON-RPC 2.0、Server-Sent Events)

安全性:解决了认证、授权、安全、隐私、追踪和监控问题

异步优先:专门为可能非常长时间运行的异步任务而设计

支持多模态:支持交换多种内容类型,包括文本、音视频、结构化数据/表单等多模态数据

不透明执行:Agent 之间基于声明的能力和交换的信息进行协作,无需共享内部想法、计划或工具实现。注意 MCP 关注的是 Agent 与工具/资源的连接,而 A2A 关注的是 Agent 与 Agent 之间的互操作。

3. 整体架构

A2A 从总体上来说,还是遵循客户端/服务端的架构。具体来说,在服务端侧,主要采用采用分层方式,从入口路由层出发,依次到请求处理层、任务生命周期层、事件系统,最后到持久化层。客户端侧比较简单些,通过 ClientFactory 创建,从http、rpc、sse三种协议中选择合适的协议,将agent异步任务或者消息发送到服务端。

具体来看下服务端的一些组件:

3.1 接入与API网关层

主要是API网关、负载均衡LB组成。主要的作用有:

  • 统一入口:接收来自源系统的请求,隐藏后端复杂性

  • 安全鉴权:负责OAuth2、JWT校验、API Key验证及白名单过滤等

  • 流量控制:实现限流、熔断,保护内部系统

3.2 核心处理与路由层

主要是一些处理一些异步任务的服务以及规则引擎。主要是作用:

  • 协议转换:将客户端传过来的的协议格式转换为服务端内部所需格式

  • 动态路由:根据请求头或内容,决定将任务分发到哪个目标队列或服务

3.3 异步任务与消息中间件层

主要是利用消息队列中间件来处理各种异步任务的,记录任务状态以及结果。主要作用如下:

  • 削峰填谷:在瞬时高并发时,将请求缓存在队列中,由后端慢慢消费

  • 解耦:客户端发送完消息即可返回,无需等待服务端处理完毕

3.4 目标系统适配层

对各类消费者进行适配的中间层。主要是负责把处理后的数据,以目标系统接受的方式(比如WebHook、SFTP、DB写入、消息队列等方式)投递出去。

3.5 管理、监控与重试控制层

这个主要是对于整个系统的监控与日志收集。主要的功能有:

  • 异常处理:管理失败的任务,进行指数退避重试。

  • 全链路监控:监控任务的延迟、吞吐量和堆积情况。

4. 核心数据模型

A2A 的数据模型主要包括以下内容:

  • AgentCard

    Agent 的身份描述,包含名称、技能、支持的传输接口、安全方案等。

  • AgentInterface

    Agent之间绑定好的传输协议(JSONRPC/HTTP+JSON/GRPC)、协议版本、URL 和可选的租户信息。

  • Message

    通信的主要载体,包含角色(USER/AGENT)、内容部件列表、消息ID。

  • Part

    内容单元,支持 text(字符串)、data(Protobuf Value)、raw(字节)、url(字符串)四种类型。

  • Task

    Agent工作的基本单元,任务跟踪对象,包含 ID、状态、制品列表、历史消息。

  • TaskState

    任务状态枚举,定义了完整的生命周期。

5. 发送异步任务消息

5.1. Send Message(发送消息)

这是启动 Agent 交互的主要操作。客户端发送消息,Agent 返回一个跟踪处理过程的 Task,或者针对简单交互直接返回 Message 响应。

输入SendMessageRequest(包含 message、configuration、metadata)

输出TaskMessage

SendMessageConfiguration中的关键配置项:

配置项说明
returnImmediately若为 true,服务器立即返回 TASK_STATE_SUBMITTED状态的任务,不等待完成
acceptedOutputModes客户端接受的输出媒体类型列表
historyLength返回的历史消息最大数量
taskPushNotificationConfig随请求一起设置的推送通知配置

5.2. Send Streaming Message(发送流式消息)

与 Send Message 类似,但建立流式连接以实时接收更新。流式响应遵循以下模式之一:

  • 纯消息流:若 Agent 返回 Message,流中只包含一个 Message 对象然后立即关闭。

  • 任务生命周期流:若 Agent 返回 Task,流以 Task 对象开始,随后是零个或多个 TaskStatusUpdateEventTaskArtifactUpdateEvent,当任务到达终态时流关闭。

6. Agent Card 与服务发现

这里我们将重点探讨异步任务在A2A中是如何进行的。

6.1. Agent Card 结构

Agent Card 是 A2A 服务发现的基础,是一个 JSON 元数据文档,描述了 Agent 的身份、能力、技能和认证要求。注意这里和我们日常工作过程中遇到的微服务注册与发现有点类似,都是基于生产者消费者的设计模式来设计的;微服务注册与发现是将相关的服务放到网关上,其他的服务就可以找到对应的ip,进行转发接口任务,而A2A也是类似的,只不过是在任务队列中获取并执行异步任务而已。

以下是一个完整的 Agent Card 示例:

{
  "name": "GeoSpatial Route Planner Agent",
"description": "Provides advanced route planning, traffic analysis, and custom map generation services.",
"version": "1.2.0",
"provider": {
    "organization": "Example Geo Services Inc.",
    "url": "https://www.examplegeoservices.com"
  },
"supportedInterfaces": [
    {
      "url": "https://georoute-agent.example.com/a2a/v1",
      "protocolBinding": "JSONRPC",
      "protocolVersion": "1.0"
    },
    {
      "url": "https://georoute-agent.example.com/a2a/grpc",
      "protocolBinding": "GRPC",
      "protocolVersion": "1.0"
    },
    {
      "url": "https://georoute-agent.example.com/a2a/json",
      "protocolBinding": "HTTP+JSON",
      "protocolVersion": "1.0"
    }
  ],
"capabilities": {
    "streaming": true,
    "pushNotifications": true,
    "extendedAgentCard": true
  },
"securitySchemes": {
    "google": {
      "openIdConnectSecurityScheme": {
        "openIdConnectUrl": "https://accounts.google.com/.well-known/openid-configuration"
      }
    }
  },
"security": [{ "google": ["openid", "profile", "email"] }],
"defaultInputModes": ["application/json", "text/plain"],
"defaultOutputModes": ["application/json", "image/png"],
"skills": [
    {
      "id": "route-optimizer-traffic",
      "name": "Traffic-Aware Route Optimizer",
      "description": "Calculates the optimal driving route considering real-time traffic.",
      "tags": ["maps", "routing", "navigation", "traffic"],
      "examples": [
        "Plan a route from Mountain View to SFO avoiding tolls."
      ],
      "inputModes": ["application/json", "text/plain"],
      "outputModes": ["application/json", "text/html"]
    }
  ]
}

6.2. 服务发现机制

从A2A官方的源码中可以看到,客户端可以通过以下方式找到 Agent Card:

  1. Well-Known URI:访问 https://{server_domain}/.well-known/agent-card.json

  2. 注册表/目录:查询 Agent 的策划目录

  3. 直接配置:预配置的 Agent Card URL 或内容

6.3. 协议选择规则

客户端在获取 Agent Card 后,必须遵循以下规则选择协议:

  • 解析 supportedInterfaces列表,选择第一个本地支持的传输协议。

  • supportedInterfaces中越靠前的条目优先级越高(代表 Agent 的偏好)。

  • 使用所选传输协议对应的 URL。

  • AgentInterface中设置了 tenant字段,客户端必须在所有请求消息的 tenant字段中包含该值。

6.4. Agent Card 签名

Agent Card 可以使用 JSON Web Signature 进行数字签名,以确保真实性和完整性。签名前需使用 JSON Canonicalization Scheme 对 Card 进行规范化处理。

7. 任务状态机与生命周期

异步任务的执行过程如下:

A2A 任务遵循严格的状态机,TaskState定义了所有可能的状态:

状态类型说明
TASK_STATE_SUBMITTED正常任务已成功提交并确认
TASK_STATE_WORKING正常任务正在被 Agent 积极处理
TASK_STATE_COMPLETED终态任务已成功完成
TASK_STATE_FAILED终态任务以错误结束
TASK_STATE_CANCELED终态任务在完成前被取消
TASK_STATE_REJECTED终态Agent 决定不执行该任务
TASK_STATE_INPUT_REQUIRED中断Agent 需要额外的用户输入才能继续
TASK_STATE_AUTH_REQUIRED中断需要认证授权才能继续

终态任务不接受进一步的消息。中断状态的任务可以通过客户端发送带有相同 taskId的新消息来恢复。

8. 任务更新传递机制

A2A 提供了三种互补的机制,供客户端接收任务进度和完成情况的更新:

机制操作优点缺点适用场景
轮询Get Task实现简单,适用于所有绑定延迟高,可能产生不必要的请求简单集成、更新不频繁、客户端在防火墙后
流式传输Send Streaming Message / Subscribe to Task低延迟,实时更新需要持久连接支持交互式应用、实时仪表板、进度监控
推送通知推送通知配置无需持久连接,异步交付客户端需可通过 HTTP 访问服务器间集成、长时间运行任务、事件驱动架构

推送通知的工作原理:客户端注册 Webhook URL 后,当任务状态发生变化时,Agent 服务器会向该 URL 发送 HTTP POST 请求,payload 为 StreamResponse格式的 JSON,并在 header 中携带 X-A2A-Notification-Token用于验证。

9. 认证与授权

A2A 的身份信息应该写到Gateway网关层做进一步处理,下发任务的时候,在请求头内添加相关安全header就可以了,而非在 A2A 语义内处理。

9.1. 支持的安全方案

安全方案说明
APIKeySecuritySchemeAPI 密钥认证
HTTPAuthSecuritySchemeHTTP 认证
OAuth2SecuritySchemeOAuth 2.0
OpenIdConnectSecuritySchemeOpenID Connect
MutualTlsSecurityScheme双向 TLS 认证

9.2. 客户端认证流程

  1. 发现要求:客户端通过 Agent Card 的 securitySchemes字段发现服务器要求的认证方案。

  2. 凭证获取:客户端通过特定于所需认证方案的带外流程获取必要凭证。

  3. 凭证传输:客户端在每个 A2A 请求的协议特定 header 或元数据中包含这些凭证。

9.3. 任务内授权

在执行任务过程中,Agent 可能需要授权才能执行操作(如调用外部 API)。A2A 提供了 TASK_STATE_AUTH_REQUIRED状态来处理这种情况:

  • Agent 将任务状态转换为 TASK_STATE_AUTH_REQUIRED,并在状态消息中说明所需的授权。

  • 客户端收到此状态后,需要通过带外方式或扩展协商的带内方式提供凭证。

  • 如果客户端本身也是一个 A2A Agent,它可以进一步将授权请求委托给其自己的客户端,形成授权请求链。

10. 协议绑定

A2A 官方支持三种核心协议绑定,所有绑定在功能上等效。

10.1. 方法映射参考

功能JSON-RPC 方法gRPC 方法REST 端点
发送消息SendMessageSendMessagePOST /message:send
发送流式消息SendStreamingMessageSendStreamingMessagePOST /message:stream
获取任务GetTaskGetTaskGET /tasks/{id}
列出任务ListTasksListTasksGET /tasks
取消任务CancelTaskCancelTaskPOST /tasks/{id}:cancel
订阅任务SubscribeToTaskSubscribeToTaskPOST /tasks/{id}:subscribe
创建推送配置CreateTaskPushNotificationConfig同左POST /tasks/{id}/pushNotificationConfigs
获取扩展 CardGetExtendedAgentCardGetExtendedAgentCardGET /extendedAgentCard

10.2. JSON-RPC 绑定

  • 协议:JSON-RPC 2.0 over HTTP(S)

  • Content-Typeapplication/json

  • 流式传输:Server-Sent Events (text/event-stream)

  • 服务参数传递:通过 HTTP 请求 header(如 A2A-VersionA2A-Extensions

// 请求示例
POST /rpc HTTP/1.1
Content-Type: application/json
A2A-Version: 1.0

{
"jsonrpc": "2.0",
"id": 1,
"method": "SendMessage",
"params": {
    "message": {
      "role": "ROLE_USER",
      "messageId": "msg-uuid",
      "parts": [{"text": "What is the weather today?"}]
    }
  }
}

// 响应示例
{
"jsonrpc": "2.0",
"id": 1,
"result": {
    "task": {
      "id": "task-uuid",
      "contextId": "context-uuid",
      "status": {"state": "TASK_STATE_COMPLETED"}
    }
  }
}

10.3. gRPC 绑定

  • 协议:gRPC over HTTP/2 with TLS

  • 序列化:Protocol Buffers v3

  • 服务定义:实现 A2AServicegRPC 服务

  • 服务参数传递:通过 gRPC metadata(header)

  • 流式传输:原生 gRPC server streaming RPC

10.4. HTTP+JSON/REST 绑定

  • 协议:HTTP(S) with JSON payloads

  • Content-Typeapplication/a2a+json

  • 流式传输:Server-Sent Events

  • 服务参数传递:通过 HTTP 请求 header

// 流式响应示例 (SSE)
HTTP/1.1 200 OK
Content-Type: text/event-stream

data: {"task": {"id": "task-uuid", "status": {"state": "TASK_STATE_WORKING"}}}

data: {"artifactUpdate": {"taskId": "task-uuid", "artifact": {"parts": [{"text": "# Report\n\n"}]}}}

data: {"statusUpdate": {"taskId": "task-uuid", "status": {"state": "TASK_STATE_COMPLETED"}}}

11. 完整示例:构建一个 Hello World Agent

通过阅读a2a-pythonSDK的使用介绍,以下是基于它构建一个完整 A2A Agent 服务的示例,同时支持 JSON-RPC、HTTP+JSON 和 gRPC 三种传输协议。

11.1. 服务端实现

# hello_world_agent.py
import asyncio
import logging
import grpc
import uvicorn
from fastapi import FastAPI

from a2a.server.agent_execution.agent_executor import AgentExecutor
from a2a.server.agent_execution.context import RequestContext
from a2a.server.events.event_queue import EventQueue
from a2a.server.request_handlers import DefaultRequestHandler, GrpcHandler
from a2a.server.routes import (
    add_a2a_routes_to_fastapi,
    create_agent_card_routes,
    create_jsonrpc_routes,
    create_rest_routes,
)
from a2a.server.tasks.inmemory_task_store import InMemoryTaskStore
from a2a.server.tasks.task_updater import TaskUpdater
from a2a.types import (
    AgentCapabilities, AgentCard, AgentInterface,
    AgentProvider, AgentSkill, Part, Task,
    TaskState, TaskStatus, a2a_pb2_grpc,
)


class HelloWorldAgentExecutor(AgentExecutor):
    """实现 AgentExecutor 接口,包含核心 Agent 逻辑"""

    asyncdef execute(self, context: RequestContext, event_queue: EventQueue) -> None:
        task_id = context.task_id
        context_id = context.context_id

        # Step 1: 发布初始 Task 事件(SUBMITTED 状态)
        await event_queue.enqueue_event(
            Task(
                id=task_id,
                context_id=context_id,
                status=TaskStatus(state=TaskState.TASK_STATE_SUBMITTED),
                history=[context.message],  # 将用户消息加入历史
            )
        )

        # Step 2: 使用 TaskUpdater 简化后续状态更新
        updater = TaskUpdater(event_queue, task_id, context_id)

        # 更新为 WORKING 状态
        await updater.start_work(
            message=updater.new_agent_message([Part(text='Processing your request...')])
        )

        # Step 3: 执行 Agent 逻辑
        user_input = context.get_user_input()
        reply = f"Hello! You said: '{user_input}'"

        # Step 4: 发布 Artifact(任务输出)
        await updater.add_artifact(
            parts=[Part(text=reply)],
            name='response',
            last_chunk=True,
        )

        # Step 5: 标记任务完成
        await updater.complete()

    asyncdef cancel(self, context: RequestContext, event_queue: EventQueue) -> None:
        updater = TaskUpdater(event_queue, context.task_id or'', context.context_id or'')
        await updater.cancel()


asyncdef serve(host: str = '127.0.0.1', port: int = 41241, grpc_port: int = 50051):
    # 定义 Agent Card
    agent_card = AgentCard(
        name='Hello World Agent',
        description='A simple hello world agent.',
        provider=AgentProvider(organization='Example', url='https://example.com'),
        version='1.0.0',
        capabilities=AgentCapabilities(streaming=True, push_notifications=False),
        default_input_modes=['text/plain'],
        default_output_modes=['text/plain'],
        skills=[
            AgentSkill(
                id='hello',
                name='Hello World',
                description='Greets the user.',
                tags=['greeting'],
                examples=['Hello!'],
            )
        ],
        supported_interfaces=[
            AgentInterface(protocol_binding='GRPC', protocol_version='1.0', url=f'{host}:{grpc_port}'),
            AgentInterface(protocol_binding='JSONRPC', protocol_version='1.0', url=f'http://{host}:{port}/a2a/jsonrpc'),
            AgentInterface(protocol_binding='HTTP+JSON', protocol_version='1.0', url=f'http://{host}:{port}/a2a/rest'),
        ],
    )

    # 组装服务端组件
    task_store = InMemoryTaskStore()
    request_handler = DefaultRequestHandler(
        agent_executor=HelloWorldAgentExecutor(),
        task_store=task_store,
        agent_card=agent_card,
    )

    # 创建 FastAPI 应用并挂载 A2A 路由
    app = FastAPI()
    add_a2a_routes_to_fastapi(
        app,
        agent_card_routes=create_agent_card_routes(agent_card=agent_card),
        jsonrpc_routes=create_jsonrpc_routes(request_handler=request_handler, rpc_url='/a2a/jsonrpc'),
        rest_routes=create_rest_routes(request_handler=request_handler, path_prefix='/a2a/rest'),
    )

    # 启动 gRPC 服务器
    grpc_server = grpc.aio.server()
    grpc_server.add_insecure_port(f'{host}:{grpc_port}')
    a2a_pb2_grpc.add_A2AServiceServicer_to_server(GrpcHandler(request_handler), grpc_server)

    # 并发启动 HTTP 和 gRPC 服务器
    await asyncio.gather(
        grpc_server.start(),
        uvicorn.Server(uvicorn.Config(app, host=host, port=port)).serve(),
    )


if __name__ == '__main__':
    logging.basicConfig(level=logging.INFO)
    asyncio.run(serve())

13.2. 客户端实现

# client_example.py
import asyncio
import uuid
import httpx
import grpc

from a2a.client import A2ACardResolver, ClientConfig, create_client
from a2a.types import Message, Part, Role, SendMessageRequest, TaskState


asyncdef main():
    agent_url = 'http://127.0.0.1:41241'

    # Step 1: 发现 Agent Card
    asyncwith httpx.AsyncClient() as httpx_client:
        resolver = A2ACardResolver(httpx_client, agent_url)
        card = await resolver.get_agent_card()
        print(f'Connected to: {card.name}')

    # Step 2: 创建客户端(自动选择传输层)
    config = ClientConfig(grpc_channel_factory=grpc.aio.insecure_channel)
    client = await create_client(card, client_config=config)

    # Step 3: 发送消息
    context_id = str(uuid.uuid4())
    current_task_id = None

    message = Message(
        role=Role.ROLE_USER,
        message_id=str(uuid.uuid4()),
        parts=[Part(text='Hello, World!')],
        context_id=context_id,
    )

    # Step 4: 处理流式响应
    asyncfor event in client.send_message(SendMessageRequest(message=message)):
        if event.HasField('task'):
            current_task_id = event.task.id
            print(f'Task created: {current_task_id}')

        elif event.HasField('status_update'):
            state = TaskState.Name(event.status_update.status.state)
            print(f'Status: {state}')
            if event.status_update.status.HasField('message'):
                text = event.status_update.status.message.parts[0].text
                print(f'  Message: {text}')

        elif event.HasField('artifact_update'):
            artifact = event.artifact_update.artifact
            print(f'Artifact [{artifact.name}]: {artifact.parts[0].text}')

    await client.close()


if __name__ == '__main__':
    asyncio.run(main())

后续如果有需要的话,可以把a2a官方项目的源码拿出来看一下,然后通过a2a做一个实际的从0到1的开源项目。

另外大家反馈的信息,我也在考虑建群中,请稍等几天,我把事情做好,直接邀请大家进来一起讨论关于agent的问题。

还有就是代码的问题,我也在整理中,很多代码都只是一种概念,没有实际的应用,我在想要结合一个开源项目来真正的实现之前聊到的agent思想。请大家稍作等待,干货很快就送上。

另外大家提到的,需要gpt-image2模型生成图片的焚诀(skill或者prompt):请你依照参考以下的内容(或者图片)生成一幅架构图(流程图/时序图等等),可以适当添加相关的xxx内容。【这就是自己整理的实际要画的图的整体流程】。

后续如果有需要,我会单独放到一个skills中,以便供大家使用。


内容效果不满意?点此反馈

输入关键词开始搜索