A2A协议与AI Agent间通信架构深度实战:从协议原理到多Agent协作系统与企业级Agent网络构建全解析

举报
江南清风起 发表于 2026/08/26 19:05:02 2026/08/26
【摘要】 A2A协议与AI Agent间通信架构深度实战:从协议原理到多Agent协作系统与企业级Agent网络构建全解析 引言Agent-to-Agent(A2A)协议是2026年AI Agent生态中最重要的通信标准之一。随着AI Agent从单机工具发展为分布式协作系统,Agent之间如何发现彼此、如何通信、如何协调任务、如何保证安全成为核心挑战。A2A协议由Google联合多家AI公司推出,...

A2A协议与AI Agent间通信架构深度实战:从协议原理到多Agent协作系统与企业级Agent网络构建全解析

引言

Agent-to-Agent(A2A)协议是2026年AI Agent生态中最重要的通信标准之一。随着AI Agent从单机工具发展为分布式协作系统,Agent之间如何发现彼此、如何通信、如何协调任务、如何保证安全成为核心挑战。A2A协议由Google联合多家AI公司推出,定义了Agent间的标准通信接口,使得不同框架(LangChain、CrewAI、AutoGen、Vertex AI等)构建的Agent能够互操作。本文将从A2A协议的核心概念出发,系统性地讲解Agent Card、任务管理、消息传递、能力发现、安全认证、多Agent编排等关键内容,通过大量可运行的Python代码示例帮助读者构建多Agent协作系统。

一、A2A协议核心概念

1.1 协议概述

A2A协议的设计目标是创建一个标准化的Agent间通信接口,类似于HTTP之于Web服务。在A2A架构中,每个Agent都是一个独立的服务,暴露标准的API端点。其他Agent(或人类用户)可以通过这些标准端点与Agent交互,无需了解Agent的内部实现。A2A协议基于HTTP/JSON-RPC构建,支持同步和异步通信模式。核心概念包括Agent Card(Agent名片,描述Agent的能力和端点)、Task(任务,Agent间协作的基本单位)、Message(消息,任务中的通信内容)、Artifact(产出物,任务的结果产物)。

1.2 Agent Card

Agent Card是A2A协议的核心——每个Agent通过一个JSON格式的Agent Card向外界声明自己的身份、能力和联系方式。以下是一个Agent Card的完整示例:

# Agent Card结构示例
agent_card = {
    "id": "code-reviewer-agent-001",
    "name": "Code Review Agent",
    "description": "专业的代码审查Agent,支持多种编程语言的代码质量分析、安全漏洞检测和改进建议生成。",
    "version": "2.1.0",
    "capabilities": {
        "streaming": True,           # 支持流式输出
        "pushNotifications": True,    # 支持推送通知
        "stateTransitionHistory": True # 支持状态历史
    },
    "skills": [
        {
            "id": "security-review",
            "name": "安全审查",
            "description": "检测SQL注入、XSS、路径遍历等安全漏洞",
            "inputModes": ["text", "code"],
            "outputModes": ["text", "markdown"],
            "examples": [
                "审查这段Python代码的安全性",
                "检查这个API端点是否有漏洞"
            ]
        },
        {
            "id": "performance-review",
            "name": "性能审查",
            "description": "分析代码性能瓶颈,提供优化建议",
            "inputModes": ["text", "code"],
            "outputModes": ["text", "markdown"],
            "examples": [
                "分析这个函数的性能",
                "找出N+1查询问题"
            ]
        },
        {
            "id": "style-review",
            "name": "代码风格审查",
            "description": "检查代码是否符合编码规范",
            "inputModes": ["text", "code"],
            "outputModes": ["text", "markdown"],
        }
    ],
    "authentication": {
        "schemes": ["bearer", "api-key"],
        "credentials": "required"
    },
    "endpoints": {
        "base_url": "https://code-reviewer.example.com/a2a",
        "tasks_send": "/tasks/send",
        "tasks_get": "/tasks/get",
        "tasks_cancel": "/tasks/cancel",
        "tasks_subscribe": "/tasks/subscribe",
        "agent_card": "/.well-known/agent.json"
    },
    "defaultInputModes": ["text", "code"],
    "defaultOutputModes": ["text", "markdown"],
}

