从零构建智能体协作网络:基于A2A协议的Python实战指南

想象一下,你正在开发一个复杂的应用,它需要处理自然语言查询、进行数据分析、生成图表,最后还要将结果通过邮件发送出去。传统的单体应用架构会让你陷入无尽的代码耦合和功能堆叠中。而今天,我们有机会用一种全新的范式来构建系统:智能体协作网络。这不再是让一个庞大的、无所不能的模型去处理所有事情,而是将任务分解,交给一群各有所长的“专家”智能体,让它们像一支训练有素的团队一样协同工作。

这种范式背后的核心协议之一,便是由Google推动的 A2A协议。它定义了一套标准,让不同的智能体能够互相发现、对话、委托任务并共享状态。这不仅仅是技术上的革新,更是一种架构思想的转变——从构建“全能巨人”转向构建“高效团队”。

在本指南中,我们将彻底抛开理论空谈,直接动手。我会带你从零开始,用Python搭建一个由多个智能体组成的协作网络。我们将创建三个各司其职的数学专家智能体,并构建一个能够理解用户意图、智能分配任务的“调度中心”。整个过程将包含完整的代码、清晰的解释以及我实际搭建时踩过的坑和解决方案。无论你是想探索下一代AI应用架构,还是希望提升现有系统的模块化与灵活性,这篇文章都将为你提供一条清晰的实践路径。

1. 理解核心:A2A协议与智能体网络架构

在敲下第一行代码之前,我们必须先厘清几个核心概念。A2A,即Agent-to-Agent协议,其目标是为异构的AI智能体提供一套通用的“社交语言”。你可以把它想象成智能体世界的HTTP协议,规定了它们如何自我介绍、如何请求服务、如何传递结果。

1.1 A2A协议解决了什么问题?

在单体AI应用中,所有逻辑都塞进一个庞大的提示词或一个复杂的函数链中。这种方式存在几个明显痛点:

  • 维护困难:任何功能改动都可能牵一发而动全身。
  • 能力瓶颈:单个模型可能不擅长所有类型的任务。
  • 资源浪费:简单的任务也需要调用庞大的模型,成本高昂。
  • 缺乏弹性:某个功能失败可能导致整个流程中断。

A2A协议通过引入“智能体”这一抽象层,将上述问题逐一化解。每个智能体都是一个独立的、专注的服务单元。它们通过A2A协议进行通信,共同完成复杂任务。这种架构带来了几个关键优势:

优势说明
模块化与解耦每个智能体独立开发、部署、更新,系统边界清晰。
能力复用优秀的智能体可以被多个不同的应用或工作流调用。
弹性与容错单个智能体故障不会导致整个系统瘫痪,任务可以被路由到备用智能体。
动态组合可以根据任务需求,动态地组合不同的智能体形成新的解决方案。

1.2 A2A与相关概念辨析

你可能会听到其他类似的概念,比如MCP。这里简单澄清一下它们的区别与联系:

  • MCP:全称是Model Context Protocol。它主要解决的是单个智能体如何与外部工具和资源进行交互的问题。比如,让智能体能够读取数据库、调用API、操作文件系统。你可以把MCP看作是智能体内部的“工具箱管理协议”。
  • A2A:全称是Agent-to-Agent Protocol。它关注的是不同智能体之间如何协作与任务委派。它定义了智能体之间“对话”和“分工”的规则。A2A更像是智能体之间的“社交网络协议”或“团队协作协议”。

它们的关系是互补的。一个A2A智能体在接收到任务后,可以内部使用MCP协议去调用各种工具来完成自己负责的那部分工作,然后再通过A2A协议将结果返回给发起者或委托给其他智能体。简单理解:A2A是对外协作的标准,MCP是对内使用工具的标准

1.3 我们的实战目标:构建一个数学助手网络

为了将概念具象化,我们本次实战的目标是构建一个“数学助手网络”。这个网络由三个专家智能体和一个调度中心组成:

  1. 正弦智能体:专精于计算任意数字的正弦值。
  2. 余弦智能体:专精于计算任意数字的余弦值。
  3. 正切智能体:专精于计算任意数字的正切值。
  4. 调度路由器:接收用户的自然语言查询(如“计算0.5的正弦值”),理解其意图,并将任务精准地路由到对应的专家智能体。

通过这个简单的例子,你将完整掌握A2A网络从搭建到运作的全过程。接下来,我们进入实战环节。

2. 环境准备与基础库安装

工欲善其事,必先利其器。我们的搭建工作将从配置一个干净的Python环境开始。

2.1 创建并激活虚拟环境

我强烈建议使用虚拟环境来管理项目依赖,这能避免不同项目间的库版本冲突。

# 创建项目目录并进入
mkdir a2a_math_network && cd a2a_math_network

