首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >AI Agent全栈工程师实战:从零构建生产级智能体系统

AI Agent全栈工程师实战:从零构建生产级智能体系统

原创
作者头像
用户12608867
发布2026-08-29 11:45:40
发布2026-08-29 11:45:40
310
举报

AI Agent全栈工程师实战:从零构建生产级智能体系统

引言

AI Agent 不再是学术玩具。2026 年,企业级智能体已从“对话玩具”进化为“数字员工”——能自主规划、调用工具、记忆上下文、甚至自我纠错。作为全栈工程师,仅会调用 LLM API 远远不够;你需要掌握状态管理、工具抽象、规划算法、流式传输、可观测性等系统性能力。

本文不教概念,直接带您从零搭建一个生产级 AI Agent 系统,包含:

  • 基于 FastAPI 的异步后端,支持 SSE 流式输出
  • ReAct + 反思 双模式规划器
  • 可插拔的工具注册中心(含本地函数 + 远程 API)
  • 分层记忆系统(短期对话 + 长期向量检索)
  • 前端极简交互(原生 JS + EventSource)

全部代码可运行,技术点直击工程痛点。


一、系统架构全景

代码语言:javascript
复制
┌─────────────────────────────────────────────────────────────┐
│                     前端 (HTML + JS)                        │
│                  EventSource 接收 SSE 流                    │
└──────────────────────────┬──────────────────────────────────┘
                           │
┌──────────────────────────▼──────────────────────────────────┐
│                    FastAPI 网关层                          │
│              - 鉴权 / 限流 / 日志                          │
│              - SSE 流式响应                               │
└──────────┬──────────────┬──────────────┬──────────────────┘
           │              │              │
┌──────────▼──────────┐ ┌─▼──────────────▼─┐ ┌─────────────▼──┐
│   规划器 (Planner)   │ │  记忆管理器       │ │  工具注册中心   │
│  - ReAct 循环        │ │  - 滑动窗口缓存   │ │  - 函数式工具   │
│  - 自我反思 (可选)   │ │  - 向量检索 (FAISS)│ │  - HTTP 工具   │
└──────────┬──────────┘ └───────────────────┘ └────────────────┘
           │
┌──────────▼──────────────────────────────────────────────────┐
│              LLM 抽象层 (支持 OpenAI / 本地)                │
│           - 异步客户端 + 重试 + 回退                       │
│           - 结构化输出 (JSON Schema 强制)                  │
└─────────────────────────────────────────────────────────────┘

二、环境与依赖

使用 Python 3.11+,依赖精简:

代码语言:javascript
复制
pip install fastapi uvicorn[standard] openai httpx pydantic faiss-cpu numpy python-dotenv

项目结构:

代码语言:javascript
复制
agent_workshop/
├── main.py                 # FastAPI 入口
├── core/
│   ├── llm.py              # LLM 统一客户端
│   ├── planner.py          # ReAct + 反思规划器
│   ├── tools.py            # 工具注册与执行
│   ├── memory.py           # 短期缓存 + 长期向量记忆
│   └── schemas.py          # Pydantic 模型
├── static/
│   └── index.html          # 前端页面
└── .env                    # API_KEY 等

三、核心模块实现

3.1 统一 LLM 客户端(core/llm.py

支持 OpenAI 兼容接口,内置指数退避重试和结构化输出。

代码语言:javascript
复制
import asyncio
import json
from typing import List, Dict, Any, Optional
from openai import AsyncOpenAI
from pydantic import BaseModel
from tenacity import retry, stop_after_attempt, wait_exponential
import os

class LLMClient:
    def __init__(self, model: str = "gpt-4o-mini", base_url: Optional[str] = None):
        self.model = model
        self.client = AsyncOpenAI(
            api_key=os.getenv("OPENAI_API_KEY"),
            base_url=base_url or os.getenv("OPENAI_BASE_URL"),
            timeout=60.0
        )
    
    @retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=1, max=10))
    async def chat_completion(
        self,
        messages: List[Dict[str, str]],
        temperature: float = 0.3,
        response_format: Optional[BaseModel] = None,
        stream: bool = False
    ) -> Any:
        """统一调用,支持 JSON Schema 结构化输出"""
        kwargs = {
            "model": self.model,
            "messages": messages,
            "temperature": temperature,
            "stream": stream
        }
        if response_format:
            # 强制结构化输出(OpenAI 最新 API)
            kwargs["response_format"] = {
                "type": "json_schema",
                "json_schema": {
                    "name": response_format.__name__,
                    "schema": response_format.model_json_schema()
                }
            }
        if stream:
            return await self.client.chat.completions.create(**kwargs)
        else:
            resp = await self.client.chat.completions.create(**kwargs)
            content = resp.choices[0].message.content
            if response_format:
                return response_format.model_validate_json(content)
            return content

