A2A (Agent-to-Agent Communication Protocol) 是 EventMesh 的一个高性能协议插件,专门设计用于支持智能体(Agents)之间的异步通信、协作和任务协调。该协议采用 MCP (Model Context Protocol) 架构理念,结合 EventMesh 的事件驱动特性,打造了一个异步、解耦、高性能的智能体协作总线。
A2A 协议不仅支持传统的 FIPA-ACL 风格语义,更全面拥抱现代大模型(LLM)生态,通过支持 JSON-RPC 2.0 标准,实现了对工具调用(Tool Use)、上下文共享(Context Sharing)等 LLM 核心场景的原生支持。
- 标准兼容: 完全支持 MCP (Model Context Protocol) 定义的
tools/call,resources/read等标准方法。 - 事件驱动: 将同步的 RPC 调用映射为异步的 Request/Response 事件流,充分利用 EventMesh 的高并发处理能力。
- 协议无关: 所有的 MCP 消息都被封装在 CloudEvents 标准信封中,可以在 HTTP, TCP, gRPC, Kafka 等任意 EventMesh 支持的传输层上运行。
- 双模支持:
- Modern Mode: 支持 MCP 标准的 JSON-RPC 2.0 消息,面向 LLM 应用。
- Legacy Mode: 兼容旧版 A2A 协议(基于
messageType和 FIPA 动词),保障存量业务平滑迁移。
- 自动识别: 协议适配器根据消息内容特征(如
jsonrpc字段)自动智能选择处理模式。
- 批量处理: 原生支持 JSON-RPC Batch 请求,EventMesh 会将其自动拆分为并行事件流,极大提升吞吐量。
- 智能路由: 支持从 MCP 请求参数(如
_agentId)中提取路由线索,自动注入 CloudEvents 扩展属性 (targetagent),实现零解包路由。
- 类型映射: 自动将 MCP 方法映射为 CloudEvent Type (例如
tools/call->org.apache.eventmesh.a2a.tools.call.req)。 - 上下文传递: 利用 CloudEvents Extension 传递
traceparent,实现跨 Agent 的全链路追踪。
┌─────────────────────────────────────────────────────────────┐
│ EventMesh A2A Protocol v2.0 │
│ (MCP over CloudEvents Architecture) │
├─────────────────────────────────────────────────────────────┤
│ ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │
│ │ MCP/JSON-RPC│ │ Legacy A2A │ │ Protocol │ │
│ │ Handler │ │ Handler │ │ Delegator │ │
│ └─────────────┘ └─────────────┘ └─────────────┘ │
├─────────────────────────────────────────────────────────────┤
│ ┌───────────────────────────────────────────────────────┐ │
│ │ Enhanced A2A Protocol Adaptor │ │
│ │ (Intelligent Parsing & CloudEvent Mapping) │ │
│ └───────────────────────────────────────────────────────┘ │
├─────────────────────────────────────────────────────────────┤
│ EventMesh Protocol Infrastructure │
│ ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │
│ │ CloudEvents │ │ HTTP │ │ gRPC │ │
│ │ Protocol │ │ Protocol │ │ Protocol │ │
│ └─────────────┘ └─────────────┘ └─────────────┘ │
└─────────────────────────────────────────────────────────────┘
protocol-a2a 模块 runtime 模块 examples 模块
┌─────────────────┐ ┌──────────────────┐ ┌──────────────┐
│ A2AMessageTransport(接口) │ A2AGatewayServer │ │ A2AGatewayDemo│
│ A2AClient (SDK) │<──HTTP──>│ (main, Netty HTTP)│<──HTTP──>│ (纯客户端) │
│ AgentCard/Topic │ │ InMemoryTransport │ └──────────────┘
└─────────────────┘ │ GatewayService │
│ TaskRegistry(TTL) │
└──────────────────┘
| 组件 | 模块 | 职责 |
|---|---|---|
A2AGatewayServer |
runtime | Netty HTTP 服务器入口,预注册 mock agent |
A2AGatewayHttpHandler |
runtime | HTTP 请求路由,支持 SSE 流式响应 |
A2AGatewayService |
runtime | 核心编排:任务提交、响应处理、SSE 推送 |
TaskRegistry |
runtime | 内存任务状态机 + TTL 自动清理(5 分钟) |
A2APublishSubscribeService |
runtime | AgentCard 注册、发现、心跳管理 |
InMemoryA2AMessageTransport |
runtime | 内存 pub/sub(可替换为 EventMesh broker) |
A2AClient |
protocol-a2a | Java SDK,提供类型化 API |
为了在事件驱动架构中支持 MCP 的 Request/Response 模型,A2A 协议定义了以下映射规则:
| MCP 概念 | CloudEvent 映射 | 说明 |
|---|---|---|
Request (tools/call) |
type: org.apache.eventmesh.a2a.tools.call.req mcptype: request |
这是一个请求事件 |
Response (result) |
type: org.apache.eventmesh.a2a.common.response mcptype: response |
这是一个响应事件 |
Correlation (id) |
extension: collaborationid / id |
用于将 Response 关联回 Request |
| Target | extension: targetagent |
路由目标 Agent ID |
{
"jsonrpc": "2.0",
"method": "tools/call",
"params": {
"name": "get_weather",
"arguments": {
"city": "Shanghai"
},
"_agentId": "weather-service" // 路由提示
},
"id": "req-123456"
}转换后的 CloudEvent:
id:req-123456type:org.apache.eventmesh.a2a.tools.call.reqsource:eventmesh-a2aextension: a2amethod:tools/callextension: mcptype:requestextension: targetagent:weather-service
{
"jsonrpc": "2.0",
"result": {
"content": [
{
"type": "text",
"text": "Shanghai: 25°C, Sunny"
}
]
},
"id": "req-123456"
}转换后的 CloudEvent:
id:uuid-new-event-idtype:org.apache.eventmesh.a2a.common.responseextension: collaborationid:req-123456(关联 ID)extension: mcptype:response
{
"protocol": "A2A",
"messageType": "PROPOSE",
"sourceAgent": { "agentId": "agent-001" },
"payload": { "task": "data-process" }
}您只需要发送标准的 JSON-RPC 格式消息到 EventMesh:
// 1. 构造 MCP Request JSON
String mcpRequest = "{"
"jsonrpc": "2.0",
"method": "tools/call",
"params": { "name": "weather", "_agentId": "weather-agent" },
"id": "req-001"
"}";
// 2. 通过 EventMesh SDK 发送
eventMeshProducer.publish(new A2AProtocolTransportObject(mcpRequest));订阅相应的主题,处理业务逻辑,并发送回响应:
// 1. 订阅 MCP Request 主题
eventMeshConsumer.subscribe("org.apache.eventmesh.a2a.tools.call.req");
// 2. 收到消息后处理...
public void handle(CloudEvent event) {
// 解包 Request
String reqJson = new String(event.getData().toBytes());
// ... 执行业务逻辑 ...
// 3. 构造 Response
String mcpResponse = "{"
"jsonrpc": "2.0",
"result": { "text": "Sunny" },
"id": """ + event.getId() + """
"}";
// 4. 发送回 EventMesh
eventMeshProducer.publish(new A2AProtocolTransportObject(mcpResponse));
}A2A Gateway 提供完整的 REST API,支持非 Java 客户端通过 HTTP 交互:
# 同步提交 task
curl -X POST 'http://localhost:10105/a2a/tasks?mode=sync' \
-H 'Content-Type: application/json' \
-d '{"targetAgent":"weather-agent","message":"Beijing"}'
# 异步提交 task
curl -X POST 'http://localhost:10105/a2a/tasks?mode=async' \
-H 'Content-Type: application/json' \
-d '{"targetAgent":"weather-agent","message":"Shanghai"}'
# 查询状态
curl http://localhost:10105/a2a/tasks/{taskId}
# 列出 tasks(支持 state/limit/offset)
curl 'http://localhost:10105/a2a/tasks?state=COMPLETED&limit=20&offset=0'
# SSE 流式推送(含 heartbeat 保活)
curl -N http://localhost:10105/a2a/tasks/{taskId}/stream
# 健康检查
curl http://localhost:10105/a2a/health
# 列出 agents
curl http://localhost:10105/a2a/agents| 方法 | 路径 | 说明 |
|---|---|---|
| POST | /a2a/tasks?mode=sync |
同步提交 task(等待结果) |
| POST | /a2a/tasks?mode=async |
异步提交 task(立即返回 taskId) |
| GET | /a2a/tasks?state=&limit=&offset= |
分页列出 tasks,可按状态过滤 |
| GET | /a2a/tasks/{taskId} |
查询 task 状态 |
| DELETE | /a2a/tasks/{taskId} |
取消 task |
| GET | /a2a/tasks/{taskId}/wait |
长轮询等待 task 结果 |
| GET | /a2a/tasks/{taskId}/stream |
SSE 流式推送 task 状态更新 |
| GET | /a2a/agents |
列出所有已注册 agents |
| POST | /a2a/heartbeat |
Agent 心跳 |
| GET | /a2a/cards/list |
列出所有 AgentCard |
| POST | /a2a/cards/card/{org}/{unit}/{agent} |
注册 AgentCard |
A2AClient client = A2AClient.builder()
.gatewayUrl("http://localhost:10105")
.namespace("global")
.agentName("my-agent")
.agentCard(card)
.heartbeatInterval(30_000)
.build();
client.start();
// 同步 task(返回类型化 TaskResult)
TaskResult result = client.sendTaskSync("weather-agent", "Beijing", null);
// 异步 task(返回 taskId)
String taskId = client.sendTaskAsync("weather-agent", "Shanghai", null);
// 查询状态
TaskResult status = client.getTaskStatus(taskId);
// 取消
boolean cancelled = client.cancelTask(taskId);
// 列出 agents(返回 List<String>)
List<String> agents = client.listAgents();
client.shutdown();A2A 协议不限制 method 的名称。您可以定义自己的业务方法,例如 agents/negotiate 或 tasks/submit。EventMesh 会自动将其映射为 CloudEvent 类型 org.apache.eventmesh.a2a.agents.negotiate.req。
由于 A2A 兼容标准的 JSON-RPC 2.0,您可以轻松编写适配器,将 LangChain 的 Tool 调用转换为 EventMesh 消息,从而让您的 LLM 应用具备分布式、异步的通信能力。
-
v2.0.0: 全面拥抱 MCP (Model Context Protocol)
- 引入
EnhancedA2AProtocolAdaptor,支持 JSON-RPC 2.0。 - 实现异步 RPC over CloudEvents 模式。
- 支持 Request/Response 自动识别与语义映射。
- 保留对 Legacy A2A 协议的完全兼容。
- 引入
-
v2.1.0: Gateway 运行时架构
- 新增
A2AGatewayServer(Netty HTTP) 独立 Gateway 服务。 - 实现
TaskRegistry任务状态机 + TTL 自动清理(5 分钟)。 - 支持 SSE 流式响应 (
GET /a2a/tasks/{taskId}/stream)。 A2AClientSDK 返回类型化对象 (TaskResult,List<String>)。- 修复
pendingTasks竞态条件(put-before-publish)。 - AgentCard 注册、发现、心跳管理。
- 73 个测试场景全部通过。
- 新增
欢迎贡献代码和文档!请参考以下步骤:
- Fork项目仓库
- 创建功能分支
- 提交代码更改
- 创建Pull Request
Apache License 2.0