1.3 任务生命周期

A2A协议定义了标准化的任务生命周期。任务从submitted状态开始,经过working状态(Agent正在处理),最终到达completed(成功完成)、failed(失败)或canceled(被取消)状态。以下是任务状态机的Python实现:

from enum import Enum
from typing import Optional, Dict, Any
from datetime import datetime
import uuid

class TaskState(Enum):
    SUBMITTED = "submitted"
    WORKING = "working"
    INPUT_REQUIRED = "input_required"  # 需要额外输入
    COMPLETED = "completed"
    FAILED = "failed"
    CANCELED = "canceled"

# 合法的状态转换
VALID_TRANSITIONS = {
    TaskState.SUBMITTED: [TaskState.WORKING, TaskState.CANCELED, TaskState.FAILED],
    TaskState.WORKING: [TaskState.INPUT_REQUIRED, TaskState.COMPLETED, TaskState.FAILED, TaskState.CANCELED],
    TaskState.INPUT_REQUIRED: [TaskState.WORKING, TaskState.CANCELED],
    TaskState.COMPLETED: [],
    TaskState.FAILED: [],
    TaskState.CANCELED: [],
}

class Task:
    """A2A任务。"""

    def __init__(self, sender: str, receiver: str, message: str, skill_id: Optional[str] = None):
        self.id = str(uuid.uuid4())
        self.sender = sender
        self.receiver = receiver
        self.message = message
        self.skill_id = skill_id
        self.state = TaskState.SUBMITTED
        self.history: list[dict] = []
        self.artifacts: list[dict] = []
        self.created_at = datetime.utcnow()
        self.updated_at = datetime.utcnow()
        self._add_history("Task created", TaskState.SUBMITTED)

    def _add_history(self, message: str, state: TaskState):
        self.history.append({
            "timestamp": datetime.utcnow().isoformat(),
            "message": message,
            "state": state.value,
            "actor": self.sender,
        })
        self.updated_at = datetime.utcnow()

    def transition(self, new_state: TaskState, message: str = ""):
        """执行状态转换。"""
        if new_state not in VALID_TRANSITIONS.get(self.state, []):
            raise ValueError(f"Invalid transition: {self.state.value} -> {new_state.value}")

        old_state = self.state
        self.state = new_state
        self._add_history(f"State: {old_state.value} -> {new_state.value}. {message}", new_state)

    def add_artifact(self, name: str, content: str, mime_type: str = "text/plain"):
        """添加产出物。"""
        artifact = {
            "id": str(uuid.uuid4()),
            "name": name,
            "content": content,
            "mimeType": mime_type,
            "createdAt": datetime.utcnow().isoformat(),
        }
        self.artifacts.append(artifact)

    def add_message(self, sender: str, content: str):
        """添加消息到任务。"""
        self.history.append({
            "timestamp": datetime.utcnow().isoformat(),
            "message": content,
            "state": self.state.value,
            "actor": sender,
        })

    def to_dict(self) -> dict:
        return {
            "id": self.id,
            "sender": self.sender,
            "receiver": self.receiver,
            "state": self.state.value,
            "message": self.message,
            "skillId": self.skill_id,
            "history": self.history,
            "artifacts": self.artifacts,
            "createdAt": self.created_at.isoformat(),
            "updatedAt": self.updated_at.isoformat(),
        }

二、A2A Server实现

2.1 基础Agent Server

以下是一个完整的A2A协议Agent Server实现,使用FastAPI构建:

from fastapi import FastAPI, HTTPException, Header
from pydantic import BaseModel
from typing import Optional, Dict, Any
import uuid
import asyncio
from datetime import datetime

app = FastAPI(title="A2A Code Review Agent")