# 创建Python虚拟环境(假设你使用Python 3.10+)
python -m venv venv

# 激活虚拟环境
# 在 macOS/Linux 上:
source venv/bin/activate
# 在 Windows 上:
# venv\Scripts\activate

激活后,你的命令行提示符前应该会出现 (venv) 字样,表示你已处于虚拟环境中。

2.2 安装核心库:python-a2a

Google官方并未提供一个叫 python_a2a 的PyPI包。在真实的A2A生态中,你可能需要使用像 agentalanggraph 这样的框架,或者直接基于gRPC/HTTP实现。为了本次教学清晰易懂,我们将创建一个简化的、模拟A2A协议的核心库

我们将手动创建这个库的核心模块。在项目根目录下,创建一个名为 a2a_core.py 的文件,并填入以下代码:

# a2a_core.py
# 一个简化的A2A协议核心实现,用于教学演示

from dataclasses import dataclass, field
from enum import Enum
from typing import Any, Dict, List, Optional, Callable
import asyncio
from contextlib import asynccontextmanager
import json
import aiohttp
from aiohttp import web
import logging

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

class TaskState(str, Enum):
    PENDING = "PENDING"
    RUNNING = "RUNNING"
    COMPLETED = "COMPLETED"
    FAILED = "FAILED"

@dataclass
class TaskStatus:
    state: TaskState
    error_message: Optional[str] = None
    progress: Optional[float] = None  # 0.0 到 1.0

@dataclass
class Artifact:
    type: str  # e.g., "text", "image", "data"
    content: Dict[str, Any]

@dataclass
class Task:
    task_id: str
    message: Dict[str, Any]  # 原始请求消息
    status: TaskStatus = field(default_factory=lambda: TaskStatus(state=TaskState.PENDING))
    artifacts: List[Artifact] = field(default_factory=list)  # 任务产出物

class AgentCard:
    """智能体的“名片”,用于声明自身能力"""
    def __init__(self, name: str, description: str, version: str, skills: List[Dict]):
        self.name = name
        self.description = description
        self.version = version
        self.skills = skills  # 技能列表,每个技能包含名称、描述、输入输出格式等
        self.endpoint: Optional[str] = None  # 智能体的服务端点

class A2AAgent:
    """A2A智能体基类"""
    def __init__(self, card: AgentCard):
        self.card = card
        self.card.endpoint = f"http://localhost:{self.port}" if hasattr(self, 'port') else None

    async def handle_task(self, task: Task) -> Task:
        """处理任务的核心方法,子类必须重写"""
        raise NotImplementedError("Subclasses must implement handle_task")

    async def invoke(self, target_agent_endpoint: str, task_message: Dict) -> Dict:
        """调用另一个智能体"""
        async with aiohttp.ClientSession() as session:
            async with session.post(f"{target_agent_endpoint}/task", json=task_message) as resp:
                return await resp.json()

class AgentNetwork:
    """智能体网络,负责注册和发现智能体"""
    def __init__(self):
        self.agents: Dict[str, AgentCard] = {}  # 智能体名称 -> 智能体名片

    def register(self, agent_name: str, agent_card: AgentCard):
        """注册一个智能体到网络"""
        self.agents[agent_name] = agent_card
        logger.info(f"Agent '{agent_name}' registered at {agent_card.endpoint}")

    def discover(self, skill_query: str) -> List[AgentCard]:
        """根据技能查询发现智能体(简化版:根据名称和描述匹配)"""
        matched = []
        for card in self.agents.values():
            if skill_query.lower() in card.name.lower() or skill_query.lower() in card.description.lower():
                matched.append(card)
            else:
                for skill in card.skills:
                    if skill_query.lower() in skill.get('name', '').lower():
                        matched.append(card)
                        break
        return matched

# 简化的HTTP服务器包装器
def run_agent_server(agent_instance: A2AAgent, port: int):
    """运行一个智能体的HTTP服务器"""
    app = web.Application()

    async def handle_task(request):
        data = await request.json()
        task = Task(task_id=data.get('task_id', 'unknown'), message=data)
        # 在实际中,这里应该异步处理,但为简化我们同步调用
        processed_task = asyncio.run(agent_instance.handle_task(task))
        return web.json_response({
            'task_id': processed_task.task_id,
            'status': {
                'state': processed_task.status.state.value,
                'error_message': processed_task.status.error_message
            },
            'artifacts': [{'type': a.type, 'content': a.content} for a in processed_task.artifacts]
        })

    app.router.add_post('/task', handle_task)

    # 添加一个端点用于获取智能体名片
    async def handle_discover(request):
        return web.json_response({
            'name': agent_instance.card.name,
            'description': agent_instance.card.description,
            'version': agent_instance.card.version,
            'skills': agent_instance.card.skills,
            'endpoint': agent_instance.card.endpoint
        })
    app.router.add_get('/discover', handle_discover)

    logger.info(f"Starting {agent_instance.card.name} on http://localhost:{port}")
    web.run_app(app, port=port)

