从零构建 AI Agent:工具调用 + 记忆管理 + 多 Agent 协作实战

一、引言

2024 年是 AI Agent 元年。从 AutoGPT 到 OpenAI Assistants,从 LangChain Agent 到 Coze 平台,"智能体"概念正在从实验室走向生产。但大多数"Agent"只是简单的 LLM + 工具调用的包装,离真正的自主决策、多步推理、记忆演化还有不小距离。

本文将带你从零构建一个具备工具调用、短期/长期记忆、多 Agent 协作能力的生产级 Agent 系统,使用 LangGraph 作为编排框架。

二、Agent 核心架构设计

2.1 四层架构

应用层 → 决策层(ReAct/Plan-Execute) → 能力层(工具调用/代码执行) → 记忆层(短期/长期/工作)

2.2 ReAct 模式实现

ReAct(Reasoning + Acting)是 Agent 决策的核心范式:

Thought → Action → Observation → Thought → Action → ... → Final Answer
from typing import List, Dict, Any, TypedDict, Annotated
from langgraph.graph import StateGraph, END
from langgraph.prebuilt import ToolNode
from langchain_core.messages import HumanMessage, AIMessage, ToolMessage
from langchain_openai import ChatOpenAI
from langchain_core.tools import tool
import operator

class AgentState(TypedDict):
    messages: Annotated[List, operator.add]
    tools_called: List[str]
    iterations: int

# === 工具定义 ===
@tool
def search_knowledge_base(query: str) -> str:
    """搜索内部知识库"""
    kb = {
        "退款政策": "用户可在购买后30天内申请全额退款。需提供订单号和退款原因。",
        "VIP权益": "VIP会员享有优先客服、免费配送、专属折扣9折。",
        "配送时效": "全国主要城市1-3天,偏远地区3-7天。"
    }
    for key, value in kb.items():
        if key in query:
            return value
    return f"未找到关于 '{query}' 的相关信息"

@tool
def calculate(expression: str) -> str:
    """计算数学表达式"""
    try:
        result = eval(expression, {"__builtins__": {}}, {})
        return str(result)
    except Exception as e:
        return f"计算错误: {e}"

@tool
def get_current_time() -> str:
    """获取当前时间"""
    from datetime import datetime
    return datetime.now().strftime("%Y-%m-%d %H:%M:%S")

class ReActAgent:
    def __init__(self, model_name: str = "gpt-4"):
        self.llm = ChatOpenAI(model=model_name, temperature=0)
        self.tools = [search_knowledge_base, calculate, get_current_time]
        self.llm_with_tools = self.llm.bind_tools(self.tools)
        self.graph = self._build_graph()
    
    def _build_graph(self):
        workflow = StateGraph(AgentState)
        workflow.add_node("agent", self._agent_node)
        workflow.add_node("tools", ToolNode(self.tools))
        workflow.set_entry_point("agent")
        workflow.add_conditional_edges("agent", self._should_continue, {"continue": "tools", "end": END})
        workflow.add_edge("tools", "agent")
        return workflow.compile()
    
    def _agent_node(self, state: AgentState) -> Dict:
        response = self.llm_with_tools.invoke(state["messages"])
        return {"messages": [response]}
    
    def _should_continue(self, state: AgentState) -> str:
        last_message = state["messages"][-1]
        if hasattr(last_message, "tool_calls") and last_message.tool_calls:
            return "continue"
        return "end"
    
    def run(self, user_input: str) -> str:
        result = self.graph.invoke({"messages": [HumanMessage(content=user_input)]})
        for msg in reversed(result["messages"]):
            if isinstance(msg, AIMessage) and msg.content:
                return msg.content
        return "无法生成回答"

# 使用示例
agent = ReActAgent()
response = agent.run("VIP有什么权益?帮我算一下100块钱打9折是多少")
print(response)

三、记忆系统设计

3.1 三层记忆架构

import chromadb
from datetime import datetime
from typing import List, Dict, Any, Optional