# === 数据存储(生产环境应使用数据库) ===
tasks_db: Dict[str, Task] = {}
agent_card = {
    "id": "code-reviewer-agent",
    "name": "Code Review Agent",
    "description": "AI-powered code review agent",
    "version": "2.0.0",
    "capabilities": {
        "streaming": True,
        "pushNotifications": False,
        "stateTransitionHistory": True,
    },
    "skills": [
        {
            "id": "security-review",
            "name": "Security Review",
            "description": "Detect security vulnerabilities in code",
        },
        {
            "id": "performance-review",
            "name": "Performance Review",
            "description": "Analyze performance bottlenecks",
        },
    ],
    "endpoints": {
        "base_url": "http://localhost:8001/a2a",
    },
    "defaultInputModes": ["text", "code"],
    "defaultOutputModes": ["text", "markdown"],
}

# === 请求/响应模型 ===
class TaskSendRequest(BaseModel):
    message: str
    skillId: Optional[str] = None
    sender: str = "anonymous"
    inputModes: list[str] = ["text"]
    outputModes: list[str] = ["text"]
    sessionId: Optional[str] = None

class TaskResponse(BaseModel):
    id: str
    state: str
    message: Optional[str] = None
    artifacts: list = []

# === API端点 ===

@app.get("/.well-known/agent.json")
async def get_agent_card():
    """Agent Card发现端点(标准路径)。"""
    return agent_card

@app.post("/a2a/tasks/send")
async def send_task(req: TaskSendRequest):
    """发送任务给Agent。"""
    # 创建任务
    task = Task(
        sender=req.sender,
        receiver=agent_card["id"],
        message=req.message,
        skill_id=req.skillId,
    )
    tasks_db[task.id] = task

    # 异步处理任务
    asyncio.create_task(process_task(task))

    return task.to_dict()

@app.get("/a2a/tasks/{task_id}")
async def get_task(task_id: str):
    """获取任务状态。"""
    task = tasks_db.get(task_id)
    if not task:
        raise HTTPException(status_code=404, detail="Task not found")
    return task.to_dict()

@app.post("/a2a/tasks/{task_id}/cancel")
async def cancel_task(task_id: str):
    """取消任务。"""
    task = tasks_db.get(task_id)
    if not task:
        raise HTTPException(status_code=404, detail="Task not found")
    try:
        task.transition(TaskState.CANCELED, "Cancelled by sender")
    except ValueError as e:
        raise HTTPException(status_code=400, detail=str(e))
    return task.to_dict()

@app.get("/a2a/tasks/{task_id}/subscribe")
async def subscribe_task(task_id: str):
    """订阅任务状态更新(SSE)。"""
    from fastapi.responses import StreamingResponse

    async def event_stream():
        task = tasks_db.get(task_id)
        if not task:
            yield f"data: {json.dumps({'error': 'Task not found'})}\n\n"
            return

        last_state = None
        while True:
            if task.state != last_state:
                last_state = task.state
                yield f"data: {json.dumps(task.to_dict())}\n\n"

            if task.state in [TaskState.COMPLETED, TaskState.FAILED, TaskState.CANCELED]:
                break

            await asyncio.sleep(1)

    return StreamingResponse(event_stream(), media_type="text/event-stream")

# === 任务处理逻辑 ===
async def process_task(task: Task):
    """异步处理任务。"""
    try:
        task.transition(TaskState.WORKING, "Starting processing")

        # 模拟AI处理
        skill = task.skill_id or "general-review"
        review_result = await perform_code_review(task.message, skill)

        # 添加产出物
        task.add_artifact(
            name="review_report",
            content=review_result,
            mime_type="text/markdown",
        )

        task.transition(TaskState.COMPLETED, "Review completed")
    except Exception as e:
        task.transition(TaskState.FAILED, str(e))