这个文件定义了我们简化版A2A协议的核心数据结构和基类。它包含了:

  • TaskTaskStatus: 表示任务及其状态。
  • AgentCard: 智能体的“名片”,用于声明自身能力。
  • A2AAgent: 所有智能体的基类,定义了处理任务的接口。
  • AgentNetwork: 一个中心化的智能体注册与发现服务。
  • run_agent_server: 一个快速启动智能体HTTP服务的工具函数。

注意:这是一个高度简化的教学实现。真实的A2A协议可能涉及更复杂的序列化、流式响应、安全认证和服务发现机制。但它的核心思想——通过标准化的消息格式进行通信——已经体现出来了。

现在,我们的“基础框架”就准备好了。接下来,让我们创建第一个专家智能体。

3. 创建专家智能体:数学三剑客

我们将创建三个独立的Python文件,分别代表正弦、余弦、正切智能体。它们将继承自我们刚定义的 A2AAgent 基类。

3.1 正弦智能体

在项目根目录创建 sine_agent.py

# sine_agent.py
import math
import re
from a2a_core import A2AAgent, AgentCard, Task, TaskStatus, TaskState, Artifact, run_agent_server

class SineAgent(A2AAgent):
    """专精于计算正弦值的智能体"""
    def __init__(self, port: int):
        # 定义智能体的名片
        card = AgentCard(
            name="SineCalculator",
            description="Calculates the sine of a given number in radians.",
            version="1.0.0",
            skills=[
                {
                    "name": "compute_sine",
                    "description": "Compute the sine value of a floating-point number.",
                    "input_schema": {"type": "object", "properties": {"number": {"type": "number"}}},
                    "output_schema": {"type": "object", "properties": {"result": {"type": "number"}, "message": {"type": "string"}}}
                }
            ]
        )
        super().__init__(card)
        self.port = port
        self.card.endpoint = f"http://localhost:{port}"

    async def handle_task(self, task: Task) -> Task:
        logger.info(f"SineAgent received task: {task.message}")
        # 从消息中提取用户查询文本
        query_text = task.message.get("content", {}).get("text", "")
        
        # 使用正则表达式提取数字
        match = re.search(r"[-+]?\d*\.?\d+", query_text)
        if not match:
            task.status = TaskStatus(state=TaskState.FAILED, error_message="No valid number found in query.")
            return task
        
        try:
            number = float(match.group())
            result = math.sin(number)
            
            # 将结果封装为产出物
            task.artifacts.append(Artifact(
                type="text",
                content={
                    "result": result,
                    "message": f"The sine of {number} is {result:.6f}",
                    "input": number
                }
            ))
            task.status = TaskStatus(state=TaskState.COMPLETED)
            logger.info(f"SineAgent computed sin({number}) = {result}")
            
        except ValueError as e:
            task.status = TaskStatus(state=TaskState.FAILED, error_message=f"Invalid number format: {e}")
        except Exception as e:
            task.status = TaskState.FAILED, error_message=f"Unexpected error: {e}"
        
        return task

if __name__ == "__main__":
    # 启动智能体服务器,监听在指定端口
    agent = SineAgent(port=8001)
    run_agent_server(agent, port=8001)

代码解读

  1. __init__ 方法中,我们定义了智能体的名片,声明了它的名字、描述、版本和唯一技能 compute_sine
  2. handle_task 方法是核心。它接收一个 Task 对象,从中解析出用户查询。
  3. 我们使用一个简单的正则表达式从文本中提取数字。在实际生产环境中,你应该使用更鲁棒的解析方法,或者依赖上游的路由器来传递结构化的参数。
  4. 计算正弦值后,将结果包装成一个 Artifact 对象,附加到任务上,并将任务状态标记为 COMPLETED
  5. 最后,在 __main__ 中,我们实例化智能体并启动一个HTTP服务器。

用同样的模式,我们创建余弦和正切智能体。

3.2 余弦智能体

创建 cosine_agent.py

# cosine_agent.py
import math
import re
from a2a_core import A2AAgent, AgentCard, Task, TaskStatus, TaskState, Artifact, run_agent_server