class MemorySystem:
    """Agent 三层记忆系统:短期/长期/工作记忆"""
    
    def __init__(self, persist_dir: str = "./agent_memory"):
        self.client = chromadb.PersistentClient(path=persist_dir)
        self.short_term: List[Dict] = []
        self.max_short_term = 20
        self.long_term = self.client.get_or_create_collection("long_term")
        self.working: Dict[str, Any] = {}
    
    def add_short_term(self, role: str, content: str):
        self.short_term.append({"role": role, "content": content, "time": datetime.now()})
        if len(self.short_term) > self.max_short_term:
            # 滑动窗口 + 摘要压缩
            system_msgs = [m for m in self.short_term if m["role"] == "system"]
            recent = self.short_term[-(self.max_short_term - len(system_msgs)):]
            self.short_term = system_msgs + recent
    
    def get_short_term_context(self, n: int = 10) -> str:
        return "\n".join([f"[{m['role']}] {m['content'][:200]}" for m in self.short_term[-n:]])
    
    def save_long_term(self, key: str, value: str, metadata: dict = None):
        self.long_term.add(
            documents=[value],
            metadatas=[{"key": key, "timestamp": datetime.now().isoformat(), **(metadata or {})}],
            ids=[f"mem_{datetime.now().timestamp()}_{hash(key) % 10000}"]
        )
    
    def recall(self, query: str, top_k: int = 5) -> List[str]:
        results = self.long_term.query(query_texts=[query], n_results=top_k)
        return results.get("documents", [[]])[0]
    
    def set_working(self, key: str, value: Any):
        self.working[key] = value
    
    def get_working(self, key: str) -> Optional[Any]:
        return self.working.get(key)
    
    def consolidate(self):
        """将短期记忆中的重要信息整合到长期记忆"""
        if len(self.short_term) < 5:
            return
        recent = "\n".join([f"[{m['role']}]: {m['content'][:300]}" for m in self.short_term[-6:]])
        self.save_long_term(f"summary_{datetime.now().timestamp()}", recent)
        print(f"记忆已整合: {len(recent)} 字符")

四、多 Agent 协作

4.1 Supervisor 监督者模式

from langgraph.graph import StateGraph, END
from langgraph.prebuilt import create_react_agent
from typing import Literal, TypedDict

class SupervisorState(TypedDict):
    task: str
    research_result: str
    code_result: str
    review_result: str
    final_output: str
    next_step: str

class SupervisorAgent:
    def __init__(self):
        self.researcher = create_react_agent(
            ChatOpenAI(model="gpt-4"), [search_knowledge_base, get_current_time],
            state_modifier="你是研究员Agent。负责搜集信息、分析数据、提供报告。"
        )
        self.coder = create_react_agent(
            ChatOpenAI(model="gpt-4"), [calculate],
            state_modifier="你是程序员Agent。负责编写代码、调试、性能优化。"
        )
        self.reviewer = create_react_agent(
            ChatOpenAI(model="gpt-4"), [],
            state_modifier="你是审核员Agent。负责审查输出质量、发现错误。"
        )
        self.graph = self._build()
    
    def _build(self):
        workflow = StateGraph(SupervisorState)
        workflow.add_node("supervisor", self._supervisor_node)
        workflow.add_node("research", self._research_node)
        workflow.add_node("coding", self._coding_node)
        workflow.add_node("review", self._review_node)
        workflow.set_entry_point("supervisor")
        workflow.add_conditional_edges("supervisor", lambda s: s["next_step"], {
            "research": "research", "coding": "coding", "review": "review", "end": END
        })
        for node in ["research", "coding", "review"]:
            workflow.add_edge(node, "supervisor")
        return workflow.compile()
    
    def _supervisor_node(self, state: SupervisorState) -> Dict:
        if not state.get("research_result"):
            return {"next_step": "research"}
        elif not state.get("code_result"):
            return {"next_step": "coding"}
        elif not state.get("review_result"):
            return {"next_step": "review"}
        else:
            return {"next_step": "end", "final_output": f"## 研究\n{state['research_result']}\n## 代码\n{state['code_result']}\n## 审查\n{state['review_result']}"}
    
    def _research_node(self, state: SupervisorState) -> Dict:
        r = self.researcher.invoke({"messages": [HumanMessage(content=f"研究: {state['task']}")]})
        return {"research_result": r["messages"][-1].content}
    
    def _coding_node(self, state: SupervisorState) -> Dict:
        r = self.coder.invoke({"messages": [HumanMessage(content=f"基于研究编写代码:\n{state['research_result']}")]})
        return {"code_result": r["messages"][-1].content}
    
    def _review_node(self, state: SupervisorState) -> Dict:
        ctx = f"## 研究\n{state['research_result']}\n## 代码\n{state['code_result']}"
        r = self.reviewer.invoke({"messages": [HumanMessage(content=f"审查:\n{ctx}")]})
        return {"review_result": r["messages"][-1].content}