async def perform_code_review(code: str, skill: str) -> str:
    """执行代码审查(这里使用模拟实现,实际应调用LLM)。"""
    await asyncio.sleep(2)  # 模拟处理时间

    if skill == "security-review":
        return f"""## 安全审查报告

### 分析的代码

{code[:500]}


### 发现
1. **SQL注入风险** - 检测到字符串拼接SQL查询
2. **XSS风险** - 未对用户输入进行转义
3. **路径遍历** - 文件路径未验证

### 修复建议
1. 使用参数化查询替代字符串拼接
2. 对所有用户输入进行HTML转义
3. 验证文件路径在允许的根目录内
"""
    elif skill == "performance-review":
        return f"""## 性能审查报告

### 分析的代码

{code[:500]}


### 发现
1. **N+1查询** - 循环中执行数据库查询
2. **内存泄漏** - 未释放的资源
3. **同步阻塞** - IO操作使用同步方式

### 优化建议
1. 使用批量查询或JOIN替代循环查询
2. 使用context manager管理资源
3.IO操作改为异步
"""
    else:
        return f"## 代码审查完成\n\n已分析代码,未发现严重问题。"

if __name__ == "__main__":
    import uvicorn
    uvicorn.run(app, host="0.0.0.0", port=8001)

三、A2A Client实现

3.1 Agent发现与通信

以下是一个A2A Client实现,能够发现Agent、发送任务和接收结果:

import httpx
import asyncio
import json
from typing import Optional, Dict, Any

class A2AClient:
    """A2A协议客户端,用于与远程Agent通信。"""

    def __init__(self, agent_url: str, auth_token: Optional[str] = None):
        self.agent_url = agent_url.rstrip("/")
        self.auth_token = auth_token
        self.client = httpx.AsyncClient(
            timeout=60.0,
            headers={"Authorization": f"Bearer {auth_token}"} if auth_token else {}
        )
        self.agent_card: Optional[dict] = None

    async def discover(self) -> dict:
        """发现Agent Card。"""
        response = await self.client.get(f"{self.agent_url}/.well-known/agent.json")
        response.raise_for_status()
        self.agent_card = response.json()
        return self.agent_card

    async def send_task(
        self,
        message: str,
        skill_id: Optional[str] = None,
        sender: str = "a2a-client",
    ) -> dict:
        """发送任务。"""
        response = await self.client.post(f"{self.agent_url}/a2a/tasks/send", json={
            "message": message,
            "skillId": skill_id,
            "sender": sender,
        })
        response.raise_for_status()
        return response.json()

    async def get_task(self, task_id: str) -> dict:
        """获取任务状态。"""
        response = await self.client.get(f"{self.agent_url}/a2a/tasks/{task_id}")
        response.raise_for_status()
        return response.json()

    async def cancel_task(self, task_id: str) -> dict:
        """取消任务。"""
        response = await self.client.post(f"{self.agent_url}/a2a/tasks/{task_id}/cancel")
        response.raise_for_status()
        return response.json()

    async def wait_for_completion(self, task_id: str, poll_interval: float = 1.0, timeout: float = 120.0) -> dict:
        """等待任务完成。"""
        elapsed = 0.0
        while elapsed < timeout:
            task = await self.get_task(task_id)
            if task["state"] in ["completed", "failed", "canceled"]:
                return task
            await asyncio.sleep(poll_interval)
            elapsed += poll_interval
        raise TimeoutError(f"Task {task_id} did not complete within {timeout}s")

    async def send_and_wait(
        self,
        message: str,
        skill_id: Optional[str] = None,
        sender: str = "a2a-client",
        timeout: float = 120.0,
    ) -> dict:
        """发送任务并等待完成。"""
        task = await self.send_task(message, skill_id, sender)
        result = await self.wait_for_completion(task["id"], timeout=timeout)
        return result

    async def close(self):
        await self.client.aclose()


# 使用示例
async def main():
    # 创建客户端
    client = A2AClient("http://localhost:8001")

    # 发现Agent
    card = await client.discover()
    print(f"发现Agent: {card['name']}")
    print(f"技能: {[s['name'] for s in card['skills']]}")

    # 发送任务并等待结果
    result = await client.send_and_wait(
        message="""
def get_user(user_id):
    query = f"SELECT * FROM users WHERE id = {user_id}"
    cursor.execute(query)
    return cursor.fetchone()
""",
        skill_id="security-review",
        sender="test-client",
    )

    print(f"\n任务状态: {result['state']}")
    if result.get("artifacts"):
        print(f"\n审查报告:\n{result['artifacts'][0]['content']}")

    await client.close()