class CosineAgent(A2AAgent):
    def __init__(self, port: int):
        card = AgentCard(
            name="CosineCalculator",
            description="Calculates the cosine of a given number in radians.",
            version="1.0.0",
            skills=[{
                "name": "compute_cosine",
                "description": "Compute the cosine value of a floating-point number.",
                "input_schema": {"type": "object", "properties": {"number": {"type": "number"}}},
                "output_schema": {"type": "object", "properties": {"result": {"type": "number"}, "message": {"type": "string"}}}
            }]
        )
        super().__init__(card)
        self.port = port
        self.card.endpoint = f"http://localhost:{port}"

    async def handle_task(self, task: Task) -> Task:
        logger.info(f"CosineAgent received task: {task.message}")
        query_text = task.message.get("content", {}).get("text", "")
        match = re.search(r"[-+]?\d*\.?\d+", query_text)
        if not match:
            task.status = TaskStatus(state=TaskState.FAILED, error_message="No valid number found in query.")
            return task
        try:
            number = float(match.group())
            result = math.cos(number)
            task.artifacts.append(Artifact(
                type="text",
                content={
                    "result": result,
                    "message": f"The cosine of {number} is {result:.6f}",
                    "input": number
                }
            ))
            task.status = TaskStatus(state=TaskState.COMPLETED)
            logger.info(f"CosineAgent computed cos({number}) = {result}")
        except Exception as e:
            task.status = TaskStatus(state=TaskState.FAILED, error_message=f"Error: {e}")
        return task

if __name__ == "__main__":
    agent = CosineAgent(port=8002)
    run_agent_server(agent, port=8002)

3.3 正切智能体

创建 tangent_agent.py

# tangent_agent.py
import math
import re
from a2a_core import A2AAgent, AgentCard, Task, TaskStatus, TaskState, Artifact, run_agent_server

class TangentAgent(A2AAgent):
    def __init__(self, port: int):
        card = AgentCard(
            name="TangentCalculator",
            description="Calculates the tangent of a given number in radians. Note: Input close to (n+0.5)*π may cause large values.",
            version="1.0.0",
            skills=[{
                "name": "compute_tangent",
                "description": "Compute the tangent value of a floating-point number.",
                "input_schema": {"type": "object", "properties": {"number": {"type": "number"}}},
                "output_schema": {"type": "object", "properties": {"result": {"type": "number"}, "message": {"type": "string"}}}
            }]
        )
        super().__init__(card)
        self.port = port
        self.card.endpoint = f"http://localhost:{port}"

    async def handle_task(self, task: Task) -> Task:
        logger.info(f"TangentAgent received task: {task.message}")
        query_text = task.message.get("content", {}).get("text", "")
        match = re.search(r"[-+]?\d*\.?\d+", query_text)
        if not match:
            task.status = TaskStatus(state=TaskState.FAILED, error_message="No valid number found in query.")
            return task
        try:
            number = float(match.group())
            result = math.tan(number)
            task.artifacts.append(Artifact(
                type="text",
                content={
                    "result": result,
                    "message": f"The tangent of {number} is {result:.6f}",
                    "input": number
                }
            ))
            task.status = TaskStatus(state=TaskState.COMPLETED)
            logger.info(f"TangentAgent computed tan({number}) = {result}")
        except Exception as e:
            task.status = TaskStatus(state=TaskState.FAILED, error_message=f"Error: {e}")
        return task

if __name__ == "__main__":
    agent = TangentAgent(port=8003)
    run_agent_server(agent, port=8003)

现在,我们已经有了三个专家智能体。每个智能体都是一个独立的HTTP服务,运行在不同的端口上(8001, 8002, 8003)。它们对外暴露了两个端点:

  • POST /task: 用于接收并处理任务。
  • GET /discover: 用于让其他服务发现自己的能力(返回AgentCard)。

接下来,我们需要让它们“组团”,并创建一个大脑来指挥它们。

4. 构建智能路由器与协作网络

智能体们已经就位,但还缺一个“指挥官”。这个指挥官需要做两件事:

  1. 服务发现:知道网络里有哪些智能体,各自有什么能力。
  2. 任务路由:根据用户请求的内容,决定将任务派发给哪个智能体。

我们将创建一个 orchestrator.py 文件来实现这个指挥官。

4.1 实现基于关键词的简单路由器

首先,我们实现一个最简单的路由器:基于查询语句中的关键词(如“sin”, “cos”, “tan”)来决定派发目标。

# orchestrator.py 第一部分:网络注册与简单路由
import asyncio
import aiohttp
import json
from typing import Dict, List, Tuple, Optional
from a2a_core import AgentNetwork, AgentCard

