2026年AI Agent开发框架实战指南:从零构建多Agent协作系统
本文为原创技术文章,基于2026年AI Agent技术生态现状,从环境搭建到生产部署,带你从零构建一套完整的多Agent协作系统。
一、概述:为什么AI Agent是2026年的技术核心
1.1 AI Agent的发展现状
2026年,AI Agent已经从概念验证阶段全面进入工程化落地阶段。与2023-2024年的"提示词工程"热潮不同,当前的AI Agent开发更强调工程化、可观测性、安全性和团队协作。
根据行业调研数据,超过60%的中大型企业已经在生产环境中部署了至少一个AI Agent应用,主要集中在以下场景:
1.2 主流框架对比
2026年的AI Agent框架生态已经形成了清晰的分层格局:
| 框架 | 定位 | 核心优势 | 适用场景 |
|------|------|----------|----------|
| LangChain / LangGraph | 通用编排框架 | 生态最完善,工具链丰富 | 复杂工作流、生产级应用 |
| AutoGen | 多Agent对话框架 | 多角色对话原生支持 | 多Agent协商、群体决策 |
| CrewAI | 角色化Agent框架 | 角色定义简洁,任务分配清晰 | 团队协作、任务型Agent |
| DeepSeek Harness | 国产推理框架 | 本地部署友好,中文优化 | 私有化部署、企业内网 |
| OpenAI Codex Harness | 代码生成框架 | 代码理解与生成能力强 | 代码Agent、DevOps场景 |
1.3 关键技术趋势
二、环境搭建:从零开始准备开发环境
2.1 基础环境要求
在开始之前,请确保你的开发环境满足以下要求:
# Python 版本要求 3.10+
python --version # Python 3.11.7
# 推荐使用虚拟环境管理工具
# 可选:venv / conda / poetry / uv2.2 核心依赖安装
本文我们将以 LangGraph(LangChain的Agent编排层) 为核心框架,结合MCP协议和本地向量存储,构建一套完整的多Agent系统。
# 创建项目目录
mkdir multi-agent-system && cd multi-agent-system
# 使用uv创建虚拟环境(推荐,速度更快)
uv venv
source .venv/bin/activate
# 安装核心依赖
uv pip install langchain langgraph langchain-openai
uv pip install langchain-community langchain-core
uv pip install pydantic python-dotenv
# 安装向量存储与记忆相关
uv pip install chromadb sentence-transformers
# 安装MCP协议相关
uv pip install mcp
# 安装工具与沙箱相关
uv pip install python-docker requests
# 安装部署相关
uv pip install fastapi uvicorn2.3 配置文件与环境变量
创建 .env 文件管理配置:
# .env
# 大模型配置(支持OpenAI兼容接口,可替换为本地模型)
LLM_API_KEY=your-api-key
LLM_BASE_URL=https://api.openai.com/v1
LLM_MODEL=gpt-4o-mini
# 嵌入模型配置
EMBEDDING_MODEL=BAAI/bge-small-zh-v1.5
# 向量数据库路径
CHROMA_DB_PATH=./data/chroma
# 沙箱配置
SANDBOX_ENABLED=true
SANDBOX_TIMEOUT=302.4 项目目录结构
推荐的项目目录结构如下:
multi-agent-system/
├── .env # 环境变量配置
├── requirements.txt # 依赖清单
├── main.py # 入口文件
├── agents/ # Agent定义
│ ├── __init__.py
│ ├── base.py # 基础Agent类
│ ├── researcher.py # 研究员Agent
│ ├── writer.py # 写作Agent
│ └── reviewer.py # 审稿Agent
├── tools/ # 工具定义
│ ├── __init__.py
│ ├── search.py # 搜索工具
│ ├── calculator.py # 计算工具
│ └── code_executor.py # 代码执行工具
├── memory/ # 记忆系统
│ ├── __init__.py
│ ├── short_term.py # 短期记忆
│ └── long_term.py # 长期记忆
├── workflows/ # 工作流编排
│ ├── __init__.py
│ └── research_flow.py # 研究工作流
├── deployment/ # 部署相关
│ ├── server.py # API服务
│ └── Dockerfile
└── data/ # 数据存储
├── chroma/ # 向量数据库
└── logs/ # 日志目录三、核心概念:理解Agent的构成要素
3.1 Agent的核心组件
一个完整的AI Agent由以下核心组件构成:
┌─────────────────────────────────────────────────┐
│ Agent │
│ ┌─────────┐ ┌─────────┐ ┌─────────────────┐ │
│ │ LLM │ │ Memory │ │ Tools / MCP │ │
│ │ 大脑 │ │ 记忆 │ │ 工具集 │ │
│ └────┬────┘ └────┬────┘ └────────┬────────┘ │
│ │ │ │ │
│ ┌────┴────────────┴────────────────┴───────┐ │
│ │ Planning / Reasoning │ │
│ │ 规划与推理引擎 │ │
│ └──────────────────┬───────────────────────┘ │
│ │ │
│ ┌──────────────────┴───────────────────────┐ │
│ │ Action Execution │ │
│ │ 动作执行器 │ │
│ └──────────────────────────────────────────┘ │
└─────────────────────────────────────────────────┘各组件职责:
3.2 ReAct范式:思考与行动的循环
当前主流的Agent推理范式是 ReAct(Reasoning + Acting),其核心循环如下:
思考(Thought) → 行动(Action) → 观察(Observation) → 思考(Thought) → ...用伪代码表示:
def react_loop(question, max_steps=10):
scratchpad = ""
for step in range(max_steps):
# 1. 思考:基于当前信息推理下一步
thought = llm.generate_thought(question, scratchpad)
# 2. 行动:决定要调用的工具和参数
action = llm.generate_action(thought, tools)
# 3. 观察:执行工具并获取结果
observation = execute_tool(action)
# 4. 记录到草稿板
scratchpad += f"
Thought: {thought}
Action: {action}
Observation: {observation}"
# 5. 判断是否完成
if is_final_answer(observation):
return extract_answer(observation)
return "超出最大步数限制"3.3 多Agent协作模式
多Agent系统常见的协作模式有三种:
本文将实现一个层级+流水线混合的多Agent系统,用于自动化研究报告生成。
四、单Agent实现:打造你的第一个智能体
4.1 基础Agent类设计
首先,我们定义一个可扩展的基础Agent类:
# agents/base.py
from __future__ import annotations
from typing import List, Dict, Any, Optional
from dataclasses import dataclass, field
from langchain_core.messages import BaseMessage, HumanMessage, SystemMessage
from langchain_openai import ChatOpenAI
from memory.short_term import ShortTermMemory
from tools import ToolRegistry
@dataclass
class AgentConfig:
"""Agent配置"""
name: str
role: str
system_prompt: str
model: str = "gpt-4o-mini"
temperature: float = 0.7
max_tokens: int = 2048
tool_names: List[str] = field(default_factory=list)
memory_window: int = 20
class BaseAgent:
"""基础Agent类"""
def __init__(self, config: AgentConfig, tool_registry: ToolRegistry):
self.config = config
self.tool_registry = tool_registry
self.llm = ChatOpenAI(
model=config.model,
temperature=config.temperature,
max_tokens=config.max_tokens,
)
self.memory = ShortTermMemory(window_size=config.memory_window)
self._init_tools()
def _init_tools(self):
"""初始化可用工具"""
self.tools = []
for tool_name in self.config.tool_names:
tool = self.tool_registry.get_tool(tool_name)
if tool:
self.tools.append(tool)
# 将工具绑定到LLM
if self.tools:
self.llm = self.llm.bind_tools(self.tools)
def _build_system_prompt(self) -> SystemMessage:
"""构建系统提示词"""
prompt = f"""你是 {self.config.name},你的角色是:{self.config.role}
{self.config.system_prompt}
请严格按照你的角色定位来回应。如果需要调用工具,请使用工具调用功能。
"""
return SystemMessage(content=prompt)
async def run(self, user_input: str) -> str:
"""运行Agent,执行一次完整的思考-行动循环"""
# 构建消息列表
messages = [self._build_system_prompt()]
messages.extend(self.memory.get_history())
messages.append(HumanMessage(content=user_input))
# 调用LLM
response = await self.llm.ainvoke(messages)
# 处理工具调用(支持多轮工具调用)
while response.tool_calls:
messages.append(response)
for tool_call in response.tool_calls:
result = await self.tool_registry.execute_tool(
tool_call["name"], tool_call["args"]
)
messages.append(
{
"role": "tool",
"tool_call_id": tool_call["id"],
"content": str(result),
}
)
response = await self.llm.ainvoke(messages)
# 保存到记忆
self.memory.add_message(HumanMessage(content=user_input))
self.memory.add_message(response)
return response.content4.2 短期记忆实现
# memory/short_term.py
from collections import deque
from typing import List
from langchain_core.messages import BaseMessage
class ShortTermMemory:
"""短期记忆:基于滑动窗口的对话历史管理"""
def __init__(self, window_size: int = 20):
self.window_size = window_size
self._history: deque = deque(maxlen=window_size)
def add_message(self, message: BaseMessage):
"""添加一条消息到记忆"""
self._history.append(message)
def get_history(self) -> List[BaseMessage]:
"""获取当前记忆中的所有消息"""
return list(self._history)
def clear(self):
"""清空记忆"""
self._history.clear()
def __len__(self) -> int:
return len(self._history)4.3 工具注册中心
# tools/__init__.py
from __future__ import annotations
from typing import Dict, Any, Callable, Optional
from langchain_core.tools import tool, BaseTool
import importlib
class ToolRegistry:
"""工具注册中心,支持MCP协议兼容的工具管理"""
def __init__(self):
self._tools: Dict[str, BaseTool] = {}
def register(self, tool_func: BaseTool):
"""注册一个工具"""
self._tools[tool_func.name] = tool_func
return tool_func
def get_tool(self, name: str) -> Optional[BaseTool]:
"""根据名称获取工具"""
return self._tools.get(name)
def list_tools(self) -> list[BaseTool]:
"""列出所有已注册的工具"""
return list(self._tools.values())
async def execute_tool(self, name: str, args: Dict[str, Any]) -> Any:
"""执行指定工具"""
tool = self._tools.get(name)
if not tool:
raise ValueError(f"Tool '{name}' not found")
return await tool.ainvoke(args)
# 全局工具注册中心实例
tool_registry = ToolRegistry()
def register_module_tools(module_path: str):
"""从模块中自动注册工具(使用@tool装饰器的函数)"""
module = importlib.import_module(module_path)
for attr_name in dir(module):
attr = getattr(module, attr_name)
if isinstance(attr, BaseTool):
tool_registry.register(attr)4.4 一个简单的计算器工具
# tools/calculator.py
from langchain_core.tools import tool
import math
@tool
def calculator(expression: str) -> str:
"""
安全的数学计算器,支持常见的数学运算。
支持的运算:加减乘除、幂运算、平方根、三角函数等。
Args:
expression: 数学表达式字符串,例如 "2 + 3 * 4" 或 "sqrt(16)"
"""
# 安全的命名空间,只允许数学相关操作
safe_namespace = {
"abs": abs,
"round": round,
"min": min,
"max": max,
"sqrt": math.sqrt,
"pow": math.pow,
"sin": math.sin,
"cos": math.cos,
"tan": math.tan,
"log": math.log,
"log10": math.log10,
"pi": math.pi,
"e": math.e,
}
try:
# 使用eval但限制命名空间,防止代码注入
result = eval(expression, {"__builtins__": {}}, safe_namespace)
return f"计算结果: {result}"
except Exception as e:
return f"计算错误: {str(e)}"4.5 测试单Agent
创建一个简单的测试入口:
# main_single_agent.py
import asyncio
from dotenv import load_dotenv
load_dotenv()
from agents.base import BaseAgent, AgentConfig
from tools import tool_registry
from tools.calculator import calculator
# 注册工具
tool_registry.register(calculator)
async def main():
# 创建一个数学助手Agent
config = AgentConfig(
name="数学助手",
role="专业的数学问题解答者",
system_prompt="你是一位专业的数学老师,擅长用简洁明了的方式解答数学问题。当遇到需要计算的问题时,请使用计算器工具。",
tool_names=["calculator"],
)
agent = BaseAgent(config, tool_registry)
# 测试
questions = [
"1234的平方根是多少?",
"如果一个圆的半径是7,它的面积是多少?",
"sin(30度)等于多少?",
]
for q in questions:
print(f"
用户: {q}")
answer = await agent.run(q)
print(f"助手: {answer}")
if __name__ == "__main__":
asyncio.run(main())运行结果示例:
用户: 1234的平方根是多少?
助手: 1234的平方根约等于 35.13(精确值为 35.12833614050059)。五、多Agent协作:构建研究报告生成系统
5.1 系统架构设计
我们将构建一个三Agent协作的研究报告生成系统:
用户请求
│
▼
┌─────────────────┐
│ 协调者Agent │ 分配任务、整合结果
│ (Coordinator) │
└───┬─────────┬───┘
│ │
▼ ▼
┌────────┐ ┌────────┐
│研究员 │ │ 写作 │
│Agent │ │ Agent │
│(搜索、 │ │(撰写、 │
│ 分析) │ │ 排版) │
└────┬───┘ └────┬───┘
│ │
▼ ▼
搜索工具 代码执行
计算器 文件操作
│ │
└─────┬────┘
▼
┌───────────┐
│ 记忆系统 │
│ (短期+长期)│
└───────────┘5.2 定义三个角色Agent
# agents/researcher.py
from .base import BaseAgent, AgentConfig
class ResearcherAgent(BaseAgent):
"""研究员Agent:负责信息搜集、数据分析和事实核查"""
def __init__(self, tool_registry):
config = AgentConfig(
name="研究员小研",
role="资深研究分析师",
system_prompt="""你是一位资深研究分析师,擅长:
1. 通过搜索工具搜集最新的行业信息和数据
2. 对数据进行深入分析和解读
3. 对关键信息进行事实核查
4. 输出结构化的研究要点
请在回答中:
- 注明信息来源
- 区分事实和观点
- 用要点形式呈现研究结果
- 如果信息不足,明确说明""",
tool_names=["web_search", "calculator"],
temperature=0.3,
)
super().__init__(config, tool_registry)# agents/writer.py
from .base import BaseAgent, AgentConfig
class WriterAgent(BaseAgent):
"""写作Agent:负责文章撰写、润色和排版"""
def __init__(self, tool_registry):
config = AgentConfig(
name="作家小文",
role="专业技术作家",
system_prompt="""你是一位专业的技术作家,擅长:
1. 将研究资料转化为结构清晰、可读性强的文章
2. 使用恰当的标题、段落和列表组织内容
3. 确保技术准确性的同时保持语言流畅
4. 根据读者群体调整语言风格
写作要求:
- 使用Markdown格式
- 结构清晰,有明确的标题层级
- 段落简洁,重点突出
- 适当使用代码示例和图表描述
- 字数不少于1000字""",
tool_names=[], # 写作Agent不直接调用外部工具
temperature=0.8,
)
super().__init__(config, tool_registry)# agents/reviewer.py
from .base import BaseAgent, AgentConfig
class ReviewerAgent(BaseAgent):
"""审稿Agent:负责质量检查和改进建议"""
def __init__(self, tool_registry):
config = AgentConfig(
name="审稿人小审",
role="严谨的技术编辑",
system_prompt="""你是一位严谨的技术编辑,负责:
1. 检查文章的技术准确性
2. 评估结构的完整性和逻辑性
3. 提出具体的修改建议
4. 给出最终的评分(1-10分)
审稿维度:
- 技术准确性:信息是否准确可靠
- 结构完整性:是否有清晰的逻辑结构
- 可读性:语言是否流畅易懂
- 实用性:内容是否有实际价值
请用结构化的方式输出审稿意见,最后给出改进后的版本。""",
tool_names=["calculator"],
temperature=0.2,
)
super().__init__(config, tool_registry)5.3 使用LangGraph构建多Agent工作流
LangGraph是LangChain官方推出的Agent编排框架,非常适合构建多Agent协作的状态机工作流。
# workflows/research_flow.py
from __future__ import annotations
from typing import TypedDict, List, Annotated
from langgraph.graph import StateGraph, END
from langchain_core.messages import BaseMessage, HumanMessage
import operator
from agents.researcher import ResearcherAgent
from agents.writer import WriterAgent
from agents.reviewer import ReviewerAgent
from tools import ToolRegistry
# 定义工作流状态
class ResearchState(TypedDict):
"""研究工作流的共享状态"""
topic: str # 研究主题
research_notes: str # 研究笔记
draft: str # 初稿
review_feedback: str # 审稿意见
final_report: str # 最终报告
messages: Annotated[List[BaseMessage], operator.add] # 消息历史
iteration: int # 当前迭代次数
max_iterations: int # 最大迭代次数
class ResearchWorkflow:
"""研究报告生成工作流"""
def __init__(self, tool_registry: ToolRegistry, max_iterations: int = 2):
self.tool_registry = tool_registry
self.max_iterations = max_iterations
# 初始化各Agent
self.researcher = ResearcherAgent(tool_registry)
self.writer = WriterAgent(tool_registry)
self.reviewer = ReviewerAgent(tool_registry)
# 构建工作流图
self.workflow = self._build_graph()
def _build_graph(self) -> StateGraph:
"""构建工作流状态图"""
workflow = StateGraph(ResearchState)
# 添加节点
workflow.add_node("research", self._research_node)
workflow.add_node("write", self._write_node)
workflow.add_node("review", self._review_node)
workflow.add_node("revise", self._revise_node)
# 设置入口
workflow.set_entry_point("research")
# 添加边
workflow.add_edge("research", "write")
workflow.add_edge("write", "review")
# 条件边:根据审稿结果决定下一步
workflow.add_conditional_edges(
"review",
self._should_revise,
{
"revise": "revise",
"end": END,
},
)
workflow.add_edge("revise", "review")
return workflow.compile()
async def _research_node(self, state: ResearchState) -> dict:
"""研究节点:搜集信息"""
topic = state["topic"]
print(f"[研究员] 开始研究主题: {topic}")
prompt = f"""请针对以下主题进行深入研究,输出详细的研究笔记:
主题:{topic}
研究要求:
1. 搜集该领域的最新发展动态
2. 整理关键数据和统计信息
3. 分析未来趋势和挑战
4. 列出3-5个核心观点
请用结构化方式输出研究笔记。"""
research_notes = await self.researcher.run(prompt)
print(f"[研究员] 研究完成,笔记长度: {len(research_notes)} 字")
return {
"research_notes": research_notes,
"iteration": state.get("iteration", 0),
"max_iterations": self.max_iterations,
}
async def _write_node(self, state: ResearchState) -> dict:
"""写作节点:生成初稿"""
research_notes = state["research_notes"]
topic = state["topic"]
print(f"[写作] 开始撰写初稿...")
prompt = f"""请根据以下研究笔记,撰写一篇完整的技术文章:
文章主题:{topic}
研究笔记:
{research_notes}
文章要求:
1. 标题醒目,有吸引力
2. 包含引言、正文(3-5个小节)、结论
3. 使用Markdown格式
4. 语言流畅,逻辑清晰
5. 不少于1500字
6. 适当加入代码示例(如果适用)"""
draft = await self.writer.run(prompt)
print(f"[写作] 初稿完成,长度: {len(draft)} 字")
return {"draft": draft}
async def _review_node(self, state: ResearchState) -> dict:
"""审稿节点:质量检查"""
draft = state["draft"]
iteration = state.get("iteration", 0) + 1
print(f"[审稿] 第 {iteration} 轮审稿...")
prompt = f"""请对以下文章进行审稿:
---
{draft}
---
请从技术准确性、结构完整性、可读性、实用性四个维度进行评估,
给出具体的修改建议,并输出一个总体评分(1-10分)。
最后提供你认为可以直接使用的修改后版本。"""
review_feedback = await self.reviewer.run(prompt)
print(f"[审稿] 审稿完成")
return {
"review_feedback": review_feedback,
"iteration": iteration,
}
async def _revise_node(self, state: ResearchState) -> dict:
"""修改节点:根据审稿意见修改文章"""
draft = state["draft"]
feedback = state["review_feedback"]
print(f"[修改] 根据审稿意见修改文章...")
prompt = f"""请根据以下审稿意见修改文章:
原始文章:
{draft}
审稿意见:
{feedback}
请输出修改后的完整文章,确保采纳了审稿意见中的合理建议。"""
revised_draft = await self.writer.run(prompt)
print(f"[修改] 修改完成,长度: {len(revised_draft)} 字")
return {"draft": revised_draft}
def _should_revise(self, state: ResearchState) -> str:
"""判断是否需要继续修改"""
iteration = state.get("iteration", 0)
max_iterations = state.get("max_iterations", self.max_iterations)
if iteration >= max_iterations:
print(f"[工作流] 已达到最大迭代次数 ({max_iterations}),结束")
# 保存最终报告
state["final_report"] = state["draft"]
return "end"
# 简单的评分判断逻辑可以在这里扩展
print(f"[工作流] 继续修改 (当前迭代: {iteration}/{max_iterations})")
return "revise"
async def run(self, topic: str) -> str:
"""运行完整的研究工作流"""
initial_state = {
"topic": topic,
"research_notes": "",
"draft": "",
"review_feedback": "",
"final_report": "",
"messages": [],
"iteration": 0,
"max_iterations": self.max_iterations,
}
result = await self.workflow.ainvoke(initial_state)
return result.get("final_report", result.get("draft", ""))5.4 运行多Agent系统
# main_multi_agent.py
import asyncio
from dotenv import load_dotenv
load_dotenv()
from tools import tool_registry
from tools.calculator import calculator
from tools.search import web_search
from workflows.research_flow import ResearchWorkflow
# 注册工具
tool_registry.register(calculator)
tool_registry.register(web_search) # 下一节实现
async def main():
# 创建工作流
workflow = ResearchWorkflow(tool_registry, max_iterations=2)
# 运行
topic = "2026年AI Agent技术发展趋势与企业应用前景"
print(f"开始生成研究报告: {topic}
")
report = await workflow.run(topic)
print("
" + "=" * 60)
print("最终报告已生成!")
print("=" * 60)
print(report[:500] + "...")
# 保存到文件
with open("output/report.md", "w", encoding="utf-8") as f:
f.write(report)
print(f"
完整报告已保存到 output/report.md")
if __name__ == "__main__":
asyncio.run(main())六、工具调用:基于MCP协议的工具系统
6.1 MCP协议简介
MCP(Model Context Protocol)是2025年兴起的模型上下文协议,它定义了一套标准的工具注册、发现和调用规范,使得不同厂商的AI模型可以用统一的方式调用外部工具。
MCP的核心优势:
6.2 实现网络搜索工具
# tools/search.py
from langchain_core.tools import tool
import requests
from typing import Optional
import os
from dotenv import load_dotenv
load_dotenv()
@tool
def web_search(query: str, num_results: int = 5) -> str:
"""
使用搜索引擎搜索网络信息。
Args:
query: 搜索关键词
num_results: 返回结果数量,最多10条
Returns:
搜索结果的标题、摘要和链接
"""
# 这里使用Tavily搜索API作为示例,也可以替换为其他搜索服务
# Tavily是专为AI Agent设计的搜索API
api_key = os.getenv("TAVILY_API_KEY", "")
if not api_key:
# 如果没有配置API key,返回模拟数据(开发调试用)
return _mock_search_results(query, num_results)
try:
response = requests.post(
"https://api.tavily.com/search",
json={
"api_key": api_key,
"query": query,
"max_results": min(num_results, 10),
"search_depth": "basic",
},
timeout=15,
)
response.raise_for_status()
data = response.json()
results = []
for i, item in enumerate(data.get("results", []), 1):
results.append(
f"{i}. [{item.get('title', '无标题')}]({item.get('url', '')})
"
f" 摘要: {item.get('content', '无摘要')[:200]}..."
)
return "
".join(results) if results else "未找到相关结果"
except Exception as e:
return f"搜索失败: {str(e)}"
def _mock_search_results(query: str, num_results: int) -> str:
"""生成模拟搜索结果(用于开发调试)"""
mock_data = [
{
"title": f"{query} - 2026年最新发展报告",
"url": f"https://example.com/report/{query}",
"content": f"关于{query}的最新行业研究报告显示,2026年该领域呈现出快速发展的趋势,多项关键技术取得突破性进展...",
},
{
"title": f"{query}:企业落地实践指南",
"url": f"https://example.com/practice/{query}",
"content": f"本文总结了多家头部企业在{query}方面的落地经验,包括技术选型、架构设计、团队建设等方面的最佳实践...",
},
{
"title": f"深度解析:{query}的核心技术原理",
"url": f"https://example.com/tech/{query}",
"content": f"从技术原理角度深入分析{query},涵盖核心算法、系统架构、性能优化等关键技术点...",
},
]
results = []
for i, item in enumerate(mock_data[:num_results], 1):
results.append(
f"{i}. [{item['title']}]({item['url']})
"
f" 摘要: {item['content']}"
)
return "
".join(results) + "
[注:以上为模拟数据,配置TAVILY_API_KEY后可获取真实搜索结果]"6.3 实现代码沙箱执行工具
安全是Agent工具调用的重中之重。代码执行工具必须运行在沙箱环境中:
# tools/code_executor.py
from langchain_core.tools import tool
import subprocess
import tempfile
import os
from typing import Optional
@tool
def python_executor(code: str, timeout: int = 30) -> str:
"""
在安全沙箱中执行Python代码并返回结果。
WARNING: 代码会在隔离环境中执行,但仍需谨慎使用。
禁止执行网络请求、文件系统操作(除临时目录外)等危险操作。
Args:
code: 要执行的Python代码
timeout: 超时时间(秒)
Returns:
代码执行的输出结果(stdout + stderr)
"""
# 创建临时目录作为沙箱工作目录
with tempfile.TemporaryDirectory() as sandbox_dir:
code_file = os.path.join(sandbox_dir, "script.py")
# 写入代码文件
with open(code_file, "w", encoding="utf-8") as f:
f.write(code)
try:
# 使用子进程执行,限制资源
result = subprocess.run(
["python", "-c", code],
capture_output=True,
text=True,
timeout=timeout,
# 限制环境变量
env={
"PATH": "/usr/bin:/bin",
"HOME": sandbox_dir,
"TMPDIR": sandbox_dir,
},
cwd=sandbox_dir,
)
output = ""
if result.stdout:
output += f"标准输出:
{result.stdout}
"
if result.stderr:
output += f"错误输出:
{result.stderr}
"
if not output:
output = "代码执行完成,无输出。"
output += f"
退出码: {result.returncode}"
return output
except subprocess.TimeoutExpired:
return f"代码执行超时(超过{timeout}秒),已终止。"
except Exception as e:
return f"执行错误: {str(e)}"6.4 MCP工具服务器
接下来我们实现一个简单的MCP兼容工具服务器,让外部Agent也能调用我们的工具:
# tools/mcp_server.py
from __future__ import annotations
from fastapi import FastAPI, HTTPException
from pydantic import BaseModel
from typing import Dict, Any, List
from . import tool_registry
app = FastAPI(title="MCP Tool Server", version="1.0.0")
class ToolCallRequest(BaseModel):
tool_name: str
arguments: Dict[str, Any]
class ToolInfo(BaseModel):
name: str
description: str
parameters: Dict[str, Any]
@app.get("/tools", response_model=List[ToolInfo])
async def list_tools():
"""列出所有可用工具(MCP标准端点)"""
tools_info = []
for tool in tool_registry.list_tools():
tools_info.append(
ToolInfo(
name=tool.name,
description=tool.description,
parameters=tool.args_schema.model_json_schema()
if hasattr(tool, "args_schema") and tool.args_schema
else {"type": "object", "properties": {}},
)
)
return tools_info
@app.post("/tools/call")
async def call_tool(request: ToolCallRequest):
"""调用工具(MCP标准端点)"""
try:
result = await tool_registry.execute_tool(
request.tool_name, request.arguments
)
return {"result": result, "status": "success"}
except ValueError as e:
raise HTTPException(status_code=404, detail=str(e))
except Exception as e:
raise HTTPException(status_code=500, detail=f"Tool execution failed: {str(e)}")
@app.get("/health")
async def health_check():
"""健康检查"""
return {"status": "healthy", "tools_count": len(tool_registry.list_tools())}启动MCP工具服务器:
# 启动工具服务器
uvicorn tools.mcp_server:app --host 0.0.0.0 --port 8000七、记忆系统:让Agent拥有长期记忆
7.1 记忆系统架构
一个完善的记忆系统应该包含三个层次:
┌─────────────────────────────────────────┐
│ 工作记忆 (Working Memory) │
│ 当前对话上下文,滑动窗口管理 │
│ 存储:内存 / 时效:即时 │
├─────────────────────────────────────────┤
│ 短期记忆 (Short-term Memory) │
│ 会话历史,支持摘要压缩 │
│ 存储:内存 / 时效:会话级 │
├─────────────────────────────────────────┤
│ 长期记忆 (Long-term Memory) │
│ 知识库 + 向量检索 + 知识图谱 │
│ 存储:向量DB / 时效:永久 │
└─────────────────────────────────────────┘7.2 长期记忆实现
# memory/long_term.py
from __future__ import annotations
from typing import List, Dict, Any, Optional
from dataclasses import dataclass
from datetime import datetime
import chromadb
from chromadb.config import Settings
from sentence_transformers import SentenceTransformer
import uuid
@dataclass
class MemoryItem:
"""记忆条目"""
id: str
content: str
metadata: Dict[str, Any]
timestamp: float
importance: float = 0.5 # 重要性评分 0-1
class LongTermMemory:
"""长期记忆:基于向量数据库的语义记忆系统"""
def __init__(
self,
collection_name: str = "agent_memory",
persist_directory: str = "./data/chroma",
embedding_model: str = "BAAI/bge-small-zh-v1.5",
):
self.collection_name = collection_name
# 初始化向量数据库
self.client = chromadb.PersistentClient(
path=persist_directory,
settings=Settings(anonymized_telemetry=False),
)
self.collection = self.client.get_or_create_collection(
name=collection_name,
metadata={"description": "Agent long-term memory"},
)
# 初始化嵌入模型
self.embedding_model = SentenceTransformer(embedding_model)
def _embed(self, text: str) -> List[float]:
"""生成文本的向量嵌入"""
return self.embedding_model.encode(text).tolist()
def add(
self,
content: str,
metadata: Optional[Dict[str, Any]] = None,
importance: float = 0.5,
) -> str:
"""添加一条记忆"""
memory_id = str(uuid.uuid4())
timestamp = datetime.now().timestamp()
meta = {
"timestamp": timestamp,
"importance": importance,
"content_length": len(content),
}
if metadata:
meta.update(metadata)
# 生成嵌入向量
embedding = self._embed(content)
# 存入向量数据库
self.collection.add(
ids=[memory_id],
embeddings=[embedding],
documents=[content],
metadatas=[meta],
)
return memory_id
def search(
self,
query: str,
top_k: int = 5,
min_score: float = 0.3,
filter_metadata: Optional[Dict[str, Any]] = None,
) -> List[MemoryItem]:
"""语义搜索记忆"""
query_embedding = self._embed(query)
results = self.collection.query(
query_embeddings=[query_embedding],
n_results=top_k,
where=filter_metadata,
)
memory_items = []
if results["ids"] and results["ids"][0]:
for i, (doc_id, doc, meta, dist) in enumerate(
zip(
results["ids"][0],
results["documents"][0],
results["metadatas"][0],
results["distances"][0],
)
):
# Chroma返回的是距离,转换为相似度分数
similarity = 1.0 - min(dist, 1.0)
if similarity >= min_score:
memory_items.append(
MemoryItem(
id=doc_id,
content=doc,
metadata=meta,
timestamp=meta.get("timestamp", 0),
importance=meta.get("importance", 0.5),
)
)
return memory_items
def update_importance(self, memory_id: str, importance: float):
"""更新记忆的重要性评分"""
# 获取现有metadata
result = self.collection.get(ids=[memory_id])
if result["metadatas"] and result["metadatas"][0]:
meta = result["metadatas"][0]
meta["importance"] = importance
self.collection.update(ids=[memory_id], metadatas=[meta])
def delete(self, memory_id: str):
"""删除一条记忆"""
self.collection.delete(ids=[memory_id])
def count(self) -> int:
"""获取记忆总数"""
return self.collection.count()
def forget_old_low_importance(self, days: int = 30, min_importance: float = 0.3):
"""遗忘机制:删除旧的低重要性记忆"""
import time
cutoff = time.time() - days * 86400
# 注意:Chroma的where查询语法
# 这里用一个简化的实现,实际使用时需要批量处理
results = self.collection.get(
where={
"$and": [
{"timestamp": {"$lt": cutoff}},
{"importance": {"$lt": min_importance}},
]
}
)
if results["ids"]:
self.collection.delete(ids=results["ids"])
return len(results["ids"])
return 07.3 记忆增强(RAG)集成到Agent
将长期记忆集成到BaseAgent中,实现记忆增强的推理:
# agents/base.py(更新版本,添加RAG支持)
# ...(保持原有代码,添加以下方法)
class BaseAgent:
# ...(保持原有代码)
def set_long_term_memory(self, ltm: LongTermMemory, memory_query_top_k: int = 5):
"""设置长期记忆"""
self.long_term_memory = ltm
self.memory_query_top_k = memory_query_top_k
def _retrieve_relevant_memory(self, query: str) -> str:
"""检索相关的长期记忆"""
if not hasattr(self, 'long_term_memory') or self.long_term_memory is None:
return ""
memories = self.long_term_memory.search(query, top_k=self.memory_query_top_k)
if not memories:
return ""
memory_text = "
".join(
[f"- [{i+1}] {m.content}" for i, m in enumerate(memories)]
)
return f"
相关记忆:
{memory_text}"
async def run(self, user_input: str) -> str:
"""运行Agent(增强版:带记忆检索)"""
# 检索相关记忆
relevant_memory = self._retrieve_relevant_memory(user_input)
# 构建消息列表
messages = [self._build_system_prompt()]
messages.extend(self.memory.get_history())
# 将记忆增强内容加入用户输入
enhanced_input = user_input + relevant_memory
messages.append(HumanMessage(content=enhanced_input))
# ...(后续逻辑保持不变)7.4 记忆写入:自动总结与存储
Agent在对话结束后,自动将重要内容存入长期记忆:
# memory/memory_manager.py
from __future__ import annotations
from typing import List
from langchain_core.messages import BaseMessage
from langchain_openai import ChatOpenAI
from .long_term import LongTermMemory
class MemoryManager:
"""记忆管理器:负责对话记忆的总结和长期存储"""
def __init__(self, llm: ChatOpenAI, ltm: LongTermMemory):
self.llm = llm
self.ltm = ltm
async def consolidate_and_store(self, conversation_history: List[BaseMessage]):
"""
将对话历史总结为关键知识点并存入长期记忆
"""
if len(conversation_history) < 4: # 对话太短,不需要总结
return
# 构建总结提示
conversation_text = "
".join(
[f"{msg.type}: {msg.content}" for msg in conversation_history]
)
prompt = f"""请从以下对话中提取3-5个关键知识点或重要事实,
每个知识点用一句话概括,并标注重要性评分(0.1-1.0)。
格式要求(每行一个):
[重要性] 知识点内容
对话内容:
{conversation_text}
"""
response = await self.llm.ainvoke(prompt)
# 解析结果并存入长期记忆
lines = response.content.strip().split("
")
for line in lines:
line = line.strip()
if not line:
continue
# 尝试解析 [重要性] 内容 格式
if line.startswith("[") and "]" in line:
try:
importance_str = line[1 : line.index("]")]
importance = float(importance_str)
content = line[line.index("]") + 1 :].strip()
self.ltm.add(
content=content,
metadata={"source": "conversation_summary"},
importance=importance,
)
except (ValueError, IndexError):
pass八、部署优化:从本地开发到生产环境
8.1 API服务化
将Agent系统封装为FastAPI服务,支持HTTP调用:
# deployment/server.py
from fastapi import FastAPI, HTTPException
from pydantic import BaseModel
from typing import Optional, List, Dict, Any
import asyncio
from dotenv import load_dotenv
load_dotenv()
from tools import tool_registry
from tools.calculator import calculator
from tools.search import web_search
from tools.code_executor import python_executor
from memory.long_term import LongTermMemory
from workflows.research_flow import ResearchWorkflow
app = FastAPI(
title="Multi-Agent System API",
description="2026 AI Agent多智能体协作系统API",
version="1.0.0",
)
# 全局初始化
_workflow: Optional[ResearchWorkflow] = None
_ltm: Optional[LongTermMemory] = None
def get_workflow() -> ResearchWorkflow:
"""获取或初始化工作流"""
global _workflow
if _workflow is None:
# 注册工具
tool_registry.register(calculator)
tool_registry.register(web_search)
tool_registry.register(python_executor)
_workflow = ResearchWorkflow(tool_registry, max_iterations=2)
return _workflow
def get_ltm() -> LongTermMemory:
"""获取或初始化长期记忆"""
global _ltm
if _ltm is None:
_ltm = LongTermMemory(
collection_name="agent_memory",
persist_directory="./data/chroma",
)
return _ltm
# ---------- 请求/响应模型 ----------
class ResearchRequest(BaseModel):
topic: str
max_iterations: int = 2
class ResearchResponse(BaseModel):
report: str
status: str
execution_time: float
class ChatRequest(BaseModel):
agent: str # researcher / writer / reviewer
message: str
class ChatResponse(BaseModel):
reply: str
class MemoryAddRequest(BaseModel):
content: str
metadata: Optional[Dict[str, Any]] = None
importance: float = 0.5
class MemorySearchRequest(BaseModel):
query: str
top_k: int = 5
# ---------- 路由 ----------
@app.post("/api/research/generate", response_model=ResearchResponse)
async def generate_research(request: ResearchRequest):
"""生成研究报告"""
import time
start_time = time.time()
try:
workflow = get_workflow()
workflow.max_iterations = request.max_iterations
report = await workflow.run(request.topic)
execution_time = time.time() - start_time
return ResearchResponse(
report=report,
status="success",
execution_time=execution_time,
)
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
@app.post("/api/chat", response_model=ChatResponse)
async def chat_with_agent(request: ChatRequest):
"""与单个Agent对话"""
workflow = get_workflow()
agent_map = {
"researcher": workflow.researcher,
"writer": workflow.writer,
"reviewer": workflow.reviewer,
}
agent = agent_map.get(request.agent)
if not agent:
raise HTTPException(
status_code=400,
detail=f"Unknown agent: {request.agent}. Available: {list(agent_map.keys())}",
)
reply = await agent.run(request.message)
return ChatResponse(reply=reply)
@app.post("/api/memory/add")
async def add_memory(request: MemoryAddRequest):
"""添加长期记忆"""
ltm = get_ltm()
memory_id = ltm.add(
content=request.content,
metadata=request.metadata,
importance=request.importance,
)
return {"id": memory_id, "status": "success"}
@app.post("/api/memory/search")
async def search_memory(request: MemorySearchRequest):
"""搜索长期记忆"""
ltm = get_ltm()
results = ltm.search(request.query, top_k=request.top_k)
return {
"results": [
{
"id": m.id,
"content": m.content,
"metadata": m.metadata,
"timestamp": m.timestamp,
"importance": m.importance,
}
for m in results
]
}
@app.get("/api/health")
async def health_check():
"""健康检查"""
return {
"status": "healthy",
"version": "1.0.0",
"tools_count": len(tool_registry.list_tools()),
}
@app.get("/api/tools")
async def list_tools():
"""列出所有可用工具"""
return {
"tools": [
{"name": t.name, "description": t.description}
for t in tool_registry.list_tools()
]
}
if __name__ == "__main__":
import uvicorn
uvicorn.run(app, host="0.0.0.0", port=8080)8.2 Docker容器化部署
# deployment/Dockerfile
FROM python:3.11-slim
WORKDIR /app
# 安装系统依赖
RUN apt-get update && apt-get install -y --no-install-recommends \
build-essential \
&& rm -rf /var/lib/apt/lists/*
# 安装Python依赖
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
# 复制应用代码
COPY . .
# 创建数据目录
RUN mkdir -p /app/data/chroma /app/data/logs /app/output
# 暴露端口
EXPOSE 8080
# 健康检查
HEALTHCHECK --interval=30s --timeout=10s --start-period=60s --retries=3 \
CMD curl -f http://localhost:8080/api/health || exit 1
# 启动服务
CMD ["uvicorn", "deployment.server:app", "--host", "0.0.0.0", "--port", "8080"]# deployment/docker-compose.yml
version: '3.8'
services:
agent-server:
build:
context: ..
dockerfile: deployment/Dockerfile
ports:
- "8080:8080"
environment:
- LLM_API_KEY=${LLM_API_KEY}
- LLM_BASE_URL=${LLM_BASE_URL}
- LLM_MODEL=${LLM_MODEL}
volumes:
- ../data:/app/data
- ../output:/app/output
restart: unless-stopped
deploy:
resources:
limits:
memory: 2G
cpus: '1.0'8.3 性能优化策略
#### 8.3.1 提示词缓存
对于重复的系统提示词和常见查询,使用缓存减少LLM调用:
# utils/cache.py
from functools import lru_cache
from typing import Dict, Any
import hashlib
import json
import time
class LLMCache:
"""LLM响应缓存"""
def __init__(self, max_size: int = 1000, ttl: int = 3600):
self.max_size = max_size
self.ttl = ttl # 缓存有效期(秒)
self._cache: Dict[str, tuple] = {} # key: (value, timestamp)
def _make_key(self, prompt: str, **kwargs) -> str:
"""生成缓存键"""
key_data = {"prompt": prompt, **kwargs}
key_str = json.dumps(key_data, sort_keys=True)
return hashlib.md5(key_str.encode()).hexdigest()
def get(self, prompt: str, **kwargs) -> Any:
"""获取缓存"""
key = self._make_key(prompt, **kwargs)
if key in self._cache:
value, timestamp = self._cache[key]
if time.time() - timestamp < self.ttl:
return value
else:
del self._cache[key] # 过期删除
return None
def set(self, prompt: str, value: Any, **kwargs):
"""设置缓存"""
key = self._make_key(prompt, **kwargs)
self._cache[key] = (value, time.time())
# LRU淘汰:超出大小时删除最旧的
if len(self._cache) > self.max_size:
oldest_key = min(self._cache, key=lambda k: self._cache[k][1])
del self._cache[oldest_key]#### 8.3.2 并发与异步优化
多Agent系统天然适合并发执行。以下是一个并发研究的例子:
# workflows/concurrent_research.py
import asyncio
from typing import List
class ConcurrentResearchFlow:
"""并发研究工作流:多个研究员同时研究不同子话题"""
async def parallel_research(self, sub_topics: List[str]) -> List[str]:
"""并行研究多个子话题"""
tasks = []
for topic in sub_topics:
# 为每个子话题创建独立的研究员Agent
task = self._research_single(topic)
tasks.append(task)
# 并发执行
results = await asyncio.gather(*tasks)
return results
async def _research_single(self, topic: str) -> str:
"""研究单个子话题"""
from agents.researcher import ResearcherAgent
researcher = ResearcherAgent(tool_registry)
result = await researcher.run(f"请研究:{topic}")
return result#### 8.3.3 流式响应
对于长对话和长文本生成,使用流式响应提升用户体验:
# deployment/streaming.py
from fastapi.responses import StreamingResponse
async def stream_chat(agent, message: str):
"""流式聊天响应"""
# 构建消息
messages = [agent._build_system_prompt()]
messages.extend(agent.memory.get_history())
messages.append({"role": "user", "content": message})
async def generate():
full_response = ""
async for chunk in agent.llm.astream(messages):
if chunk.content:
full_response += chunk.content
yield f"data: {chunk.content}
"
yield "data: [DONE]
"
# 保存到记忆(流式完成后)
agent.memory.add_message({"role": "user", "content": message})
agent.memory.add_message({"role": "assistant", "content": full_response})
return StreamingResponse(generate(), media_type="text/event-stream")8.4 可观测性与监控
生产环境的Agent系统必须具备完善的可观测性:
# utils/observability.py
from __future__ import annotations
import time
import logging
from functools import wraps
from typing import Callable, Any
from dataclasses import dataclass, field
from datetime import datetime
# 配置日志
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s',
handlers=[
logging.FileHandler('./data/logs/agent.log'),
logging.StreamHandler(),
]
)
logger = logging.getLogger("agent_system")
@dataclass
class AgentMetrics:
"""Agent性能指标"""
total_requests: int = 0
total_tokens: int = 0
total_latency: float = 0.0
tool_calls: int = 0
errors: int = 0
def record_request(self, tokens: int, latency: float, tool_calls: int = 0, error: bool = False):
self.total_requests += 1
self.total_tokens += tokens
self.total_latency += latency
self.tool_calls += tool_calls
if error:
self.errors += 1
@property
def avg_latency(self) -> float:
return self.total_latency / self.total_requests if self.total_requests else 0.0
@property
def error_rate(self) -> float:
return self.errors / self.total_requests if self.total_requests else 0.0
metrics = AgentMetrics()
def trace_agent(func: Callable) -> Callable:
"""Agent调用追踪装饰器"""
@wraps(func)
async def wrapper(*args, **kwargs):
agent_name = args[0].config.name if args else "unknown"
start_time = time.time()
tool_call_count = 0
error = False
logger.info(f"[{agent_name}] 开始执行")
try:
result = await func(*args, **kwargs)
return result
except Exception as e:
error = True
logger.error(f"[{agent_name}] 执行错误: {str(e)}", exc_info=True)
raise
finally:
latency = time.time() - start_time
metrics.record_request(
tokens=0, # 实际使用时从LLM响应中获取token数
latency=latency,
tool_calls=tool_call_count,
error=error,
)
logger.info(
f"[{agent_name}] 执行完成 - 耗时: {latency:.2f}s, "
f"错误: {error}"
)
return wrapper九、总结与展望
9.1 本文回顾
本文从零开始,构建了一套完整的多Agent协作系统,涵盖:
9.2 下一步扩展方向
9.3 写在最后
2026年是AI Agent从概念走向工程化落地的关键一年。作为开发者,我们不仅要掌握框架的使用,更要理解Agent系统的设计原理:思考-行动-观察的循环逻辑、记忆的分层管理、工具的安全调用、多Agent的协作模式。
希望本文能帮助你快速上手AI Agent开发,在这个充满机遇的技术浪潮中找到自己的位置。
参考资源:
本文为原创技术文章,欢迎转载,请注明出处。
💬 评论区 (0)
暂无评论,快来抢沙发吧!