asyncio.run(main())

四、多Agent协作编排

4.1 Agent网络

以下是一个多Agent协作编排框架,实现Agent间的任务委派和协作:

import asyncio
from typing import Dict, List, Optional
from dataclasses import dataclass, field

@dataclass
class AgentInfo:
    """Agent信息。"""
    id: str
    name: str
    url: str
    skills: List[str] = field(default_factory=list)
    client: Optional[A2AClient] = None

class AgentNetwork:
    """多Agent协作网络。"""

    def __init__(self):
        self.agents: Dict[str, AgentInfo] = {}
        self._lock = asyncio.Lock()

    async def register_agent(self, agent_url: str) -> AgentInfo:
        """注册并发现Agent。"""
        client = A2AClient(agent_url)
        card = await client.discover()

        info = AgentInfo(
            id=card["id"],
            name=card["name"],
            url=agent_url,
            skills=[s["id"] for s in card.get("skills", [])],
            client=client,
        )

        async with self._lock:
            self.agents[info.id] = info

        print(f"注册Agent: {info.name} (技能: {info.skills})")
        return info

    async def find_agent_for_skill(self, skill_id: str) -> Optional[AgentInfo]:
        """找到拥有指定技能的Agent。"""
        for agent in self.agents.values():
            if skill_id in agent.skills:
                return agent
        return None

    async def delegate_task(
        self,
        message: str,
        skill_id: str,
        sender: str = "orchestrator",
    ) -> dict:
        """委派任务给合适的Agent。"""
        agent = await self.find_agent_for_skill(skill_id)
        if not agent:
            raise ValueError(f"No agent found for skill: {skill_id}")

        print(f"委派任务给 {agent.name} (skill: {skill_id})")
        result = await agent.client.send_and_wait(message, skill_id, sender)
        return result

    async def broadcast_task(
        self,
        message: str,
        skill_ids: List[str],
        sender: str = "orchestrator",
    ) -> Dict[str, dict]:
        """向多个Agent广播任务。"""
        tasks = {}
        for skill_id in skill_ids:
            agent = await self.find_agent_for_skill(skill_id)
            if agent:
                tasks[skill_id] = agent.client.send_and_wait(message, skill_id, sender)

        results = {}
        completed = await asyncio.gather(*tasks.values(), return_exceptions=True)
        for skill_id, result in zip(tasks.keys(), completed):
            if isinstance(result, Exception):
                results[skill_id] = {"error": str(result)}
            else:
                results[skill_id] = result

        return results

    async def pipeline(
        self,
        initial_message: str,
        steps: List[Dict[str, str]],
    ) -> Dict[str, Any]:
        """执行多Agent流水线。

        Args:
            initial_message: 初始输入
            steps: 步骤列表,每个步骤包含skill_id和可选的transform
        """
        results = {}
        current_input = initial_message

        for i, step in enumerate(steps):
            skill_id = step["skill_id"]
            transform = step.get("transform")

            if transform:
                current_input = transform(current_input, results)

            print(f"\n步骤 {i+1}: 委派给 {skill_id}")
            result = await self.delegate_task(current_input, skill_id)

            # 提取产出物内容作为下一步输入
            if result.get("artifacts"):
                current_input = result["artifacts"][0]["content"]

            results[f"step_{i+1}_{skill_id}"] = result

        return results

    async def close_all(self):
        """关闭所有Agent连接。"""
        for agent in self.agents.values():
            if agent.client:
                await agent.client.aclose()