class SimpleOrchestrator:
    def __init__(self):
        self.network = AgentNetwork()
        self.agent_clients = {}  # 缓存智能体的HTTP客户端会话

    async def register_agent(self, name: str, endpoint: str):
        """向网络注册一个智能体(在实际中,可能通过服务发现自动完成)"""
        # 这里我们模拟获取智能体的名片
        async with aiohttp.ClientSession() as session:
            try:
                async with session.get(f"{endpoint}/discover", timeout=2) as resp:
                    if resp.status == 200:
                        card_data = await resp.json()
                        card = AgentCard(
                            name=card_data['name'],
                            description=card_data['description'],
                            version=card_data['version'],
                            skills=card_data['skills']
                        )
                        card.endpoint = endpoint
                        self.network.register(name, card)
                        self.agent_clients[name] = endpoint
                        print(f"[Orchestrator] Successfully registered agent: {name} at {endpoint}")
                    else:
                        print(f"[Orchestrator] Failed to discover agent at {endpoint}")
            except Exception as e:
                print(f"[Orchestrator] Error registering agent {name}: {e}")

    def _route_by_keyword(self, query: str) -> Tuple[Optional[str], float]:
        """基于关键词的简单路由逻辑"""
        query_lower = query.lower()
        routing_map = [
            (["sin", "sine"], "SineCalculator"),
            (["cos", "cosine"], "CosineCalculator"),
            (["tan", "tangent"], "TangentCalculator"),
        ]
        for keywords, agent_name in routing_map:
            for kw in keywords:
                if kw in query_lower:
                    return agent_name, 0.9  # 返回智能体名称和置信度
        return None, 0.0

    async def process_query(self, user_query: str) -> Dict:
        """处理用户查询的核心流程"""
        # 1. 路由决策
        target_agent_name, confidence = self._route_by_keyword(user_query)
        
        if not target_agent_name:
            return {
                "success": False,
                "error": f"No suitable agent found for query: '{user_query}'",
                "suggested_agents": [card.name for card in self.network.discover("calculate")]
            }
        
        print(f"[Orchestrator] Routing query '{user_query}' to '{target_agent_name}' (confidence: {confidence:.2f})")
        
        # 2. 获取目标智能体端点
        target_endpoint = self.agent_clients.get(target_agent_name)
        if not target_endpoint:
            return {"success": False, "error": f"Agent '{target_agent_name}' is not available."}
        
        # 3. 构造任务消息并发送
        task_message = {
            "task_id": f"task_{asyncio.current_task().get_name()}_{id(user_query)}",
            "content": {
                "text": user_query,
                "type": "math_calculation_request"
            }
        }
        
        # 4. 调用目标智能体
        async with aiohttp.ClientSession() as session:
            try:
                async with session.post(f"{target_endpoint}/task", json=task_message, timeout=10) as resp:
                    if resp.status == 200:
                        result = await resp.json()
                        return {
                            "success": True,
                            "agent": target_agent_name,
                            "result": result
                        }
                    else:
                        return {
                            "success": False,
                            "error": f"Agent returned error status: {resp.status}",
                            "agent": target_agent_name
                        }
            except asyncio.TimeoutError:
                return {"success": False, "error": f"Request to {target_agent_name} timed out."}
            except Exception as e:
                return {"success": False, "error": f"Communication error: {e}"}

这个 SimpleOrchestrator 类做了以下几件事:

  1. 维护一个 AgentNetwork 来记录所有注册的智能体。
  2. 提供了 register_agent 方法,通过调用智能体的 /discover 端点来获取其名片并注册。
  3. 实现了一个非常基础的 _route_by_keyword 路由逻辑。
  4. process_query 方法封装了完整的处理流程:路由 -> 调用 -> 返回结果。

4.2 启动网络并测试

现在,让我们创建一个主脚本来启动整个系统并进行测试。创建 main.py

# main.py
import asyncio
import sys
import subprocess
import time
from orchestrator import SimpleOrchestrator

async def main():
    # 定义我们的智能体服务(在实际部署中,这些可能是独立的容器或进程)
    agents = {
        "SineCalculator": "http://localhost:8001",
        "CosineCalculator": "http://localhost:8002",
        "TangentCalculator": "http://localhost:8003",
    }
    
    # 启动智能体服务(在实际中,你可能用 systemd, docker-compose 等管理)
    # 这里为了演示,我们假设智能体服务已经在运行。
    # 你可以通过分别运行 `python sine_agent.py`, `python cosine_agent.py`, `python tangent_agent.py` 来启动它们。
    
    print("[Main] Assuming agent servers are already running on ports 8001, 8002, 8003.")
    print("[Main] If not, please start them in separate terminal windows.")
    time.sleep(2)  # 给一点时间让用户确认服务已启动
    
    # 初始化指挥器
    orchestrator = SimpleOrchestrator()
    
    # 注册所有智能体
    print("\n[Main] Registering agents with the orchestrator...")
    for name, endpoint in agents.items():
        await orchestrator.register_agent(name, endpoint)
    
    # 测试查询
    test_queries = [
        "What is the sine of 0.7854?",
        "Calculate cosine for 1.0472",
        "Compute the tangent of 0.5236",
        "I need the sin of 1.5708",
        "What's cos(0)?",
        "Please give me tan(0.785)",
        "Add 5 and 3",  # 这个查询应该无法路由
    ]
    
    print("\n" + "="*50)
    print("Starting query tests...")
    print("="*50)
    
    for query in test_queries:
        print(f"\n[Query] '{query}'")
        result = await orchestrator.process_query(query)
        
        if result["success"]:
            agent_response = result["result"]
            # 从响应中提取人类可读的消息
            if agent_response.get('artifacts'):
                for artifact in agent_response['artifacts']:
                    if artifact['type'] == 'text':
                        print(f"  [Response from {result['agent']}]: {artifact['content'].get('message', 'No message')}")
                        print(f"  [Raw Result]: {artifact['content'].get('result')}")
        else:
            print(f"  [Error]: {result['error']}")
    
    print("\n" + "="*50)
    print("Test completed.")
    print("="*50)

