2026年AI Agent开发框架实战指南:从零构建多Agent协作系统

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应用,主要集中在以下场景:

  • 智能客服与工单处理:多Agent协同处理复杂用户问题

  • 代码开发与审查:编码Agent + 审查Agent + 测试Agent的流水线

  • 数据分析与报告:数据采集Agent + 分析Agent + 撰写Agent协作

  • 企业知识管理:检索Agent + 摘要Agent + 问答Agent的知识链路
  • 1.2 主流框架对比

    2026年的AI Agent框架生态已经形成了清晰的分层格局:

    | 框架 | 定位 | 核心优势 | 适用场景 |
    |------|------|----------|----------|
    | LangChain / LangGraph | 通用编排框架 | 生态最完善,工具链丰富 | 复杂工作流、生产级应用 |
    | AutoGen | 多Agent对话框架 | 多角色对话原生支持 | 多Agent协商、群体决策 |
    | CrewAI | 角色化Agent框架 | 角色定义简洁,任务分配清晰 | 团队协作、任务型Agent |
    | DeepSeek Harness | 国产推理框架 | 本地部署友好,中文优化 | 私有化部署、企业内网 |
    | OpenAI Codex Harness | 代码生成框架 | 代码理解与生成能力强 | 代码Agent、DevOps场景 |

    1.3 关键技术趋势


  • MCP(Model Context Protocol)成为工具调用新标准:取代了传统的Function Calling,提供统一的工具注册、发现和调用协议

  • 长期记忆(Long-term Memory)成为标配:向量数据库 + 知识图谱的混合记忆方案逐渐成熟

  • 沙箱执行环境:代码执行、文件操作等危险动作全部在沙箱中完成

  • 本地/私有化部署:企业对数据安全的要求推动本地大模型和Agent框架的普及

  • 二、环境搭建:从零开始准备开发环境

    2.1 基础环境要求

    在开始之前,请确保你的开发环境满足以下要求:

    bash
    # Python 版本要求 3.10+
    python --version  # Python 3.11.7
    
    # 推荐使用虚拟环境管理工具
    # 可选:venv / conda / poetry / uv

    2.2 核心依赖安装

    本文我们将以 LangGraph(LangChain的Agent编排层) 为核心框架,结合MCP协议和本地向量存储,构建一套完整的多Agent系统。

    bash
    # 创建项目目录
    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 uvicorn

    2.3 配置文件与环境变量

    创建 .env 文件管理配置:

    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=30

    2.4 项目目录结构

    推荐的项目目录结构如下:

    text
    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由以下核心组件构成:

    text
    ┌─────────────────────────────────────────────────┐
    │                    Agent                        │
    │  ┌─────────┐  ┌─────────┐  ┌─────────────────┐ │
    │  │  LLM    │  │  Memory │  │   Tools / MCP   │ │
    │  │ 大脑    │  │ 记忆    │  │   工具集        │ │
    │  └────┬────┘  └────┬────┘  └────────┬────────┘ │
    │       │            │                │          │
    │  ┌────┴────────────┴────────────────┴───────┐  │
    │  │          Planning / Reasoning             │  │
    │  │          规划与推理引擎                   │  │
    │  └──────────────────┬───────────────────────┘  │
    │                     │                          │
    │  ┌──────────────────┴───────────────────────┐  │
    │  │          Action Execution                 │  │
    │  │          动作执行器                       │  │
    │  └──────────────────────────────────────────┘  │
    └─────────────────────────────────────────────────┘

    各组件职责:

  • LLM(大语言模型):Agent的"大脑",负责理解、推理和生成

  • Memory(记忆系统):存储对话历史、知识和经验

  • Tools(工具集):通过MCP协议调用外部能力

  • Planning(规划引擎):将复杂任务拆解为子任务

  • Action(执行器):执行具体动作并观察结果
  • 3.2 ReAct范式:思考与行动的循环

    当前主流的Agent推理范式是 ReAct(Reasoning + Acting),其核心循环如下:

    text
    思考(Thought) → 行动(Action) → 观察(Observation) → 思考(Thought) → ...

    用伪代码表示:

    python
    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系统常见的协作模式有三种:

  • 层级式(Hierarchical):一个主管Agent分配任务给下属Agent

  • 流水线式(Pipeline):任务按顺序在Agent间传递

  • 群聊式(Group Chat):多个Agent自由讨论,由仲裁者决策
  • 本文将实现一个层级+流水线混合的多Agent系统,用于自动化研究报告生成。


    四、单Agent实现:打造你的第一个智能体

    4.1 基础Agent类设计

    首先,我们定义一个可扩展的基础Agent类:

    python
    # 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.content

    4.2 短期记忆实现

    python
    # 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 工具注册中心

    python
    # 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 一个简单的计算器工具

    python
    # 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

    创建一个简单的测试入口:

    python
    # 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())

    运行结果示例:

    text
    用户: 1234的平方根是多少?
    助手: 1234的平方根约等于 35.13(精确值为 35.12833614050059)。


    五、多Agent协作:构建研究报告生成系统

    5.1 系统架构设计

    我们将构建一个三Agent协作的研究报告生成系统:

    text
    用户请求
        │
        ▼
    ┌─────────────────┐
    │   协调者Agent   │  分配任务、整合结果
    │  (Coordinator)  │
    └───┬─────────┬───┘
        │         │
        ▼         ▼
    ┌────────┐ ┌────────┐
    │研究员  │ │ 写作   │
    │Agent   │ │ Agent  │
    │(搜索、 │ │(撰写、 │
    │ 分析)  │ │ 排版) │
    └────┬───┘ └────┬───┘
         │          │
         ▼          ▼
      搜索工具    代码执行
      计算器      文件操作
         │          │
         └─────┬────┘
               ▼
         ┌───────────┐
         │ 记忆系统   │
         │ (短期+长期)│
         └───────────┘

    5.2 定义三个角色Agent

    python
    # 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)

    python
    # 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)

    python
    # 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协作的状态机工作流。

    python
    # 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系统

    python
    # 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的核心优势:

  • 标准化:统一的工具描述和调用格式

  • 可发现性:Agent可以动态发现可用工具

  • 安全性:内置权限控制和审计机制

  • 跨平台:支持HTTP、WebSocket、本地进程等多种传输方式
  • 6.2 实现网络搜索工具

    python
    # 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工具调用的重中之重。代码执行工具必须运行在沙箱环境中:

    python
    # 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也能调用我们的工具:

    python
    # 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工具服务器:

    bash
    # 启动工具服务器
    uvicorn tools.mcp_server:app --host 0.0.0.0 --port 8000


    七、记忆系统:让Agent拥有长期记忆

    7.1 记忆系统架构

    一个完善的记忆系统应该包含三个层次:

    text
    ┌─────────────────────────────────────────┐
    │           工作记忆 (Working Memory)      │
    │   当前对话上下文,滑动窗口管理            │
    │   存储:内存 / 时效:即时                │
    ├─────────────────────────────────────────┤
    │           短期记忆 (Short-term Memory)   │
    │   会话历史,支持摘要压缩                 │
    │   存储:内存 / 时效:会话级              │
    ├─────────────────────────────────────────┤
    │           长期记忆 (Long-term Memory)    │
    │   知识库 + 向量检索 + 知识图谱           │
    │   存储:向量DB / 时效:永久              │
    └─────────────────────────────────────────┘

    7.2 长期记忆实现

    python
    # 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 0

    7.3 记忆增强(RAG)集成到Agent

    将长期记忆集成到BaseAgent中,实现记忆增强的推理:

    python
    # 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在对话结束后,自动将重要内容存入长期记忆:

    python
    # 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调用:

    python
    # 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容器化部署

    dockerfile
    # 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"]

    yaml
    # 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调用:

    python
    # 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系统天然适合并发执行。以下是一个并发研究的例子:

    python
    # 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 流式响应

    对于长对话和长文本生成,使用流式响应提升用户体验:

    python
    # 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系统必须具备完善的可观测性:

    python
    # 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协作系统,涵盖:

  • 基础架构:基于LangGraph的Agent编排框架

  • 核心组件:单Agent实现、工具注册中心、记忆系统

  • 多Agent协作:研究员、写作者、审稿者三角色协作

  • 工具系统:MCP协议兼容的工具调用机制

  • 记忆系统:短期记忆 + 长期向量记忆

  • 部署优化:API服务化、Docker部署、性能优化、可观测性
  • 9.2 下一步扩展方向


  • 本地模型部署:将LLM替换为本地部署的开源模型(如Qwen、DeepSeek等)

  • 知识图谱集成:在向量记忆基础上增加知识图谱,提升复杂推理能力

  • 多模态支持:增加图片理解、文档解析等多模态能力

  • Agent市场:构建可复用的Agent组件市场,实现Agent的即插即用

  • 人机协作:增加人类反馈环节,实现Human-in-the-loop的协作模式
  • 9.3 写在最后

    2026年是AI Agent从概念走向工程化落地的关键一年。作为开发者,我们不仅要掌握框架的使用,更要理解Agent系统的设计原理:思考-行动-观察的循环逻辑、记忆的分层管理、工具的安全调用、多Agent的协作模式。

    希望本文能帮助你快速上手AI Agent开发,在这个充满机遇的技术浪潮中找到自己的位置。


    参考资源:

  • LangGraph官方文档:https://langchain-ai.github.io/langgraph/

  • MCP协议规范:https://modelcontextprotocol.io/

  • Chroma向量数据库:https://www.trychroma.com/

  • AutoGen框架:https://microsoft.github.io/autogen/

  • CrewAI框架:https://www.crewai.com/

  • 本文为原创技术文章,欢迎转载,请注明出处。

    💬 评论区 (0)

    暂无评论,快来抢沙发吧!