# === 使用示例 ===
async def main():
    network = AgentNetwork()

    # 注册多个Agent
    await network.register_agent("http://localhost:8001")  # 代码审查Agent
    await network.register_agent("http://localhost:8002")  # 文档生成Agent
    await network.register_agent("http://localhost:8003")  # 测试生成Agent

    # 1. 单任务委派
    result = await network.delegate_task(
        message="def login(user, pwd): cursor.execute(f'SELECT * FROM users WHERE name={user}')",
        skill_id="security-review",
    )
    print("安全审查结果:", result.get("artifacts", [{}])[0].get("content", ""))

    # 2. 并行广播
    results = await network.broadcast_task(
        message="async def fetch(url): response = await client.get(url); return response.json()",
        skill_ids=["security-review", "performance-review"],
    )
    for skill, result in results.items():
        print(f"\n{skill}:", result.get("artifacts", [{}])[0].get("content", ""))

    # 3. 流水线:审查 -> 生成文档 -> 生成测试
    pipeline_result = await network.pipeline(
        initial_message="""
def process_payment(amount, card_num):
    # Process payment logic
    import requests
    response = requests.post('https://api.payment.com/charge', json={
        'amount': amount,
        'card': card_num
    })
    return response.json()
""",
        steps=[
            {"skill_id": "security-review"},
            {
                "skill_id": "doc-generation",
                "transform": lambda code, prev: f"为以下代码生成技术文档:\n{code}",
            },
            {
                "skill_id": "test-generation",
                "transform": lambda doc, prev: f"为以下代码生成测试用例:\n{prev.get('step_1_security-review', {}).get('artifacts', [{}])[0].get('content', '')}",
            },
        ],
    )

    await network.close_all()

asyncio.run(main())

五、安全与认证

5.1 认证机制

A2A协议支持多种认证方案。以下是实现Bearer Token认证的示例:

from fastapi import FastAPI, HTTPException, Header, Depends
import jwt
from datetime import datetime, timedelta

app = FastAPI()

SECRET_KEY = "your-secret-key"
ALGORITHM = "HS256"

def create_token(agent_id: str) -> str:
    """创建JWT token。"""
    payload = {
        "sub": agent_id,
        "iat": datetime.utcnow(),
        "exp": datetime.utcnow() + timedelta(hours=24),
    }
    return jwt.encode(payload, SECRET_KEY, algorithm=ALGORITHM)

async def verify_token(authorization: str = Header(...)):
    """验证Bearer Token。"""
    if not authorization.startswith("Bearer "):
        raise HTTPException(status_code=401, detail="Invalid authorization header")
    token = authorization.split(" ")[1]
    try:
        payload = jwt.decode(token, SECRET_KEY, algorithms=[ALGORITHM])
        return payload
    except jwt.ExpiredSignatureError:
        raise HTTPException(status_code=401, detail="Token expired")
    except jwt.InvalidTokenError:
        raise HTTPException(status_code=401, detail="Invalid token")

@app.post("/a2a/tasks/send")
async def send_task(req: TaskSendRequest, token: dict = Depends(verify_token)):
    """需要认证的任务发送端点。"""
    req.sender = token["sub"]  # 使用token中的agent ID
    # ... 任务处理逻辑

5.2 能力授权

class CapabilityChecker:
    """Agent能力授权检查器。"""

    def __init__(self):
        # agent_id -> allowed_skills
        self.permissions: Dict[str, set] = {
            "orchestrator-agent": {"security-review", "performance-review", "doc-generation"},
            "monitor-agent": {"performance-review"},
            "test-agent": {"test-generation"},
        }

    def can_use_skill(self, agent_id: str, skill_id: str) -> bool:
        """检查Agent是否有权使用某技能。"""
        allowed = self.permissions.get(agent_id, set())
        return skill_id in allowed

    def grant_permission(self, agent_id: str, skill_id: str):
        """授予权限。"""
        if agent_id not in self.permissions:
            self.permissions[agent_id] = set()
        self.permissions[agent_id].add(skill_id)

    def revoke_permission(self, agent_id: str, skill_id: str):
        """撤销权限。"""
        if agent_id in self.permissions:
            self.permissions[agent_id].discard(skill_id)

六、A2A与MCP的协作

6.1 A2A + MCP架构