if __name__ == "__main__":
    # 检查必要的端口是否被占用(简易检查)
    import socket
    ports = [8001, 8002, 8003]
    for port in ports:
        sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
        result = sock.connect_ex(('localhost', port))
        if result != 0:
            print(f"[Warning] Port {port} is not open. Make sure the corresponding agent is running.")
        sock.close()
    
    asyncio.run(main())

4.3 运行你的第一个A2A网络

现在,让我们把整个系统跑起来。你需要打开四个独立的终端窗口(或终端标签页)。

第一步:启动专家智能体 在第一个终端,激活虚拟环境,运行正弦智能体:

python sine_agent.py

你应该看到输出:Starting SineCalculator on http://localhost:8001

在第二个终端,激活虚拟环境,运行余弦智能体:

python cosine_agent.py

输出:Starting CosineCalculator on http://localhost:8002

在第三个终端,激活虚拟环境,运行正切智能体:

python tangent_agent.py

输出:Starting TangentCalculator on http://localhost:8003

第二步:运行指挥器 在第四个终端,激活虚拟环境,运行主程序:

python main.py

如果一切顺利,你将看到类似以下的输出:

[Main] Assuming agent servers are already running on ports 8001, 8002, 8003.
[Main] If not, please start them in separate terminal windows.

[Main] Registering agents with the orchestrator...
[Orchestrator] Successfully registered agent: SineCalculator at http://localhost:8001
[Orchestrator] Successfully registered agent: CosineCalculator at http://localhost:8002
[Orchestrator] Successfully registered agent: TangentCalculator at http://localhost:8003

==================================================
Starting query tests...
==================================================

[Query] 'What is the sine of 0.7854?'
[Orchestrator] Routing query 'What is the sine of 0.7854?' to 'SineCalculator' (confidence: 0.90)
  [Response from SineCalculator]: The sine of 0.7854 is 0.707077
  [Raw Result]: 0.7070775558117769

[Query] 'Calculate cosine for 1.0472'
[Orchestrator] Routing query 'Calculate cosine for 1.0472' to 'CosineCalculator' (confidence: 0.90)
  [Response from CosineCalculator]: The cosine of 1.0472 is 0.500171
  [Raw Result]: 0.5001710747505848
...
[Query] 'Add 5 and 3'
  [Error]: No suitable agent found for query: 'Add 5 and 3'

恭喜!你已经成功搭建并运行了一个基于A2A协议的微型智能体协作网络。用户向指挥器提出自然语言查询,指挥器理解意图并将其路由到正确的专家智能体,专家智能体完成计算后返回结果。

5. 进阶:从玩具系统到生产级架构

我们构建的系统虽然能跑通流程,但距离一个健壮、可扩展的生产级系统还有很大差距。下面,我们来探讨如何将这个“玩具”升级。

5.1 引入真正的智能路由:集成LLM

我们之前的路由器 _route_by_keyword 太脆弱了。它无法理解“给我0.5的正弦值”和“sin(0.5)是多少?”是同一个意思。在生产环境中,路由决策应该由一个大语言模型来驱动

让我们升级 orchestrator.py,集成一个LLM(这里以OpenAI API为例)来实现真正的意图理解。首先,确保你安装了openai库:pip install openai

# orchestrator.py 第二部分:集成LLM的智能路由器
import openai
import os
from typing import List

