从零构建AI Agent框架:多智能体协作与MCP/A2A协议实战
·
这次我们来深入探讨如何从零构建一个完整的AI Agent框架,重点是多智能体协作、MCP协议集成和A2A协议实战。这个框架基于DeepSeek模型,可以在普通开发环境中运行,不需要高端硬件配置。
AI Agent开发正成为大模型应用的主流方向,但很多开发者面临协议复杂、框架臃肿、部署困难等问题。本文将通过实战演示如何构建一个轻量级但功能完整的AI Agent框架,支持多智能体协作、工具调用和跨平台通信。
1. 核心能力速览
| 能力项 | 具体说明 |
|---|---|
| 框架类型 | 多智能体协作框架,支持MCP+A2A协议 |
| 核心模型 | DeepSeek系列模型(支持API调用和本地部署) |
| 协议支持 | MCP协议(工具调用)、A2A协议(智能体间通信) |
| 硬件要求 | 普通开发机即可,CPU/GPU均可运行 |
| 显存需求 | 根据DeepSeek模型版本调整,最小2GB显存可运行基础版本 |
| 启动方式 | 命令行启动 + WebUI界面 |
| API支持 | 完整的RESTful API接口 |
| 批量任务 | 支持并发多智能体任务处理 |
| 适合场景 | 企业内部自动化、多系统集成、复杂工作流处理 |
2. AI Agent框架架构设计
我们的框架采用分层架构设计,确保各组件职责清晰且易于扩展。
2.1 核心架构组件
# 框架核心类结构示例
class AIAgentFramework:
def __init__(self):
self.agents = {} # 智能体管理器
self.mcp_servers = {} # MCP工具服务器
self.a2a_broker = A2ABroker() # A2A通信代理
self.task_queue = TaskQueue() # 任务队列
class BaseAgent:
def __init__(self, agent_id, llm_config):
self.agent_id = agent_id
self.llm = LLMClient(llm_config) # 大模型客户端
self.memory = MemoryManager() # 记忆管理
self.tools = ToolRegistry() # 工具注册表
2.2 多智能体协作流程
框架支持多种协作模式:
- 主从模式 :一个主智能体协调多个从属智能体
- 对等模式 :智能体间平等协作,共同完成任务
- 流水线模式 :任务按阶段在不同智能体间传递
3. 环境准备与依赖安装
3.1 基础环境要求
# 检查Python环境
python --version # 需要Python 3.8+
pip --version
# 安装核心依赖
pip install fastapi uvicorn pydantic
pip install requests httpx
pip install python-dotenv
3.2 DeepSeek模型配置
# config.py - 模型配置
DEEPSEEK_CONFIG = {
"api_key": "your_deepseek_api_key",
"base_url": "https://api.deepseek.com/v1",
"model": "deepseek-chat",
"max_tokens": 4096,
"temperature": 0.7
}
# 本地模型配置(如果使用本地部署)
LOCAL_MODEL_CONFIG = {
"model_path": "/path/to/local/model",
"device": "cuda", # 或 "cpu"
"gpu_memory": "4GB" # 显存配置
}
3.3 MCP协议环境搭建
MCP协议需要特定的运行时环境,我们提供两种安装方式:
# 方式1:使用uvx(Python环境)
pip install uv
uvx --version # 验证安装
# 方式2:使用npx(Node.js环境)
npm install -g npx
npx --version # 验证安装
4. MCP协议集成实战
MCP协议是智能体与外部工具交互的核心,我们来实现完整的MCP服务器和客户端。
4.1 MCP服务器实现
# mcp_server.py - 基础MCP服务器
import asyncio
from mcp import MCPServer, ToolRegistry
class MyMCPServer(MCPServer):
def __init__(self):
self.tools = ToolRegistry()
self.setup_tools()
def setup_tools(self):
# 注册搜索工具
@self.tools.register("web_search")
async def web_search(query: str) -> str:
"""执行网页搜索"""
# 实现搜索逻辑
return f"搜索结果: {query}"
# 注册文件操作工具
@self.tools.register("file_operation")
async def file_operation(action: str, path: str) -> str:
"""文件操作工具"""
# 实现文件操作逻辑
return f"文件操作完成: {action} {path}"
# 启动MCP服务器
async def main():
server = MyMCPServer()
await server.serve()
if __name__ == "__main__":
asyncio.run(main())
4.2 MCP客户端集成
# mcp_client.py - MCP客户端
class MCPClient:
def __init__(self, server_url):
self.server_url = server_url
self.session = requests.Session()
async def list_tools(self):
"""获取服务器支持的工具列表"""
response = await self.session.post(
f"{self.server_url}/tools/list"
)
return response.json()
async def call_tool(self, tool_name, arguments):
"""调用特定工具"""
payload = {
"name": tool_name,
"arguments": arguments
}
response = await self.session.post(
f"{self.server_url}/tools/call",
json=payload
)
return response.json()
5. A2A协议实现多智能体通信
A2A协议使得不同智能体之间能够自然协作,我们来实现完整的A2A通信层。
5.1 A2A协议基础实现
# a2a_protocol.py - A2A协议实现
from typing import Dict, Any, List
import json
class A2AMessage:
def __init__(self, sender: str, receiver: str, content: Dict):
self.sender = sender
self.receiver = receiver
self.content = content
self.timestamp = datetime.now().isoformat()
def to_dict(self) -> Dict[str, Any]:
return {
"sender": self.sender,
"receiver": self.receiver,
"content": self.content,
"timestamp": self.timestamp
}
class A2ABroker:
def __init__(self):
self.agents: Dict[str, A2AAgent] = {}
self.message_queue: asyncio.Queue = asyncio.Queue()
def register_agent(self, agent_id: str, agent: 'A2AAgent'):
"""注册智能体"""
self.agents[agent_id] = agent
async def send_message(self, message: A2AMessage):
"""发送消息到目标智能体"""
if message.receiver in self.agents:
await self.agents[message.receiver].receive_message(message)
else:
print(f"目标智能体不存在: {message.receiver}")
5.2 多智能体协作示例
# multi_agent_example.py - 多智能体协作实战
class TravelPlannerAgent:
def __init__(self, agent_id: str, a2a_broker: A2ABroker):
self.agent_id = agent_id
self.broker = a2a_broker
self.broker.register_agent(agent_id, self)
async def plan_travel(self, destination: str, dates: str):
"""旅行规划主逻辑"""
# 1. 查询天气
weather_msg = A2AMessage(
self.agent_id, "weather_agent",
{"action": "get_weather", "destination": destination, "dates": dates}
)
await self.broker.send_message(weather_msg)
# 2. 查询交通
transport_msg = A2AMessage(
self.agent_id, "transport_agent",
{"action": "find_transport", "destination": destination}
)
await self.broker.send_message(transport_msg)
class WeatherAgent:
def __init__(self, agent_id: str, a2a_broker: A2ABroker):
self.agent_id = agent_id
self.broker = a2a_broker
self.broker.register_agent(agent_id, self)
async def receive_message(self, message: A2AMessage):
"""处理天气查询请求"""
if message.content.get("action") == "get_weather":
# 调用天气API
weather_data = await self.get_weather_data(
message.content["destination"],
message.content["dates"]
)
# 回复消息
response = A2AMessage(
self.agent_id, message.sender,
{"weather_data": weather_data}
)
await self.broker.send_message(response)
6. DeepSeek模型集成与优化
6.1 模型客户端封装
# deepseek_client.py - DeepSeek模型客户端
import openai
from typing import List, Dict
class DeepSeekClient:
def __init__(self, config: Dict):
self.config = config
self.client = openai.OpenAI(
api_key=config["api_key"],
base_url=config["base_url"]
)
async def chat_completion(self, messages: List[Dict], **kwargs) -> str:
"""聊天补全接口"""
try:
response = self.client.chat.completions.create(
model=self.config["model"],
messages=messages,
max_tokens=kwargs.get("max_tokens", 2048),
temperature=kwargs.get("temperature", 0.7)
)
return response.choices[0].message.content
except Exception as e:
print(f"DeepSeek API调用失败: {e}")
return None
async def stream_chat(self, messages: List[Dict], **kwargs):
"""流式聊天接口"""
response = self.client.chat.completions.create(
model=self.config["model"],
messages=messages,
stream=True,
**kwargs
)
for chunk in response:
if chunk.choices[0].delta.content is not None:
yield chunk.choices[0].delta.content
6.2 智能体推理引擎
# agent_brain.py - 智能体推理核心
class AgentBrain:
def __init__(self, llm_client, tools_registry):
self.llm = llm_client
self.tools = tools_registry
self.conversation_history = []
async def process_query(self, query: str, context: Dict = None) -> str:
"""处理用户查询"""
# 构建系统提示词
system_prompt = self._build_system_prompt()
# 构建消息历史
messages = [
{"role": "system", "content": system_prompt},
*self.conversation_history,
{"role": "user", "content": query}
]
# 调用LLM
response = await self.llm.chat_completion(messages)
# 解析工具调用
tool_calls = self._parse_tool_calls(response)
if tool_calls:
results = await self._execute_tools(tool_calls)
response = await self._integrate_tool_results(response, results)
# 更新对话历史
self.conversation_history.append({"role": "user", "content": query})
self.conversation_history.append({"role": "assistant", "content": response})
return response
def _build_system_prompt(self) -> str:
"""构建系统提示词"""
available_tools = "\n".join([f"- {name}: {desc}" for name, desc in self.tools.list_tools()])
return f"""你是一个AI助手,可以调用以下工具:
{available_tools}
如果用户请求需要工具调用,请按照格式响应。
保持回答专业、准确。"""
7. 框架部署与启动
7.1 一键启动脚本
#!/bin/bash
# start_framework.sh - 框架启动脚本
echo "正在启动AI Agent框架..."
# 检查环境变量
if [ -z "$DEEPSEEK_API_KEY" ]; then
echo "错误: 请设置DEEPSEEK_API_KEY环境变量"
exit 1
fi
# 启动MCP服务器
echo "启动MCP服务器..."
python mcp_server.py &
# 启动A2A消息代理
echo "启动A2A消息代理..."
python a2a_broker.py &
# 启动Web界面
echo "启动Web界面..."
uvicorn web_interface:app --host 0.0.0.0 --port 8000 &
echo "框架启动完成!"
echo "Web界面: http://localhost:8000"
echo "API文档: http://localhost:8000/docs"
7.2 Docker部署配置
# Dockerfile
FROM python:3.9-slim
WORKDIR /app
# 复制依赖文件
COPY requirements.txt .
RUN pip install -r requirements.txt
# 复制源代码
COPY . .
# 设置环境变量
ENV DEEPSEEK_API_KEY=your_api_key_here
ENV PYTHONPATH=/app
# 暴露端口
EXPOSE 8000
# 启动命令
CMD ["python", "main.py"]
8. 功能测试与效果验证
8.1 基础功能测试
# test_basic_functionality.py - 基础功能测试
import asyncio
import pytest
class TestAIAgentFramework:
@pytest.fixture
async def framework(self):
"""框架测试夹具"""
framework = AIAgentFramework()
await framework.initialize()
return framework
async def test_agent_creation(self, framework):
"""测试智能体创建"""
agent = await framework.create_agent("test_agent", "deepseek-chat")
assert agent is not None
assert agent.agent_id == "test_agent"
async def test_mcp_tool_calling(self, framework):
"""测试MCP工具调用"""
agent = await framework.create_agent("tool_agent", "deepseek-chat")
result = await agent.call_tool("web_search", {"query": "今天天气"})
assert "搜索" in result or "天气" in result
async def test_a2a_communication(self, framework):
"""测试A2A通信"""
agent1 = await framework.create_agent("agent1", "deepseek-chat")
agent2 = await framework.create_agent("agent2", "deepseek-chat")
# 发送测试消息
message = A2AMessage("agent1", "agent2", {"test": "data"})
await framework.a2a_broker.send_message(message)
# 验证消息接收
# 这里需要实现消息接收验证逻辑
8.2 集成测试场景
# test_integration_scenarios.py - 集成测试
async def test_travel_planning_scenario():
"""测试旅行规划场景"""
framework = AIAgentFramework()
await framework.initialize()
# 创建旅行规划相关智能体
travel_agent = await framework.create_agent("travel_planner", "deepseek-chat")
weather_agent = await framework.create_agent("weather_agent", "deepseek-chat")
transport_agent = await framework.create_agent("transport_agent", "deepseek-chat")
# 执行旅行规划
result = await travel_agent.process_query(
"帮我规划本周五到周日的北京旅行,需要知道天气和交通信息"
)
# 验证结果包含关键信息
assert "北京" in result
assert "天气" in result or "交通" in result
print("旅行规划测试通过!")
9. 性能优化与资源管理
9.1 并发处理优化
# performance_optimizer.py - 性能优化
import asyncio
from concurrent.futures import ThreadPoolExecutor
class PerformanceOptimizer:
def __init__(self, max_workers=10):
self.executor = ThreadPoolExecutor(max_workers=max_workers)
self.semaphore = asyncio.Semaphore(5) # 控制并发数
async def process_batch_requests(self, requests: List[Dict]):
"""批量处理请求"""
async with self.semaphore:
tasks = [self._process_single_request(req) for req in requests]
results = await asyncio.gather(*tasks, return_exceptions=True)
return results
async def _process_single_request(self, request: Dict):
"""处理单个请求"""
# 实现请求处理逻辑
pass
9.2 内存和显存管理
# resource_manager.py - 资源管理
import psutil
import GPUtil
class ResourceManager:
@staticmethod
def get_system_resources():
"""获取系统资源信息"""
memory = psutil.virtual_memory()
gpus = GPUtil.getGPUs()
return {
"memory_used": memory.used,
"memory_total": memory.total,
"gpu_memory": [gpu.memoryUsed for gpu in gpus] if gpus else []
}
@staticmethod
def check_resource_limits():
"""检查资源限制"""
resources = ResourceManager.get_system_resources()
memory_usage = resources["memory_used"] / resources["memory_total"]
if memory_usage > 0.8:
print("警告: 内存使用率过高,建议优化")
return False
return True
10. 实际应用案例
10.1 企业内部自动化流程
# enterprise_automation.py - 企业自动化案例
class EnterpriseAutomation:
def __init__(self, framework):
self.framework = framework
self.setup_enterprise_agents()
def setup_enterprise_agents(self):
"""设置企业级智能体"""
# CRM智能体
self.crm_agent = self.framework.create_agent("crm_agent", "deepseek-chat")
# 财务智能体
self.finance_agent = self.framework.create_agent("finance_agent", "deepseek-chat")
# 报告生成智能体
self.report_agent = self.framework.create_agent("report_agent", "deepseek-chat")
async def generate_quarterly_report(self):
"""生成季度报告"""
# 1. 从CRM获取客户数据
crm_data = await self.crm_agent.call_tool("get_customer_data", {"period": "Q1"})
# 2. 从财务系统获取财务数据
finance_data = await self.finance_agent.call_tool("get_finance_data", {"quarter": "Q1"})
# 3. 生成综合报告
report = await self.report_agent.process_query(
f"基于以下数据生成季度报告:\nCRM数据: {crm_data}\n财务数据: {finance_data}"
)
return report
10.2 多系统集成示例
# system_integration.py - 多系统集成
class SystemIntegration:
def __init__(self):
self.setup_integration_agents()
def setup_integration_agents(self):
"""设置系统集成智能体"""
# GitLab集成
self.gitlab_agent = self.create_agent_with_tools(
"gitlab_agent",
["gitlab_merge", "gitlab_deploy"]
)
# Jenkins集成
self.jenkins_agent = self.create_agent_with_tools(
"jenkins_agent",
["jenkins_build", "jenkins_deploy"]
)
# 通知集成
self.notification_agent = self.create_agent_with_tools(
"notification_agent",
["send_slack", "send_email"]
)
async def automated_deployment_pipeline(self, project_name, environment):
"""自动化部署流水线"""
# 1. GitLab代码合并
merge_result = await self.gitlab_agent.call_tool(
"gitlab_merge",
{"project": project_name, "branch": "main"}
)
# 2. Jenkins构建部署
build_result = await self.jenkins_agent.call_tool(
"jenkins_build",
{"project": project_name, "environment": environment}
)
# 3. 发送部署通知
await self.notification_agent.call_tool(
"send_slack",
{"message": f"项目{project_name}已成功部署到{environment}"}
)
return {"merge": merge_result, "build": build_result}
11. 常见问题与解决方案
11.1 部署问题排查
| 问题现象 | 可能原因 | 解决方案 |
|---|---|---|
| MCP服务器启动失败 | 端口被占用或依赖缺失 | 检查端口占用,重新安装依赖 |
| DeepSeek API调用失败 | API密钥错误或网络问题 | 验证API密钥,检查网络连接 |
| 智能体通信超时 | 消息队列阻塞或配置错误 | 调整超时设置,检查消息队列 |
| 内存使用过高 | 并发任务过多或内存泄漏 | 减少并发数,检查代码内存使用 |
11.2 性能优化建议
# troubleshooting.py - 问题排查工具
import logging
from datetime import datetime
class Troubleshooter:
def __init__(self):
self.logger = logging.getLogger("ai_agent_framework")
async def diagnose_issue(self, error_message: str, context: Dict) -> str:
"""诊断问题并提供解决方案"""
common_issues = {
"api_key": "检查API密钥配置",
"network": "验证网络连接和代理设置",
"memory": "检查系统内存使用情况",
"timeout": "调整请求超时设置"
}
for issue, solution in common_issues.items():
if issue in error_message.lower():
return f"检测到{issue}问题: {solution}"
return "请查看详细日志或联系技术支持"
12. 扩展开发与自定义
12.1 自定义工具开发
# custom_tools.py - 自定义工具示例
class CustomToolRegistry:
def __init__(self):
self.tools = {}
def register_tool(self, name: str, description: str, function: callable):
"""注册自定义工具"""
self.tools[name] = {
"description": description,
"function": function
}
async def execute_tool(self, name: str, arguments: Dict) -> str:
"""执行自定义工具"""
if name not in self.tools:
return f"工具不存在: {name}"
try:
result = await self.tools[name]["function"](**arguments)
return str(result)
except Exception as e:
return f"工具执行失败: {str(e)}"
# 示例自定义工具
async def database_query_tool(query: str, database: str = "default") -> str:
"""数据库查询工具"""
# 实现数据库查询逻辑
return f"查询结果: 执行了 {query} 在 {database} 数据库"
async def external_api_tool(api_endpoint: str, parameters: Dict) -> str:
"""外部API调用工具"""
# 实现API调用逻辑
return f"API调用结果: {api_endpoint}"
12.2 框架扩展接口
# extension_interface.py - 扩展接口
from abc import ABC, abstractmethod
class AgentExtension(ABC):
"""智能体扩展基类"""
@abstractmethod
async def initialize(self, agent):
"""初始化扩展"""
pass
@abstractmethod
async def process_message(self, message, context):
"""处理消息"""
pass
class MonitoringExtension(AgentExtension):
"""监控扩展"""
async def initialize(self, agent):
self.agent = agent
self.metrics = {}
async def process_message(self, message, context):
# 记录监控指标
self.metrics["message_count"] = self.metrics.get("message_count", 0) + 1
return message
class LoggingExtension(AgentExtension):
"""日志扩展"""
async def initialize(self, agent):
self.logger = logging.getLogger(f"agent_{agent.agent_id}")
async def process_message(self, message, context):
self.logger.info(f"处理消息: {message}")
return message
这个AI Agent框架提供了从基础协议实现到高级多智能体协作的完整解决方案。通过MCP协议标准化工具调用,通过A2A协议实现智能体间自然通信,结合DeepSeek模型的强大推理能力,可以构建出真正实用的AI应用系统。
框架设计注重实用性和可扩展性,开发者可以根据具体需求灵活定制智能体行为、工具集和协作模式。无论是简单的自动化任务还是复杂的企业级工作流,都能在这个框架中找到合适的实现方案。
建议从基础的单智能体工具调用开始实践,逐步扩展到多智能体协作场景,最终实现完整的AI Agent应用生态系统。
更多推荐

所有评论(0)