A2A和MCP是互补的协议。MCP用于Agent与工具/数据源的通信,A2A用于Agent与Agent之间的通信。一个完整的AI Agent系统可以同时使用两者——通过MCP获取工具能力,通过A2A与其他Agent协作。以下是一个同时支持A2A和MCP的Agent实现:

from fastapi import FastAPI
import asyncio

app = FastAPI(title="Hybrid Agent (A2A + MCP)")

class HybridAgent:
    """同时支持A2A通信和MCP工具调用的Agent。"""

    def __init__(self):
        self.a2a_clients: Dict[str, A2AClient] = {}  # 其他Agent的A2A客户端
        self.mcp_tools: Dict[str, callable] = {}     # MCP工具

    def register_mcp_tool(self, name: str, handler: callable):
        """注册MCP工具。"""
        self.mcp_tools[name] = handler

    async def register_peer_agent(self, agent_url: str):
        """注册对等Agent。"""
        client = A2AClient(agent_url)
        card = await client.discover()
        self.a2a_clients[card["id"]] = client
        return card

    async def process_with_tools_and_agents(
        self,
        task: str,
        use_tools: bool = True,
        delegate_to_agents: bool = True,
    ) -> dict:
        """综合处理:使用本地工具和其他Agent。"""
        results = {"task": task, "tool_results": {}, "agent_results": {}}

        # 1. 使用本地MCP工具
        if use_tools:
            for tool_name, handler in self.mcp_tools.items():
                try:
                    result = await handler(task)
                    results["tool_results"][tool_name] = result
                except Exception as e:
                    results["tool_results"][tool_name] = {"error": str(e)}

        # 2. 委派给其他Agent
        if delegate_to_agents:
            delegate_tasks = []
            for agent_id, client in self.a2a_clients.items():
                delegate_tasks.append(
                    client.send_and_wait(task, sender="hybrid-agent")
                )

            agent_results = await asyncio.gather(*delegate_tasks, return_exceptions=True)
            for agent_id, result in zip(self.a2a_clients.keys(), agent_results):
                if isinstance(result, Exception):
                    results["agent_results"][agent_id] = {"error": str(result)}
                else:
                    results["agent_results"][agent_id] = {
                        "state": result.get("state"),
                        "artifacts": result.get("artifacts", []),
                    }

        return results


# 使用示例
agent = HybridAgent()

# 注册MCP工具
agent.register_mcp_tool("code_search", lambda q: {"matches": ["file1.py", "file2.py"]})
agent.register_mcp_tool("db_query", lambda q: {"rows": [{"id": 1, "name": "test"}]})

# 注册对等Agent(需要在各自的端口运行A2A Server)
# await agent.register_peer_agent("http://localhost:8001")  # 代码审查Agent
# await agent.register_peer_agent("http://localhost:8002")  # 文档Agent

七、实战:构建软件开发Agent团队

7.1 完整的多Agent开发团队

以下是一个完整的软件开发Agent团队实现,包含项目经理、开发者、测试和审查四个Agent角色,通过A2A协议协作:

import asyncio
from typing import Dict