class LLMEnhancedOrchestrator(SimpleOrchestrator):
    def __init__(self, openai_api_key: str):
        super().__init__()
        openai.api_key = openai_api_key
        self.llm_client = openai.OpenAI()
        # 为LLM准备智能体能力描述
        self._agent_descriptions = None

    def _build_agent_prompt(self) -> str:
        """构建描述所有可用智能体能力的提示词"""
        if self._agent_descriptions:
            return self._agent_descriptions
        
        descriptions = []
        for name, card in self.network.agents.items():
            desc = f"- {name}: {card.description}. Skills: {', '.join([s['name'] for s in card.skills])}."
            descriptions.append(desc)
        
        self._agent_descriptions = "\n".join(descriptions)
        return self._agent_descriptions

    async def _route_with_llm(self, query: str) -> Tuple[Optional[str], float, Optional[str]]:
        """使用LLM进行路由决策,并返回推理过程"""
        agent_list_text = self._build_agent_prompt()
        
        prompt = f"""你是一个智能任务路由器。请根据用户查询,从以下智能体列表中选择最合适的一个来处理。
如果没有任何智能体适合,请回答“NONE”。

可用智能体:
{agent_list_text}

用户查询:"{query}"

请严格按以下JSON格式回复:
{{
  "selected_agent": "智能体名称 或 NONE",
  "confidence": 0.0到1.0之间的数字,
  "reasoning": "你的推理过程,简短说明为什么选择这个智能体"
}}
"""
        try:
            response = self.llm_client.chat.completions.create(
                model="gpt-3.5-turbo",  # 或 gpt-4
                messages=[{"role": "system", "content": "你是一个精准的任务分配助手。"},
                         {"role": "user", "content": prompt}],
                temperature=0.1,
                max_tokens=200
            )
            import json
            result_text = response.choices[0].message.content.strip()
            # 尝试解析JSON
            result = json.loads(result_text)
            agent_name = result.get("selected_agent")
            confidence = result.get("confidence", 0.0)
            reasoning = result.get("reasoning", "")
            
            if agent_name == "NONE":
                return None, confidence, reasoning
            return agent_name, confidence, reasoning
            
        except json.JSONDecodeError:
            print(f"[LLM Router] Failed to parse LLM response: {result_text}")
            # 降级到关键词路由
            return super()._route_by_keyword(query)[0], 0.5, "Fallback to keyword routing due to LLM parse error."
        except Exception as e:
            print(f"[LLM Router] Error: {e}")
            return super()._route_by_keyword(query)[0], 0.3, f"LLM routing failed: {e}. Using fallback."

    async def process_query_enhanced(self, user_query: str) -> Dict:
        """使用LLM增强版路由处理查询"""
        target_agent_name, confidence, reasoning = await self._route_with_llm(user_query)
        
        if not target_agent_name:
            return {
                "success": False,
                "error": f"No suitable agent found for query: '{user_query}'",
                "llm_reasoning": reasoning,
                "suggested_agents": [card.name for card in self.network.discover("")]
            }
        
        print(f"[LLM Router] Selected '{target_agent_name}' for query: '{user_query}'")
        print(f"           Confidence: {confidence:.2f}, Reasoning: {reasoning}")
        
        # 剩余流程与父类相同...
        # ... [调用智能体并返回结果的代码,与之前类似]

这个增强版路由器利用LLM来理解用户查询的语义,并从智能体列表中做出选择。它甚至能提供选择理由,这极大地提高了系统的可解释性和可靠性。对于“sin(0.5)是多少?”和“计算0.5的正弦值”这样的同义查询,LLM都能准确路由到 SineCalculator

5.2 实现工作流与复杂任务分解

真正的威力在于处理复杂任务。假设用户查询是:“计算30度角的正弦和余弦值,然后告诉我它们的平方和。”

单个智能体无法处理这个任务。我们需要一个工作流管理器来分解任务:

  1. 将“30度”转换为弧度。
  2. 调用 SineCalculator 计算正弦值。
  3. 调用 CosineCalculator 计算余弦值。
  4. 将两个结果平方后相加。
  5. 返回最终答案。

我们可以创建一个新的智能体 MathWorkflowManager,它本身不进行计算,但负责协调其他智能体。

# workflow_manager.py
import asyncio
import math
from a2a_core import A2AAgent, AgentCard, Task, TaskStatus, TaskState, Artifact
from orchestrator import LLMEnhancedOrchestrator