4.2 Plan-Execute 层级模式

import json

class HierarchicalAgent:
    def plan_and_execute(self, goal: str) -> str:
        # 1. 规划
        planner_prompt = f"""将以下目标分解为可执行的子任务,输出JSON:
目标: {goal}
格式: [{{"step": 1, "description": "...", "agent": "researcher|coder", "depends_on": []}}]"""
        
        plan_response = self.llm.invoke(planner_prompt)
        plan = json.loads(plan_response.content)
        
        # 2. 按依赖顺序执行
        results = {}
        for step in plan:
            deps_met = all(d in results for d in step.get("depends_on", []))
            if not deps_met:
                continue
            if step["agent"] == "researcher":
                result = self.researcher.invoke(step["description"])
            else:
                result = self.coder.invoke(step["description"])
            results[step["step"]] = result
        
        # 3. 汇总
        summary = self.llm.invoke(f"汇总以下执行结果:\n{json.dumps(results, ensure_ascii=False, indent=2)}")
        return summary.content

五、工具系统设计

from pydantic import BaseModel, Field
from typing import Type, Callable, Dict, List
import inspect

class ToolRegistry:
    def __init__(self):
        self.tools: Dict[str, Callable] = {}
        self.schemas: Dict[str, Type[BaseModel]] = {}
    
    def register(self, name: str, description: str, schema: Type[BaseModel] = None):
        def decorator(func: Callable):
            if schema is None:
                sig = inspect.signature(func)
                fields = {}
                for pname, param in sig.parameters.items():
                    if pname in ('self', 'cls'): continue
                    annotation = param.annotation if param.annotation != inspect.Parameter.empty else str
                    fields[pname] = (annotation, Field(description=f"参数: {pname}"))
                dynamic_schema = type(f"{name}Schema", (BaseModel,), fields)
            else:
                dynamic_schema = schema
            
            self.tools[name] = func
            self.schemas[name] = dynamic_schema
            func.__tool_name__ = name
            func.__tool_description__ = description
            func.__tool_schema__ = dynamic_schema
            return func
        return decorator
    
    def get_tool_schemas(self) -> List[Dict]:
        return [{"type": "function", "function": {"name": n, "description": f.__tool_description__, "parameters": s.model_json_schema()}} for n, f in self.tools.items() for s in [self.schemas[n]]]
    
    def execute(self, tool_name: str, arguments: dict) -> str:
        if tool_name not in self.tools:
            return f"错误: 未找到工具 '{tool_name}'"
        try:
            return str(self.tools[tool_name](**arguments))
        except Exception as e:
            return f"工具执行错误: {e}"

# 使用示例
registry = ToolRegistry()

@registry.register(name="send_email", description="发送邮件")
def send_email(to: str, subject: str, body: str) -> str:
    return f"邮件已发送至 {to},主题: {subject}"

@registry.register(name="create_issue", description="创建GitHub Issue")
def create_issue(repo: str, title: str, body: str = "", labels: list = None) -> str:
    return f"Issue '{title}' 已创建于 {repo}"

六、FastAPI 部署

from fastapi import FastAPI
from pydantic import BaseModel
import uuid

app = FastAPI(title="AI Agent Service")

class AgentRequest(BaseModel):
    task: str
    agent_type: str = "react"
    max_iterations: int = 10

@app.post("/agent/run")
async def run_agent(request: AgentRequest):
    task_id = str(uuid.uuid4())[:8]
    if request.agent_type == "supervisor":
        agent = SupervisorAgent()
        result = agent.graph.invoke({"task": request.task})
        answer = result.get("final_output", "")
    else:
        agent = ReActAgent()
        answer = agent.run(request.task)
    return {"task_id": task_id, "answer": answer}

七、总结

本文从零构建了生产级 AI Agent 的四大核心能力:

  1. ReAct 决策:Thought→Action→Observation 循环,LangGraph 编排
  2. 三层记忆:短期滑动窗口、长期向量库持久化、工作记忆
  3. 多 Agent 协作:Supervisor 监督者模式 + 层级 Plan-Execute
  4. 工具系统:类型安全的工具注册中心,兼容 Function Calling

这套架构可以直接嵌入实际产品中,从简单问答到多步骤任务自动化全覆盖。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