技术亮点:通过 response_format 强制 LLM 输出符合 Pydantic Schema 的 JSON,避免解析失败,是生产环境必备。


3.2 工具注册中心(core/tools.py

支持两种工具:本地 Python 函数(如计算、时间)和 远程 HTTP API(如天气查询)。采用注册表模式,便于动态添加。

代码语言:javascript
复制
from typing import Callable, Dict, Any, Optional, Awaitable
from pydantic import BaseModel, Field
import httpx
import json

class ToolDefinition(BaseModel):
    name: str
    description: str
    parameters: Dict[str, Any]  # JSON Schema
    handler: Optional[Callable[..., Awaitable[Any]]] = None
    http_endpoint: Optional[str] = None
    http_method: str = "GET"

class ToolRegistry:
    def __init__(self):
        self._tools: Dict[str, ToolDefinition] = {}
        self._http_client = httpx.AsyncClient(timeout=30.0)
    
    def register(self, tool_def: ToolDefinition):
        self._tools[tool_def.name] = tool_def
    
    def get_schema_for_llm(self) -> List[Dict]:
        """返回给 LLM 的工具描述列表(符合 function calling 格式)"""
        return [
            {
                "type": "function",
                "function": {
                    "name": t.name,
                    "description": t.description,
                    "parameters": t.parameters
                }
            }
            for t in self._tools.values()
        ]
    
    async def execute(self, name: str, arguments: Dict) -> str:
        tool = self._tools.get(name)
        if not tool:
            return f"Error: Tool '{name}' not found"
        try:
            if tool.handler:
                # 本地异步函数调用
                result = await tool.handler(**arguments)
                return json.dumps(result, ensure_ascii=False)
            elif tool.http_endpoint:
                # 远程 HTTP 调用
                resp = await self._http_client.request(
                    method=tool.http_method,
                    url=tool.http_endpoint,
                    params=arguments if tool.http_method == "GET" else None,
                    json=arguments if tool.http_method == "POST" else None
                )
                resp.raise_for_status()
                return resp.text
            else:
                return "Error: No handler or endpoint"
        except Exception as e:
            return f"Error: {str(e)}"

# ---------- 内置工具示例 ----------
async def get_current_time(timezone: str = "UTC") -> str:
    from datetime import datetime
    import pytz
    tz = pytz.timezone(timezone)
    return datetime.now(tz).isoformat()

async def calculate(expression: str) -> float:
    # 安全计算(仅允许基本运算)
    allowed = {"+", "-", "*", "/", "(", ")", " ", "0","1","2","3","4","5","6","7","8","9","."}
    if not all(c in allowed for c in expression):
        raise ValueError("Invalid expression")
    return eval(expression)

# 注册
registry = ToolRegistry()
registry.register(ToolDefinition(
    name="get_current_time",
    description="获取指定时区的当前时间",
    parameters={
        "type": "object",
        "properties": {"timezone": {"type": "string", "default": "UTC"}}
    },
    handler=get_current_time
))
registry.register(ToolDefinition(
    name="calculate",
    description="计算数学表达式,如 '2+3*4'",
    parameters={
        "type": "object",
        "properties": {"expression": {"type": "string"}},
        "required": ["expression"]
    },
    handler=calculate
))

生产级考量:远程工具支持重试、超时、熔断(可扩展);本地函数用 eval 做了严格字符过滤,防止注入。


3.3 分层记忆系统(core/memory.py

短期记忆:滑动窗口(最近 N 轮对话);长期记忆:基于 FAISS 的向量检索,用于历史相关事实召回。

代码语言:javascript
复制
import faiss
import numpy as np
from typing import List, Dict, Tuple
from collections import deque
import hashlib
import json

class ShortTermMemory:
    """滑动窗口对话缓存"""
    def __init__(self, max_turns: int = 20):
        self.buffer = deque(maxlen=max_turns)
    
    def add(self, role: str, content: str):
        self.buffer.append({"role": role, "content": content})
    
    def get_messages(self) -> List[Dict]:
        return list(self.buffer)

class LongTermMemory:
    """基于 FAISS 的向量记忆,存储事实型知识"""
    def __init__(self, embedding_dim: int = 1536, max_items: int = 1000):
        self.dim = embedding_dim
        self.index = faiss.IndexFlatL2(embedding_dim)
        self.metadata: List[Dict] = []  # 存储原始文本和时间戳
        self.max_items = max_items
    
    def _embed(self, text: str) -> np.ndarray:
        # 实际项目可调用 OpenAI embedding 或本地模型
        # 此处用伪随机确定性向量演示(生产请替换)
        np.random.seed(hashlib.md5(text.encode()).hexdigest())
        return np.random.randn(self.dim).astype(np.float32)
    
    def add(self, text: str, meta: Dict = None):
        if len(self.metadata) >= self.max_items:
            # 移除最早的一半
            self.index = faiss.IndexFlatL2(self.dim)
            self.metadata = self.metadata[-self.max_items//2:]
            # 重建索引(简化)
            for item in self.metadata:
                vec = self._embed(item["text"])
                self.index.add(vec.reshape(1, -1))
        vec = self._embed(text)
        self.index.add(vec.reshape(1, -1))
        self.metadata.append({"text": text, "meta": meta or {}})
    
    def search(self, query: str, k: int = 3) -> List[str]:
        vec = self._embed(query)
        distances, indices = self.index.search(vec.reshape(1, -1), k)
        results = []
        for idx in indices[0]:
            if idx < len(self.metadata) and idx >= 0:
                results.append(self.metadata[idx]["text"])
        return results

# 统一记忆管理器
class MemoryManager:
    def __init__(self, short_term_max=20):
        self.short = ShortTermMemory(short_term_max)
        self.long = LongTermMemory()
    
    def add_interaction(self, user_msg: str, assistant_msg: str):
        self.short.add("user", user_msg)
        self.short.add("assistant", assistant_msg)
        # 对长记忆进行关键事实提取(此处简化为存入完整对话)
        # 生产环境可用 LLM 做摘要后存入
        self.long.add(f"User: {user_msg}\nAssistant: {assistant_msg}")
    
    def build_context(self, current_query: str) -> Tuple[List[Dict], List[str]]:
        """返回 (短期消息列表, 长期相关记忆列表)"""
        short_msgs = self.short.get_messages()
        long_facts = self.long.search(current_query, k=2)
        return short_msgs, long_facts

说明:本 demo 的 embedding 用确定性随机模拟,真实场景请接入 text-embedding-3-smallsentence-transformers


3.4 规划器:ReAct + 反思(core/planner.py

核心是 plan_loop,交替执行“思考-行动-观察”,并支持反思模式(自我校验后重试)。

代码语言:javascript
复制
from core.llm import LLMClient
from core.tools import registry
from core.memory import MemoryManager
from core.schemas import ActionStep, FinalAnswer
import json
import asyncio

class AgentPlanner:
    def __init__(self, llm: LLMClient, memory: MemoryManager, max_steps: int = 5):
        self.llm = llm
        self.memory = memory
        self.max_steps = max_steps
    
    async def plan(self, user_query: str, enable_reflection: bool = True) -> str:
        # 1. 构建上下文
        short_msgs, long_facts = self.memory.build_context(user_query)
        system_prompt = self._build_system_prompt(long_facts)
        messages = [{"role": "system", "content": system_prompt}] + short_msgs
        messages.append({"role": "user", "content": user_query})
        
        step_count = 0
        while step_count < self.max_steps:
            step_count += 1
            # 调用 LLM 获取下一步动作(结构化输出)
            action = await self.llm.chat_completion(
                messages=messages,
                response_format=ActionStep,
                temperature=0.2
            )
            # action 是 ActionStep 实例,包含 thought, action_name, action_input
            if action.action_name == "final_answer":
                # 最终回答
                final = action.action_input  # 直接作为答案
                # 存入记忆
                self.memory.add_interaction(user_query, final)
                return final
            
            # 执行工具
            if action.action_name in registry._tools:
                observation = await registry.execute(action.action_name, action.action_input)
            else:
                observation = f"Error: Unknown action {action.action_name}"
            
            # 将思考、行动、观察加入消息历史
            messages.append({
                "role": "assistant",
                "content": f"Thought: {action.thought}\nAction: {action.action_name}\nAction Input: {json.dumps(action.action_input)}"
            })
            messages.append({
                "role": "user",
                "content": f"Observation: {observation}"
            })
            
            # ---- 反思机制 ----
            if enable_reflection and step_count >= 2:
                # 检查是否需要重新规划
                reflect_prompt = f"你上一步的行动得到了观察:{observation}。请判断这个观察是否解决了用户问题?如果未解决,请提出修正计划,否则继续。"
                messages.append({"role": "user", "content": reflect_prompt})
                reflect_decision = await self.llm.chat_completion(messages, temperature=0.3)
                if "修正" in reflect_decision or "重新" in reflect_decision:
                    # 添加一条修正指令,强制重新思考
                    messages.append({"role": "assistant", "content": "收到,我将重新规划。"})
                    continue
        
        # 超过最大步数,强制返回
        return "抱歉,我未能完成您的请求,请简化问题。"

    def _build_system_prompt(self, long_facts: List[str]) -> str:
        tools_desc = registry.get_schema_for_llm()
        tool_names = [t["function"]["name"] for t in tools_desc]
        facts_str = "\n".join([f"- {f}" for f in long_facts]) if long_facts else "暂无相关记忆"
        return f"""你是一个AI助手,可以调用工具解决问题。
可用工具:{', '.join(tool_names)}
工具详情:{json.dumps(tools_desc, ensure_ascii=False)}
长期记忆中有以下相关事实(仅供参考):
{facts_str}

请按以下JSON格式输出(必须是合法JSON):
{{"thought": "你的思考", "action_name": "工具名或'final_answer'", "action_input": {{参数}}}}
如果不再需要工具,action_name设为'final_answer',action_input为最终回复文本。
"""

关键设计

  • 使用 ActionStep Pydantic 模型保证输出结构化。
  • 反思层:通过额外的 LLM 调用判断观察是否有效,若无效则追加修正指令,增加容错。
  • 步数限制防死循环。

3.5 数据模型(core/schemas.py

代码语言:javascript
复制
from pydantic import BaseModel, Field
from typing import Dict, Any

class ActionStep(BaseModel):
    thought: str = Field(description="当前思考过程")
    action_name: str = Field(description="工具名称或 'final_answer'")
    action_input: Dict[str, Any] = Field(description="工具参数或最终答案文本")

3.6 FastAPI 主入口(main.py

集成所有模块,提供流式 SSE 接口。

代码语言:javascript
复制
from fastapi import FastAPI, Request
from fastapi.responses import StreamingResponse, HTMLResponse
from fastapi.staticfiles import StaticFiles
import json
import asyncio
from core.llm import LLMClient
from core.memory import MemoryManager
from core.planner import AgentPlanner

app = FastAPI(title="AI Agent 全栈系统")

# 初始化全局组件(单例)
llm = LLMClient()
memory = MemoryManager()
planner = AgentPlanner(llm, memory, max_steps=5)

@app.post("/agent/stream")
async def agent_stream(request: Request):
    body = await request.json()
    query = body.get("query", "")
    enable_reflection = body.get("enable_reflection", True)
    
    async def event_generator():
        # 由于 planner.plan 是普通协程,我们希望流式输出中间步骤
        # 这里采用简化方式:先收集完整答案,但为了 SSE 流式,我们可以改造 plan 为生成器
        # 为演示,我们将规划过程包装为生成器(实际可改造 plan 为异步生成器)
        # 这里快速实现:直接用 plan 获取结果,同时模拟流式输出
        result = await planner.plan(query, enable_reflection)
        # 模拟逐字流式输出
        for char in result:
            yield f"data: {json.dumps({'type': 'token', 'content': char})}\n\n"
            await asyncio.sleep(0.02)
        yield f"data: {json.dumps({'type': 'done'})}\n\n"
    
    return StreamingResponse(event_generator(), media_type="text/event-stream")

@app.get("/", response_class=HTMLResponse)
async def serve_frontend():
    with open("static/index.html", "r", encoding="utf-8") as f:
        return f.read()

# 挂载静态文件
app.mount("/static", StaticFiles(directory="static"), name="static")

if __name__ == "__main__":
    import uvicorn
    uvicorn.run(app, host="0.0.0.0", port=8000)

流式优化:此处为演示每字符输出,生产环境建议将 planner.plan 改造为异步生成器,在每个步骤(思考/行动/观察)后 yield 事件,实现真正的流式 ReAct 过程。


3.7 前端极简交互(static/index.html

使用原生 EventSource 接收 SSE,实时渲染。

代码语言:javascript
复制
<!DOCTYPE html>
<html>
<head>
    <meta charset="UTF-8">
    <title>AI Agent 全栈演示</title>
    <style>
        body { font-family: system-ui; max-width: 800px; margin: 40px auto; background: #f5f7fb; }
        #chat { background: white; border-radius: 12px; padding: 20px; min-height: 400px; box-shadow: 0 2px 10px rgba(0,0,0,0.1); }
        .msg { margin: 8px 0; padding: 8px 12px; border-radius: 8px; }
        .user { background: #e0f0ff; text-align: right; }
        .agent { background: #f0f0f0; }
        #input-area { display: flex; gap: 10px; margin-top: 16px; }
        #query { flex: 1; padding: 10px; border: 1px solid #ccc; border-radius: 8px; }
        button { padding: 10px 20px; background: #1a73e8; color: white; border: none; border-radius: 8px; cursor: pointer; }
        button:hover { background: #1557b0; }
        .status { color: #666; font-size: 0.9em; }
    </style>
</head>
<body>
    <h1>🤖 AI Agent 全栈工程师实战</h1>
    <div id="chat"></div>
    <div id="input-area">
        <input id="query" placeholder="输入你的问题,例如:现在北京时间几点?或者计算 (3+5)*2" />
        <button onclick="sendQuery()">发送</button>
    </div>
    <div id="status" class="status"></div>
    <script>
        let currentEventSource = null;
        const chatDiv = document.getElementById('chat');
        const statusDiv = document.getElementById('status');

        function appendMessage(role, content) {
            const div = document.createElement('div');
            div.className = `msg ${role}`;
            div.textContent = content;
            chatDiv.appendChild(div);
            chatDiv.scrollTop = chatDiv.scrollHeight;
        }

        function sendQuery() {
            const query = document.getElementById('query').value.trim();
            if (!query) return;
            document.getElementById('query').value = '';
            appendMessage('user', query);
            statusDiv.textContent = 'Agent 思考中...';

            if (currentEventSource) {
                currentEventSource.close();
            }

            // 使用 fetch + ReadableStream 处理 SSE(因为 EventSource 不支持 POST)
            fetch('/agent/stream', {
                method: 'POST',
                headers: { 'Content-Type': 'application/json' },
                body: JSON.stringify({ query, enable_reflection: true })
            }).then(response => {
                const reader = response.body.getReader();
                const decoder = new TextDecoder();
                let agentMsg = '';
                // 创建 agent 消息容器
                const msgDiv = document.createElement('div');
                msgDiv.className = 'msg agent';
                chatDiv.appendChild(msgDiv);

                function readStream() {
                    return reader.read().then(({ done, value }) => {
                        if (done) {
                            statusDiv.textContent = '✅ 完成';
                            return;
                        }
                        const chunk = decoder.decode(value);
                        const lines = chunk.split('\n');
                        for (const line of lines) {
                            if (line.startsWith('data: ')) {
                                try {
                                    const data = JSON.parse(line.slice(6));
                                    if (data.type === 'token') {
                                        agentMsg += data.content;
                                        msgDiv.textContent = agentMsg;
                                    } else if (data.type === 'done') {
                                        statusDiv.textContent = '✅ 完成';
                                    }
                                } catch(e) { /* ignore */ }
                            }
                        }
                        return readStream();
                    });
                }
                return readStream();
            }).catch(err => {
                statusDiv.textContent = '❌ 错误: ' + err.message;
            });
        }

        // 支持回车
        document.getElementById('query').addEventListener('keydown', (e) => {
            if (e.key === 'Enter') sendQuery();
        });
    </script>
</body>
</html>

四、运行与测试

  1. 创建 .env 文件:

代码语言:javascript
复制
OPENAI_API_KEY=sk-xxxx
OPENAI_BASE_URL=https://api.openai.com/v1   # 可替换为兼容接口
  1. 启动服务:

代码语言:javascript
复制
python main.py
  1. 访问 http://localhost:8000,输入问题。

测试用例

  • 现在北京时间几点? → 触发 get_current_time 工具。
  • 计算 (3+5)*2 - 4 → 触发 calculate 工具。
  • 连续提问,验证短期记忆保留上下文。

五、生产级增强点(部署到腾讯云)

若要将此系统上线,建议做以下加固(腾讯云友好):

5.1 可观测性

接入 腾讯云 CLS(日志服务)Prometheus + Grafana 监控。示例:在 FastAPI 中间件中记录每次请求的 querystepslatency

代码语言:javascript
复制
from prometheus_fastapi_instrumentator import Instrumentator
Instrumentator().instrument(app).expose(app)

5.2 向量数据库替换

将本地 FAISS 替换为 腾讯云向量数据库(Tencent Cloud VectorDB),支持海量记忆和高并发。

代码语言:javascript
复制
# 使用 SDK
from tcvectordb import VectorDBClient
client = VectorDBClient(url=..., key=...)

5.3 异步工具超时控制

对远程 HTTP 工具增加超时和熔断,使用 asyncio.timeouttenacity

5.4 流式 ReAct 中间步骤

改造 planner.plan 为异步生成器,每完成一步 yield 事件,前端可实时展示“思考→行动→观察”过程,提升用户体验。

5.5 部署到腾讯云

  • 使用 Cloud Run(Serverless)CVM + 容器
  • 配置 API 网关 进行限流。
  • 将 LLM 切换为 腾讯混元大模型(兼容 OpenAI API),降低延迟和成本。

六、总结

本文从零构建了一个生产级 AI Agent 全栈系统,技术覆盖:

  • 结构化输出(JSON Schema)保障稳定性
  • 工具注册与远程调用抽象
  • 分层记忆(短期+长期向量检索)
  • ReAct + 反思规划器,提升复杂任务成功率
  • SSE 流式交互,前端实时响应

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

目录
  • AI Agent全栈工程师实战:从零构建生产级智能体系统
    • 引言
    • 一、系统架构全景
    • 二、环境与依赖
    • 三、核心模块实现
      • 3.1 统一 LLM 客户端(core/llm.py)
      • 3.2 工具注册中心(core/tools.py)
      • 3.3 分层记忆系统(core/memory.py)
      • 3.4 规划器:ReAct + 反思(core/planner.py)
      • 3.5 数据模型(core/schemas.py)
      • 3.6 FastAPI 主入口(main.py)
      • 3.7 前端极简交互(static/index.html)
    • 四、运行与测试
    • 五、生产级增强点(部署到腾讯云)
      • 5.1 可观测性
      • 5.2 向量数据库替换
      • 5.3 异步工具超时控制
      • 5.4 流式 ReAct 中间步骤
      • 5.5 部署到腾讯云
    • 六、总结
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档