1. 生产级Agent系统架构解析
在构建一个生产级Agent系统时,核心不在于代码行数的多少,而在于架构的完整性和可扩展性。一个典型的500行左右的MVP(最小可行产品)Agent框架应该包含以下核心组件:
- Planner:负责任务分解和规划
- Tool Router:工具选择和路由
- RAG Memory:检索增强生成的内存系统
- Reflection:结果评估和反馈机制
- Agent Runtime:核心执行循环
- Environment:与外部系统的交互接口
- Task Graph:任务依赖关系管理
提示:生产环境中的Agent系统通常会采用模块化设计,每个功能组件独立成文件,便于维护和扩展。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 项目目录结构设计
一个标准的Agent系统项目结构如下:
code复制agent_system/
│
├── agent_runtime.py # 核心执行逻辑
├── planner.py # 任务规划模块
├── memory.py # 记忆系统
├── tools.py # 工具集合
├── reflection.py # 结果评估
├── rag.py # 检索增强模块
├── environment.py # 环境交互
├── workflow.py # 工作流引擎
├── llm.py # LLM接口封装
└── main.py # 入口文件
这种结构设计有以下几个优点:
- 模块职责清晰,便于团队协作开发
- 单个文件代码量控制在合理范围(50-100行)
- 易于单元测试和功能扩展
- 符合Python项目的最佳实践
3. LLM接口封装实现
llm.py文件负责统一管理LLM调用,基础实现如下:
python复制import openai
from tenacity import retry, stop_after_attempt, wait_exponential
class LLMClient:
def __init__(self, model="gpt-4", max_retries=3):
self.model = model
self.max_retries = max_retries
@retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=4, max=10))
def chat(self, prompt, temperature=0.7):
try:
response = openai.ChatCompletion.create(
model=self.model,
messages=[{"role":"user","content":prompt}],
temperature=temperature
)
return response.choices[0].message["content"]
except Exception as e:
print(f"LLM调用失败: {str(e)}")
raise
生产环境需要考虑的增强功能:
- 重试机制:使用tenacity库实现指数退避重试
- 流式响应:支持大文本的流式返回
- 成本跟踪:记录每次调用的token消耗
- 缓存机制:对相同prompt的结果进行缓存
- 限流控制:防止API调用频率过高
4. 任务规划系统实现
planner.py负责将用户目标分解为可执行的任务步骤:
python复制import json
from llm import LLMClient
class Planner:
def __init__(self):
self.llm = LLMClient()
self.system_prompt = """
你是一个专业的任务规划AI。请将用户目标分解为可执行的步骤。
输出要求:
1. 每个步骤应该是具体的、可执行的动作
2. 步骤数量控制在3-5个
3. 使用JSON数组格式返回
"""
def plan(self, goal, context=None):
prompt = f"{self.system_prompt}\n目标:{goal}\n"
if context:
prompt += f"上下文:\n{context}\n"
response = self.llm.chat(prompt, temperature=0.3)
try:
steps = json.loads(response)
if not isinstance(steps, list):
steps = [goal]
except json.JSONDecodeError:
steps = [goal]
return steps
典型输出示例:
json复制[
"搜索东京的人口数据",
"将人口数据乘以2",
"格式化输出最终结果"
]
5. 工具系统设计与实现
tools.py包含Agent可用的工具集和路由逻辑:
python复制from typing import Dict, Callable
from llm import LLMClient
class ToolSystem:
def __init__(self):
self.llm = LLMClient()
self.tools: Dict[str, Callable] = {
"search": self.search,
"calculator": self.calculator,
"shell": self.run_shell
}
def search(self, query: str) -> str:
"""模拟搜索引擎查询"""
# 实际项目中可接入Google Search API等
return f"关于'{query}'的搜索结果..."
def calculator(self, expr: str) -> str:
"""数学表达式计算"""
try:
return str(eval(expr))
except:
return "计算错误: 无效的表达式"
def run_shell(self, command: str) -> str:
"""执行Shell命令"""
import subprocess
try:
result = subprocess.run(
command,
shell=True,
check=True,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
text=True
)
return result.stdout
except subprocess.CalledProcessError as e:
return f"命令执行失败: {e.stderr}"
def select_tool(self, task: str) -> str:
"""选择最适合当前任务的工具"""
prompt = f"""
请为以下任务选择最合适的工具:
任务: {task}
可用工具:
{", ".join(self.tools.keys())}
只需返回工具名称,如果没有合适工具则返回NONE
"""
tool_name = self.llm.chat(prompt).strip()
return tool_name if tool_name in self.tools else None
def execute(self, tool_name: str, input_data: str) -> str:
"""执行指定工具"""
if tool_name not in self.tools:
return f"错误: 未知工具 {tool_name}"
return self.tools[tool_name](input_data)
6. 记忆系统实现
memory.py实现Agent的记忆功能:
python复制from typing import List
from sentence_transformers import SentenceTransformer
import numpy as np
class MemorySystem:
def __init__(self, max_history=20):
self.history: List[str] = []
self.max_history = max_history
self.embedder = SentenceTransformer('paraphrase-multilingual-MiniLM-L12-v2')
self.embeddings = []
def add(self, item: str):
"""添加记忆项"""
self.history.append(item)
if len(self.history) > self.max_history:
self.history.pop(0)
# 生成嵌入向量
embedding = self.embedder.encode(item)
self.embeddings.append(embedding)
if len(self.embeddings) > self.max_history:
self.embeddings.pop(0)
def get_context(self, k=5) -> str:
"""获取最近的k条记忆"""
return "\n".join(self.history[-k:])
def semantic_search(self, query: str, k=3) -> List[str]:
"""语义搜索相关记忆"""
query_embed = self.embedder.encode(query)
similarities = [
np.dot(query_embed, embed)
for embed in self.embeddings
]
indices = np.argsort(similarities)[-k:]
return [self.history[i] for i in indices[::-1]]
7. RAG检索系统实现
rag.py实现检索增强生成功能:
python复制from typing import List
import chromadb
from sentence_transformers import SentenceTransformer
class RAGSystem:
def __init__(self):
self.embedder = SentenceTransformer('paraphrase-multilingual-MiniLM-L12-v2')
self.client = chromadb.Client()
self.collection = self.client.create_collection("knowledge")
def add_documents(self, documents: List[str]):
"""添加文档到知识库"""
embeddings = self.embedder.encode(documents)
self.collection.add(
embeddings=embeddings.tolist(),
documents=documents,
ids=[str(i) for i in range(len(documents))]
)
def retrieve(self, query: str, k=3) -> List[str]:
"""检索相关文档"""
query_embed = self.embedder.encode(query).tolist()
results = self.collection.query(
query_embeddings=[query_embed],
n_results=k
)
return results['documents'][0]
8. 反思系统实现
reflection.py评估任务执行结果:
python复制from llm import LLMClient
class ReflectionSystem:
def __init__(self):
self.llm = LLMClient()
self.evaluation_prompt = """
请评估以下任务执行结果:
任务: {task}
结果: {result}
评估要求:
1. 结果是否准确完成了任务?
2. 结果格式是否符合要求?
3. 是否存在明显错误?
请用JSON格式返回评估结果,包含以下字段:
- "is_correct": 布尔值,表示结果是否正确
- "feedback": 字符串,提供改进建议
"""
def evaluate(self, task: str, result: str) -> dict:
prompt = self.evaluation_prompt.format(
task=task, result=result
)
response = self.llm.chat(prompt, temperature=0.2)
try:
evaluation = json.loads(response)
return {
"is_correct": evaluation.get("is_correct", False),
"feedback": evaluation.get("feedback", "")
}
except:
return {"is_correct": False, "feedback": "评估失败"}
9. 环境交互系统
environment.py实现与外部环境的交互:
python复制import subprocess
from typing import Optional
class Environment:
def execute_command(self, command: str) -> str:
"""执行系统命令"""
try:
result = subprocess.run(
command,
shell=True,
check=True,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
text=True,
timeout=30
)
return result.stdout
except subprocess.TimeoutExpired:
return "错误: 命令执行超时"
except subprocess.CalledProcessError as e:
return f"错误: {e.stderr}"
def read_file(self, filepath: str) -> Optional[str]:
"""读取文件内容"""
try:
with open(filepath, 'r', encoding='utf-8') as f:
return f.read()
except Exception as e:
return f"文件读取失败: {str(e)}"
def write_file(self, filepath: str, content: str) -> bool:
"""写入文件"""
try:
with open(filepath, 'w', encoding='utf-8') as f:
f.write(content)
return True
except Exception as e:
print(f"文件写入失败: {str(e)}")
return False
10. 工作流引擎实现
workflow.py实现任务的有向无环图(DAG)执行:
python复制from typing import List, Dict, Callable, Optional
import networkx as nx
class WorkflowEngine:
def __init__(self):
self.graph = nx.DiGraph()
self.task_handlers: Dict[str, Callable] = {}
def add_task(self, task_id: str, handler: Callable,
dependencies: Optional[List[str]] = None):
"""添加任务到工作流"""
self.graph.add_node(task_id)
self.task_handlers[task_id] = handler
if dependencies:
for dep in dependencies:
self.graph.add_edge(dep, task_id)
def execute(self, initial_data: dict = None) -> dict:
"""执行工作流"""
if not initial_data:
initial_data = {}
results = {}
for node in nx.topological_sort(self.graph):
# 准备任务输入
inputs = {
**initial_data,
**{
dep: results[dep]
for dep in self.graph.predecessors(node)
}
}
# 执行任务
try:
results[node] = self.task_handlers[node](**inputs)
except Exception as e:
print(f"任务{node}执行失败: {str(e)}")
results[node] = None
return results
11. Agent核心运行时实现
agent_runtime.py实现Agent的核心执行循环:
python复制from typing import Optional, Dict, Any
from planner import Planner
from tools import ToolSystem
from reflection import ReflectionSystem
from memory import MemorySystem
from rag import RAGSystem
from workflow import WorkflowEngine
class AgentRuntime:
def __init__(self):
self.planner = Planner()
self.tools = ToolSystem()
self.reflection = ReflectionSystem()
self.memory = MemorySystem()
self.rag = RAGSystem()
self.workflow = WorkflowEngine()
# 初始化工作流
self._setup_workflow()
def _setup_workflow(self):
"""设置默认工作流"""
self.workflow.add_task("plan", self._plan_task)
self.workflow.add_task("execute", self._execute_task)
self.workflow.add_task("evaluate", self._evaluate_task)
# 设置任务依赖
self.workflow.graph.add_edge("plan", "execute")
self.workflow.graph.add_edge("execute", "evaluate")
def _plan_task(self, goal: str) -> List[str]:
"""规划任务步骤"""
context = self.memory.get_context()
docs = self.rag.retrieve(goal)
full_context = f"{context}\n{docs}"
return self.planner.plan(goal, full_context)
def _execute_task(self, tasks: List[str]) -> Dict[str, str]:
"""执行任务列表"""
results = {}
for task in tasks:
tool_name = self.tools.select_tool(task)
if tool_name:
result = self.tools.execute(tool_name, task)
else:
result = self.tools.llm.chat(task)
results[task] = result
self.memory.add(f"任务: {task}\n结果: {result}")
return results
def _evaluate_task(self, execution_results: Dict[str, str]) -> Dict[str, dict]:
"""评估任务结果"""
evaluations = {}
for task, result in execution_results.items():
evaluation = self.reflection.evaluate(task, result)
evaluations[task] = evaluation
if not evaluation["is_correct"]:
print(f"任务评估未通过: {task}")
print(f"反馈: {evaluation['feedback']}")
return evaluations
def run(self, goal: str) -> Dict[str, Any]:
"""运行Agent处理目标"""
workflow_input = {"goal": goal}
workflow_result = self.workflow.execute(workflow_input)
return {
"tasks": workflow_result["plan"],
"results": workflow_result["execute"],
"evaluations": workflow_result["evaluate"]
}
12. 系统入口与示例执行
main.py作为系统入口:
python复制from agent_runtime import AgentRuntime
def main():
# 初始化Agent
agent = AgentRuntime()
# 添加示例知识到RAG系统
agent.rag.add_documents([
"东京是日本的首都,人口约1400万",
"1+1=2是最基本的数学公式",
"Python是一种流行的编程语言"
])
# 运行Agent处理用户目标
goal = "查询东京人口并计算其两倍值"
result = agent.run(goal)
# 打印执行结果
print("\n执行结果:")
for task, task_result in result["results"].items():
print(f"\n任务: {task}")
print(f"结果: {task_result}")
evaluation = result["evaluations"][task]
print(f"评估: {'通过' if evaluation['is_correct'] else '未通过'}")
if evaluation["feedback"]:
print(f"反馈: {evaluation['feedback']}")
if __name__ == "__main__":
main()
13. 生产级Agent的扩展方向
13.1 多Agent协作系统
类似Microsoft AutoGen的实现方式:
python复制class MultiAgentSystem:
def __init__(self):
self.agents = {
"planner": PlannerAgent(),
"researcher": ResearchAgent(),
"analyst": AnalysisAgent(),
"executor": ExecutionAgent()
}
self.workflow = WorkflowEngine()
# 配置多Agent工作流
self._setup_workflow()
def _setup_workflow(self):
self.workflow.add_task("plan", self.agents["planner"].plan)
self.workflow.add_task("research", self.agents["researcher"].execute)
self.workflow.add_task("analyze", self.agents["analyst"].process)
self.workflow.add_task("execute", self.agents["executor"].run)
# 建立任务依赖
self.workflow.graph.add_edges_from([
("plan", "research"),
("research", "analyze"),
("analyze", "execute")
])
13.2 复杂工作流引擎
基于LangGraph的增强实现:
python复制from langgraph.graph import Graph
from langgraph.nodes import ToolNode, ConditionalEdge
class AdvancedWorkflow:
def __init__(self):
self.graph = Graph()
# 定义节点
self.graph.add_node("plan", self._plan)
self.graph.add_node("execute", self._execute)
self.graph.add_node("evaluate", self._evaluate)
self.graph.add_node("revise", self._revise)
# 定义边
self.graph.add_edge("plan", "execute")
self.graph.add_conditional_edges(
"execute",
self._should_evaluate,
{
"evaluate": "evaluate",
"end": None
}
)
self.graph.add_edge("evaluate", "revise")
self.graph.add_edge("revise", "execute")
# 编译工作流
self.app = self.graph.compile()
def _should_evaluate(self, state):
"""决定是否需要评估"""
return "evaluate" if state.get("requires_eval", False) else "end"
13.3 代码环境集成
类似OpenDevin的代码环境交互:
python复制class CodeEnvironment:
def __init__(self, workspace_dir="workspace"):
self.workspace = Path(workspace_dir)
self.workspace.mkdir(exist_ok=True)
def edit_file(self, filepath: str, changes: dict):
"""应用代码变更"""
full_path = self.workspace / filepath
if not full_path.exists():
full_path.write_text("", encoding="utf-8")
content = full_path.read_text(encoding="utf-8")
# 应用变更
for change in changes.get("edits", []):
start = change["range"]["start"]
end = change["range"]["end"]
new_text = change["text"]
# 实现具体的文本替换逻辑
# ...
full_path.write_text(new_content, encoding="utf-8")
return True
def run_tests(self, command="pytest"):
"""运行测试"""
try:
result = subprocess.run(
command,
cwd=str(self.workspace),
shell=True,
check=True,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
text=True,
timeout=120
)
return {
"success": True,
"output": result.stdout
}
except subprocess.CalledProcessError as e:
return {
"success": False,
"output": e.stderr
}
14. 性能优化与生产实践
14.1 异步执行优化
python复制import asyncio
from concurrent.futures import ThreadPoolExecutor
class AsyncAgentRuntime:
def __init__(self, max_workers=4):
self.executor = ThreadPoolExecutor(max_workers)
self.loop = asyncio.get_event_loop()
async def execute_tasks(self, tasks):
"""异步执行任务列表"""
futures = [
self.loop.run_in_executor(
self.executor,
self._execute_single,
task
)
for task in tasks
]
return await asyncio.gather(*futures)
def _execute_single(self, task):
"""执行单个任务"""
# 实际任务执行逻辑
pass
14.2 缓存策略实现
python复制from functools import lru_cache
from datetime import datetime, timedelta
class LLMCache:
def __init__(self, maxsize=1000, ttl=3600):
self.maxsize = maxsize
self.ttl = timedelta(seconds=ttl)
self.cache = {}
def _get_key(self, prompt, temperature):
return hash((prompt, temperature))
def get(self, prompt, temperature):
key = self._get_key(prompt, temperature)
if key in self.cache:
entry = self.cache[key]
if datetime.now() - entry["timestamp"] < self.ttl:
return entry["response"]
return None
def set(self, prompt, temperature, response):
key = self._get_key(prompt, temperature)
if len(self.cache) >= self.maxsize:
self.cache.pop(next(iter(self.cache)))
self.cache[key] = {
"response": response,
"timestamp": datetime.now()
}
14.3 监控与日志系统
python复制import logging
from dataclasses import dataclass
from typing import List
@dataclass
class AgentEvent:
timestamp: str
event_type: str
details: dict
class MonitoringSystem:
def __init__(self):
self.events: List[AgentEvent] = []
self.logger = logging.getLogger("agent")
self.logger.setLevel(logging.INFO)
# 配置日志处理器
handler = logging.StreamHandler()
formatter = logging.Formatter(
"%(asctime)s - %(name)s - %(levelname)s - %(message)s"
)
handler.setFormatter(formatter)
self.logger.addHandler(handler)
def log_event(self, event_type, details):
"""记录Agent事件"""
event = AgentEvent(
timestamp=datetime.now().isoformat(),
event_type=event_type,
details=details
)
self.events.append(event)
self.logger.info(f"{event_type}: {details}")
def get_metrics(self):
"""获取性能指标"""
return {
"total_events": len(self.events),
"llm_calls": sum(
1 for e in self.events
if e.event_type == "llm_call"
),
"tool_executions": sum(
1 for e in self.events
if e.event_type == "tool_execution"
)
}
15. 安全与合规考虑
15.1 输入验证与过滤
python复制import re
class SecurityFilter:
def __init__(self):
self.sensitive_patterns = [
r"rm\s+-rf",
r"([\"'])(?:(?=(\\?))\2.)*?\1[|&;]",
r"\.\.\/",
r"\/etc\/passwd"
]
def sanitize_input(self, input_str):
"""净化用户输入"""
if not input_str or not isinstance(input_str, str):
return ""
# 移除潜在危险字符
sanitized = input_str.strip()
for pattern in self.sensitive_patterns:
sanitized = re.sub(pattern, "", sanitized, flags=re.IGNORECASE)
return sanitized
def validate_command(self, command):
"""验证命令安全性"""
if not command:
return False
blocked_commands = [
"rm", "shutdown", "reboot",
"chmod", "chown", "dd"
]
first_word = command.split()[0].lower()
return first_word not in blocked_commands
15.2 权限控制系统
python复制from enum import Enum, auto
class PermissionLevel(Enum):
GUEST = auto()
USER = auto()
ADMIN = auto()
class AccessControl:
def __init__(self):
self.permissions = {
"execute_tool": {
PermissionLevel.GUEST: ["search"],
PermissionLevel.USER: ["search", "calculator"],
PermissionLevel.ADMIN: ["search", "calculator", "shell"]
},
"file_access": {
PermissionLevel.GUEST: False,
PermissionLevel.USER: True,
PermissionLevel.ADMIN: True
}
}
def check_permission(self, user_level, action, resource=None):
"""检查用户权限"""
if action not in self.permissions:
return False
required_levels = self.permissions[action]
if resource and isinstance(required_levels, dict):
if resource in required_levels:
return user_level in required_levels[resource]
return user_level in required_levels
16. 测试与验证策略
16.1 单元测试示例
python复制import unittest
from unittest.mock import patch
from agent_runtime import AgentRuntime
class TestAgentRuntime(unittest.TestCase):
def setUp(self):
self.agent = AgentRuntime()
@patch.object(Planner, 'plan')
def test_planning(self, mock_plan):
mock_plan.return_value = ["step1", "step2"]
tasks = self.agent._plan_task("test goal")
self.assertEqual(tasks, ["step1", "step2"])
@patch.object(ToolSystem, 'select_tool')
@patch.object(ToolSystem, 'execute')
def test_execution(self, mock_execute, mock_select):
mock_select.return_value = "search"
mock_execute.return_value = "results"
results = self.agent._execute_task(["test task"])
self.assertEqual(results, {"test task": "results"})
16.2 集成测试框架
python复制class AgentIntegrationTest:
def __init__(self):
self.test_cases = [
{
"name": "simple_calculation",
"input": "计算3的平方",
"expected": "9"
},
{
"name": "information_query",
"input": "搜索Python的最新版本",
"expected": "Python"
}
]
def run_tests(self):
agent = AgentRuntime()
results = []
for case in self.test_cases:
actual = agent.run(case["input"])
passed = case["expected"] in str(actual)
results.append({
"name": case["name"],
"passed": passed,
"expected": case["expected"],
"actual": actual
})
return results
17. 部署与扩展架构
17.1 容器化部署
dockerfile复制# Dockerfile示例
FROM python:3.9-slim
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY . .
# 安装生产依赖
RUN apt-get update && apt-get install -y \
gcc \
python3-dev \
&& rm -rf /var/lib/apt/lists/*
# 设置环境变量
ENV PYTHONUNBUFFERED=1
ENV OPENAI_API_KEY="your-api-key"
CMD ["python", "main.py"]
17.2 水平扩展设计
python复制import redis
from fastapi import FastAPI
from fastapi.middleware.cors import CORSMiddleware
app = FastAPI()
app.add_middleware(
CORSMiddleware,
allow_origins=["*"],
allow_methods=["*"],
allow_headers=["*"],
)
# Redis连接池
redis_pool = redis.ConnectionPool(
host='redis',
port=6379,
db=0,
decode_responses=True
)
@app.post("/execute")
async def execute_goal(goal: str):
"""分布式执行端点"""
# 检查缓存
cache_key = f"goal:{hash(goal)}"
cached = redis_pool.get(cache_key)
if cached:
return {"result": cached, "cached": True}
# 实际执行
agent = AgentRuntime()
result = agent.run(goal)
# 设置缓存
redis_pool.setex(cache_key, 3600, str(result))
return {"result": result, "cached": False}
18. 实际应用案例
18.1 数据分析Agent
python复制class DataAnalysisAgent:
def __init__(self):
self.agent = AgentRuntime()
# 注册数据分析专用工具
self.agent.tools.register_tool(
"analyze_csv",
self.analyze_csv
)
def analyze_csv(self, filepath):
"""分析CSV文件"""
import pandas as pd
try:
df = pd.read_csv(filepath)
summary = df.describe().to_string()
return f"分析结果:\n{summary}"
except Exception as e:
return f"分析失败: {str(e)}"
def run_analysis(self, query):
"""运行分析任务"""
return self.agent.run(query)
18.2 自动化测试Agent
python复制class TestAutomationAgent:
def __init__(self):
self.agent = AgentRuntime()
self.code_env = CodeEnvironment()
# 注册测试相关工具
self.agent.tools.register_tool(
"write_test",
self.write_test_case
)
self.agent.tools.register_tool(
"run_test",
self.execute_test
)
def write_test_case(self, spec):
"""根据规范编写测试用例"""
prompt = f"""
根据以下规范编写Python测试用例:
{spec}
要求:
1. 使用pytest风格
2. 包含合理的断言
3. 代码格式规范
"""
code = self.agent.llm.chat(prompt, temperature=0.3)
return code
def execute_test(self, test_file):
"""执行测试文件"""
self.code_env.edit_file("test_temp.py", test_file)
return self.code_env.run_tests("pytest test_temp.py")
19. 性能基准测试
19.1 测试指标定义
python复制@dataclass
class BenchmarkResult:
total_tasks: int
success_rate: float
avg_response_time: float
max_memory_usage: float
llm_calls_per_task: float
class AgentBenchmark:
def __init__(self, agent):
self.agent = agent
self.test_cases = self._load_test_cases()
def _load_test_cases(self):
"""加载基准测试用例"""
return [
{"goal": "计算10的阶乘", "expected": "3628800"},
{"goal": "搜索AI的最新发展", "expected": "AI"},
{"goal": "查询伦敦人口", "expected": "人口"}
]
def run(self, iterations=10):
"""运行基准测试"""
results = []
for _ in range(iterations):
for case in self.test_cases:
start_time = time.time()
result = self.agent.run(case["goal"])
elapsed = time.time() - start_time
success = case["expected"] in str(result)
results.append({
"success": success,
"time": elapsed
})
success_rate = sum(
1 for r in results if r["success"]
) / len(results)
avg_time = sum(r["time"] for r in results) / len(results)
return BenchmarkResult(
total_tasks=len(results),
success_rate=success_rate,
avg_response_time=avg_time,
max_memory_usage=0, # 实际实现中需要测量
llm_calls_per_task=2.5 # 估算值
)
20. 持续改进方向
20.1 模型微调策略
python复制class FineTuningManager:
def __init__(self, agent):
self.agent = agent
self.feedback_data = []
def collect_feedback(self, task, result, is_correct, feedback):
"""收集用户反馈数据"""
self.feedback_data.append({
"task": task,
"result": result,
"is_correct": is_correct,
"feedback": feedback,
"timestamp": datetime.now().isoformat()
})
def prepare_training_data(self):
"""准备微调数据"""
training_examples = []
for item in self.feedback_data:
if item["is_correct"]:
continue
training_examples.append({
"prompt": item["task"],
"completion": item["feedback"],
"metadata": {
"original_result": item["result"],
"timestamp": item["timestamp"]
}
})
return training_examples
def schedule_fine_tuning(self):
"""安排模型微调任务"""
training_data = self.prepare_training_data()
if not training_data:
return False
# 实际实现中调用微调API
# openai.FineTuningJob.create(...)
return True
20.2 自适应学习机制
python复制class AdaptiveLearner:
def __init__(self, agent):
self.agent = agent
self.performance_log = {}
def log_performance(self, task_type, success):
"""记录任务性能"""
if task_type not in self.performance_log:
self.performance_log[task_type] = {
"success": 0,
"total": 0
}
self.performance_log[task_type]["total"] += 1
if success:
self.performance_log[task_type]["success"] += 1
def get_success_rate(self, task_type):
"""获取任务类型成功率"""
if task_type not in self.performance_log:
return 0.5 # 默认值
stats = self.performance_log[task_type]
return stats["success"] / stats["total"]
def adjust_parameters(self, task_type):
"""根据历史表现调整参数"""
success_rate = self.get_success_rate(task_type)
if success_rate < 0.3:
# 低成功率任务增加反思步骤
self.agent.reflection.enabled = True
self.agent.reflection.retry_count = 2
elif success_rate > 0.8:
# 高成功率任务减少冗余检查
self.agent.reflection.enabled = False
构建生产级Agent系统是一个持续迭代的过程,关键在于建立坚实的架构基础,同时保持各个组件的模块化和可扩展性。通过本文介绍的500行左右的核心框架,开发者可以快速搭建一个功能完整的Agent系统原型,然后根据实际需求逐步扩展和完善各个功能模块。