class MathWorkflowManager(A2AAgent):
    """数学工作流管理器,能分解复杂任务并协调多个专家智能体"""
    def __init__(self, port: int, orchestrator: LLMEnhancedOrchestrator):
        card = AgentCard(
            name="MathWorkflowManager",
            description="Orchestrates complex mathematical calculations by decomposing them and coordinating multiple specialist agents.",
            version="1.0.0",
            skills=[{
                "name": "orchestrate_calculation",
                "description": "Decompose a complex math query, call relevant agents, and synthesize results.",
                "input_schema": {"type": "object", "properties": {"query": {"type": "string"}}},
                "output_schema": {"type": "object", "properties": {"final_result": {"type": "number"}, "steps": {"type": "array"}}}
            }]
        )
        super().__init__(card)
        self.port = port
        self.card.endpoint = f"http://localhost:{port}"
        self.orchestrator = orchestrator

    async def handle_task(self, task: Task) -> Task:
        query = task.message.get("content", {}).get("text", "")
        print(f"[WorkflowManager] Processing complex query: {query}")
        
        # 这里可以集成一个LLM来动态生成执行计划。
        # 为了简化,我们硬编码处理一种特定类型的复杂查询。
        if "sine and cosine" in query.lower() and "square sum" in query.lower():
            # 示例:处理“计算30度角的正弦和余弦值,然后告诉我它们的平方和。”
            try:
                # 1. 提取角度(简化:假设是数字)
                import re
                angle_match = re.search(r"(\d+)", query)
                if not angle_match:
                    task.status = TaskStatus(state=TaskState.FAILED, error_message="Could not extract angle from query.")
                    return task
                angle_deg = float(angle_match.group(1))
                angle_rad = math.radians(angle_deg)
                
                # 2. 分解任务并调用子智能体
                subtasks = [
                    (f"Calculate sine of {angle_rad}", "SineCalculator"),
                    (f"Calculate cosine of {angle_rad}", "CosineCalculator"),
                ]
                
                results = {}
                for subquery, agent_name in subtasks:
                    # 使用指挥器来路由和调用子任务
                    sub_result = await self.orchestrator.process_query(subquery)
                    if sub_result["success"]:
                        # 从响应中提取数值结果
                        agent_resp = sub_result["result"]
                        for artifact in agent_resp.get('artifacts', []):
                            if artifact['type'] == 'text':
                                results[agent_name] = artifact['content'].get('result')
                                print(f"  [Subtask Result] {agent_name}: {results[agent_name]}")
                    else:
                        task.status = TaskStatus(state=TaskState.FAILED, error_message=f"Subtask failed: {sub_result.get('error')}")
                        return task
                
                # 3. 合成最终结果
                sine_val = results.get("SineCalculator", 0)
                cosine_val = results.get("CosineCalculator", 0)
                final_result = sine_val**2 + cosine_val**2
                
                # 4. 返回结果
                task.artifacts.append(Artifact(
                    type="text",
                    content={
                        "final_result": final_result,
                        "message": f"For angle {angle_deg}°, sin² + cos² = {final_result:.10f} (Should be close to 1.0 due to trigonometric identity).",
                        "steps": [
                            f"Extracted angle: {angle_deg}° ({angle_rad:.4f} rad)",
                            f"Sine({angle_rad:.4f}) = {sine_val:.6f} (via SineCalculator)",
                            f"Cosine({angle_rad:.4f}) = {cosine_val:.6f} (via CosineCalculator)",
                            f"Sum of squares: {sine_val:.6f}² + {cosine_val:.6f}² = {final_result:.10f}"
                        ]
                    }
                ))
                task.status = TaskStatus(state=TaskState.COMPLETED)
                
            except Exception as e:
                task.status = TaskStatus(state=TaskState.FAILED, error_message=f"Workflow execution error: {e}")
        else:
            task.status = TaskStatus(state=TaskState.FAILED, error_message="This workflow manager currently only handles 'sine and cosine square sum' queries.")
        
        return task

这个工作流管理器展示了A2A架构的真正力量:任务分解与协同。它自己并不做具体计算,而是作为协调者,将复杂问题拆解,调度合适的专家智能体解决子问题,最后汇总结果。你可以在此基础上扩展,集成LLM来动态生成任意复杂的工作流。

5.3 生产级考量与最佳实践

将我们的玩具系统推向生产环境,还需要考虑许多方面:

1. 服务发现与健康检查 我们的注册是手动的。在生产中,应使用服务发现机制(如Consul, etcd, 或Kubernetes Services)。智能体启动时应自动向注册中心注册,并定期发送心跳。路由器应能感知智能体的健康状态,避免将任务路由到故障节点。

2. 安全与认证 智能体间的通信必须加密(HTTPS/WSS)。需要实现身份认证和授权机制,确保只有合法的智能体可以加入网络,并且智能体只能访问其被授权的服务。

3. 可观测性与监控 需要记录所有任务的执行链路、耗时、成功率。集成像Prometheus和Grafana这样的监控工具,以及像Jaeger这样的分布式追踪系统,对于调试和保障SLA至关重要。

4. 异步与流式处理 对于长任务,应支持异步处理和流式进度更新。A2A协议支持流式响应,允许智能体在计算过程中持续返回中间状态。

5. 错误处理与重试 网络调用可能失败。需要实现完善的错误处理、重试逻辑和回退策略。例如,当首选智能体无响应时,可以尝试路由到具有相似能力的备用智能体。

6. 配置与版本管理 智能体的能力(AgentCard)可能随版本更新而变化。路由器需要能处理不同版本的智能体,并可能根据版本进行路由决策。

构建一个成熟的A2A智能体网络是一项系统工程,但它带来的好处是巨大的:极高的模块化、可扩展性和灵活性。你可以像搭积木一样,通过组合不同的智能体,快速构建出应对各种复杂场景的AI应用。

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