class DevTeamOrchestrator:
    """软件开发团队编排器。"""

    def __init__(self):
        self.network = AgentNetwork()
        self.team: Dict[str, str] = {}  # role -> agent_id

    async def setup(self):
        """注册团队成员Agent。"""
        # 注册各个Agent
        pm_card = await self.network.register_agent("http://localhost:8010")
        dev_card = await self.network.register_agent("http://localhost:8011")
        test_card = await self.network.register_agent("http://localhost:8012")
        review_card = await self.network.register_agent("http://localhost:8013")

        self.team = {
            "pm": pm_card.id,
            "dev": dev_card.id,
            "tester": test_card.id,
            "reviewer": review_card.id,
        }

    async def develop_feature(self, feature_request: str) -> Dict:
        """开发一个完整功能。"""
        results = {}

        # 1. PM分析需求
        print("=== 步骤1: 需求分析 ===")
        pm_result = await self.network.delegate_task(
            message=feature_request,
            skill_id="requirements-analysis",
        )
        requirements = pm_result.get("artifacts", [{}])[0].get("content", "")
        results["requirements"] = requirements

        # 2. 开发者实现
        print("\n=== 步骤2: 代码实现 ===")
        dev_result = await self.network.delegate_task(
            message=f"基于以下需求实现功能:\n{requirements}",
            skill_id="code-implementation",
        )
        code = dev_result.get("artifacts", [{}])[0].get("content", "")
        results["code"] = code

        # 3. 并行:测试 + 审查
        print("\n=== 步骤3: 测试和审查(并行) ===")
        parallel_results = await self.network.broadcast_task(
            message=f"审查/测试以下代码:\n{code}",
            skill_ids=["test-generation", "security-review", "performance-review"],
        )
        results["tests"] = parallel_results.get("test-generation", {})
        results["security_review"] = parallel_results.get("security-review", {})
        results["performance_review"] = parallel_results.get("performance-review", {})

        # 4. 如果审查有问题,循环修复
        issues_found = any(
            r.get("artifacts") and "问题" in r["artifacts"][0].get("content", "")
            for r in [results["security_review"], results["performance_review"]]
        )

        if issues_found:
            print("\n=== 步骤4: 修复问题 ===")
            fix_input = f"""
代码:
{code}

安全问题:
{results['security_review'].get('artifacts', [{}])[0].get('content', '')}

性能问题:
{results['performance_review'].get('artifacts', [{}])[0].get('content', '')}

请修复以上问题并输出修复后的代码。
"""
            fix_result = await self.network.delegate_task(
                message=fix_input,
                skill_id="code-fix",
            )
            results["fixed_code"] = fix_result.get("artifacts", [{}])[0].get("content", "")

        return results

    async def close(self):
        await self.network.close_all()


# 使用示例
async def main():
    orchestrator = DevTeamOrchestrator()
    await orchestrator.setup()

    result = await orchestrator.develop_feature(
        "实现一个用户注册API,支持邮箱验证、密码强度检查和防重复注册"
    )

    print("\n=== 最终结果 ===")
    for key, value in result.items():
        if isinstance(value, dict):
            content = value.get("artifacts", [{}])[0].get("content", "")
            print(f"\n[{key}]:\n{content[:200]}...")
        else:
            print(f"\n[{key}]:\n{str(value)[:200]}...")

    await orchestrator.close()

asyncio.run(main())

总结

A2A协议作为AI Agent间通信的标准化协议,正在成为多Agent协作系统的基础设施。本文系统性地覆盖了A2A协议的核心概念(Agent Card、任务生命周期、消息传递)、Agent Server和Client的实现、多Agent协作编排(委派、广播、流水线)、安全认证、A2A与MCP的协作以及完整的软件开发Agent团队实战。关键要点包括:Agent Card是Agent能力发现的标准机制;任务生命周期管理确保协作的可靠性;多Agent编排支持委派、广播和流水线等多种协作模式;A2A与MCP是互补协议,前者解决Agent间通信,后者解决Agent与工具通信;安全认证在生产部署中不可或缺。随着AI Agent从单体走向分布式协作,A2A协议将成为连接不同框架、不同供应商Agent的桥梁,掌握A2A开发能力将成为AI应用架构师的核心竞争力。

【声明】本内容来自华为云开发者社区博主,不代表华为云及华为云开发者社区的观点和立场。转载时必须标注文章的来源(华为云社区)、文章链接、文章作者等基本信息,否则作者和本社区有权追究责任。如果您发现本社区中有涉嫌抄袭的内容,欢迎发送邮件进行举报,并提供相关证据,一经查实,本社区将立刻删除涉嫌侵权内容,举报邮箱: cloudbbs@huaweicloud.com
  • 点赞
  • 收藏
  • 关注作者

评论(0

0/1000
抱歉,系统识别当前为高风险访问,暂不支持该操作

全部回复

上滑加载中

设置昵称

在此一键设置昵称,即可参与社区互动!

*长度不超过10个汉字或20个英文字符,设置后3个月内不可修改。

*长度不超过10个汉字或20个英文字符,设置后3个月内不可修改。