ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

从零构建生产级记忆型AI Agent:AgentScope+DDD+SSE实战

从零构建生产级记忆型AI Agent:AgentScope+DDD+SSE实战 1. 为什么我要从零手搓一个记忆型 AI Agent去年下半年开始我陆续接手了几个 AI Agent 相关的项目从最简单的问答机器人到带工具调用的复杂工作流都摸了一遍。踩坑最多的不是模型选型也不是 Prompt 调优而是记忆管理这件事。大部分教程教你搭一个 Agent跑通一次对话就结束了但真实的生产环境里用户今天跟你说的话明天再聊时你得记得用户上周提到的偏好这周做推荐时你得用上。没有记忆的 Agent本质上就是个高级一点的搜索框。我试过几种方案把历史对话全塞进上下文窗口token 消耗爆炸不说稍微长一点的对话就开始丢信息用简单的向量库做检索召回质量忽高忽低经常答非所问自己手写一套记忆管理逻辑代码写到三千行的时候发现维护成本已经失控了。后来接触到AgentScope这个框架它的设计思路让我眼前一亮——用DDD领域驱动设计的思想来组织 Agent 的各个模块把记忆、工具、规划、执行拆成清晰的领域边界再通过SSEServer-Sent Events做流式输出整体架构非常干净。这篇文章我会从零开始完整拆解一个生产级记忆型 AI Agent 的构建过程。涉及的核心技术点包括 AgentScope 框架的使用、DDD 分层架构设计、SSE 流式通信、MCP 协议集成、记忆系统的分层实现等。适合有一定编程基础、想深入理解 AI Agent 内部机制的开发者也适合正在做 AI Agent 中台选型的技术负责人参考。我不会只讲概念每个关键环节都会给出可运行的代码和参数配置你跟着做就能跑起来。2. 整体架构设计与技术选型思路2.1 为什么选 AgentScope 而不是自己造轮子市面上 AI Agent 框架不少LangChain、AutoGen、CrewAI 各有各的定位。我最终选 AgentScope 主要看中三点第一它的抽象层次恰到好处不像 LangChain 那样封装过厚导致调试困难也不像裸调 API 那样什么都要自己写第二原生支持多 Agent 协作消息传递机制设计得很自然后面扩展多 Agent 场景不用重构第三中文文档完善社区活跃度在国内框架里算第一梯队遇到问题能快速找到答案。AgentScope 的核心抽象就几个Agent基类、Msg消息对象、Memory记忆接口、Toolkit工具集。这种极简的抽象设计让整个框架的学习曲线很平缓我大概花了一个下午就把核心 API 摸清楚了。对比之下LangChain 的 Chain、Runnable、AgentExecutor 那一套概念体系光理清关系就得花两三天。2.2 DDD 分层架构在 Agent 项目中的落地DDD 听起来很玄乎但在 Agent 项目里其实特别适用。我把它拆成四层接口层Interface负责对外暴露 API处理 SSE 连接、请求参数校验、鉴权等应用层Application编排用例比如处理用户消息这个用例要调用记忆检索、Agent 推理、工具执行、结果持久化领域层Domain核心业务逻辑包括 Agent 实体、记忆领域服务、工具领域服务基础设施层Infrastructure具体的技术实现比如向量数据库客户端、LLM API 封装、Redis 缓存这样分层的好处是当我想把向量库从 Chroma 换成 Milvus 时只需要改基础设施层的实现领域层和应用层的代码一行不用动。我在实际项目中就做过这种切换整个过程不到半天。2.3 SSE 还是 WebSocket流式通信的选型对比AI Agent 的输出是逐 token 生成的必须用流式通信才能让用户有实时响应的体验。SSE 和 WebSocket 都能做流式但适用场景不同对比维度SSEWebSocket通信方向单向服务端到客户端双向协议基于 HTTP独立协议自动重连浏览器原生支持需手动实现实现复杂度低中适用场景流式输出、通知推送实时双向交互AI Agent 的典型场景是用户发一条消息Agent 流式返回结果这是典型的单向推送SSE 完全够用。而且 SSE 基于 HTTP穿透代理和防火墙的能力更强部署时少很多麻烦。我实测下来SSE 在弱网环境下的表现也比 WebSocket 稳定因为浏览器会自动重连。注意SSE 有一个坑默认的 idle timeout 通常是 60 秒如果 Agent 推理时间较长比如调用多个工具连接可能被中间层断开报错stream disconnected before completion: idle timeout waiting for SSE。解决方案是每隔 15-30 秒发送一个心跳注释行: heartbeat\n\n保持连接活跃。2.4 MCP 协议让 Agent 接入外部工具的标准方式MCPModel Context Protocol是 Anthropic 提出的开放协议用于标准化 AI 模型与外部工具、数据源的交互方式。你可以把它理解成AI 世界的 USB 接口——只要工具实现了 MCP Server任何支持 MCP 的 Agent 都能直接调用不用为每个工具写适配代码。AgentScope 对 MCP 的支持很完善通过McpClient就能接入各种 MCP Server。我目前接入过的 MCP Server 包括 Playwright MCP浏览器自动化、文件系统 MCP、数据库 MCP 等配置都很简单。MCP 的核心价值在于解耦工具开发者专注实现工具能力Agent 开发者专注编排逻辑两边通过标准协议通信。3. 记忆系统的核心设计与实现细节3.1 三层记忆架构短期、长期、工作记忆记忆是这篇文章的核心也是生产级 Agent 和玩具 Agent 的分水岭。我设计的记忆系统分三层短期记忆Short-term Memory存储当前会话的对话历史用滑动窗口 摘要压缩的方式管理。窗口大小设为最近 10 轮对话超出部分用 LLM 压缩成摘要。这样既保留了近期上下文又不会让 token 无限增长。长期记忆Long-term Memory存储跨会话的重要信息比如用户偏好、历史事实、关键决策。用向量数据库存储检索时根据当前 query 做语义相似度匹配。我用的 embedding 模型是text-embedding-3-small维度 1536检索 top-k 设为 5。工作记忆Working Memory存储当前任务执行过程中的中间状态比如工具调用的中间结果、推理链的中间步骤。这部分用 Redis 存储设置 TTL 为 1 小时任务完成后自动清理。三层记忆的读写时机不同短期记忆每轮对话都读写长期记忆在对话开始时检索、对话结束时写入工作记忆在任务执行过程中频繁读写。3.2 记忆的写入策略什么该记什么不该记这是很多人忽略的问题。如果什么都往长期记忆里塞检索质量会急剧下降。我的策略是用户显式声明的事实比如我是做后端的、我偏好 Python这类信息直接写入Agent 推理出的结论比如用户多次询问 Java 相关问题推断出用户可能在做 Java 项目这类信息加置信度标记后写入临时性信息比如今天天气怎么样这类信息不写入长期记忆敏感信息比如密码、身份证号绝对不写入写入时还要做去重和冲突检测。我遇到过用户先说我用 MySQL后来说我们迁移到 PostgreSQL 了如果不做冲突检测检索时会同时召回两条矛盾的信息。解决方案是给每条记忆加时间戳和版本号检索时优先返回最新的。3.3 记忆检索的优化从关键词到混合检索纯向量检索的问题在于它对精确匹配不敏感。比如用户问我上次说的那个项目向量检索可能召回一堆不相关的项目信息。我的优化方案是混合检索先用 BM25 做关键词检索召回 top-20再用向量检索召回 top-20用 RRFReciprocal Rank Fusion算法融合两个结果最后用 Cross-Encoder 做精排返回 top-5这套流程下来检索准确率比纯向量检索提升了大概 30%。RRF 的公式很简单score Σ 1/(k rank_i)k 通常取 60。3.4 记忆的遗忘机制不是所有记忆都值得保留人脑会遗忘Agent 也应该会。我设计了一个基于重要性评分 时间衰减的遗忘机制重要性评分 访问频率 × 0.4 用户反馈 × 0.4 信息密度 × 0.2 最终得分 重要性评分 × exp(-λ × 天数)λ 取 0.01意味着 100 天后记忆权重衰减到初始的 37%。当最终得分低于阈值我设为 0.1时记忆被归档到冷存储不再参与检索。这样既控制了检索范围又保留了历史数据以备不时之需。4. 从零搭建完整实操流程4.1 环境准备与依赖安装先建一个 Python 3.10 的虚拟环境然后安装核心依赖pip install agentscope0.1.0 pip install fastapi uvicorn sse-starlette pip install chromadb sentence-transformers pip install redis pydantic pip install mcpAgentScope 的版本建议用 0.1.0 以上早期版本在 MCP 支持上有一些 bug。Redis 用于工作记忆ChromaDB 用于长期记忆的向量存储sse-starlette 提供了比原生 FastAPI 更好用的 SSE 封装。4.2 领域层定义 Agent 和 Memory 的核心接口先定义领域层的抽象接口这是 DDD 的核心from abc import ABC, abstractmethod from typing import List, Optional from pydantic import BaseModel from datetime import datetime class MemoryItem(BaseModel): content: str timestamp: datetime importance: float 0.5 metadata: dict {} class MemoryRepository(ABC): abstractmethod async def save(self, item: MemoryItem) - str: pass abstractmethod async def retrieve(self, query: str, top_k: int 5) - List[MemoryItem]: pass abstractmethod async def forget(self, item_id: str) - None: pass这个接口定义好后具体用 Chroma 还是 Milvus 实现领域层完全不关心。这就是 DDD 的价值——依赖倒置。4.3 基础设施层向量库与 Redis 的具体实现长期记忆用 ChromaDB 实现import chromadb from chromadb.utils import embedding_functions class ChromaMemoryRepository(MemoryRepository): def __init__(self, persist_dir: str ./memory_db): self.client chromadb.PersistentClient(pathpersist_dir) self.embedding_fn embedding_functions.SentenceTransformerEmbeddingFunction( model_nameBAAI/bge-small-zh-v1.5 ) self.collection self.client.get_or_create_collection( namelong_term_memory, embedding_functionself.embedding_fn, metadata{hnsw:space: cosine} ) async def save(self, item: MemoryItem) - str: item_id fmem_{int(item.timestamp.timestamp())}_{hash(item.content) % 10000} self.collection.add( ids[item_id], documents[item.content], metadatas[{ timestamp: item.timestamp.isoformat(), importance: item.importance, **item.metadata }] ) return item_id async def retrieve(self, query: str, top_k: int 5) - List[MemoryItem]: results self.collection.query( query_texts[query], n_resultstop_k ) items [] for i, doc in enumerate(results[documents][0]): meta results[metadatas][0][i] items.append(MemoryItem( contentdoc, timestampdatetime.fromisoformat(meta[timestamp]), importancemeta.get(importance, 0.5), metadatameta )) return items这里 embedding 模型我选的是bge-small-zh-v1.5中文效果好模型体积小约 100MB推理速度快。如果你的场景以英文为主可以换成all-MiniLM-L6-v2。工作记忆用 Redis 实现import redis.asyncio as redis import json class RedisWorkingMemory: def __init__(self, redis_url: str redis://localhost:6379): self.client redis.from_url(redis_url, decode_responsesTrue) async def set_state(self, session_id: str, key: str, value: dict, ttl: int 3600): full_key fwm:{session_id}:{key} await self.client.setex(full_key, ttl, json.dumps(value, ensure_asciiFalse)) async def get_state(self, session_id: str, key: str) - Optional[dict]: full_key fwm:{session_id}:{key} data await self.client.get(full_key) return json.loads(data) if data else None4.4 应用层编排记忆检索与 Agent 推理应用层的核心用例是处理用户消息流程是检索长期记忆 → 组装上下文 → 调用 Agent → 流式返回 → 写入记忆。from agentscope.agent import ReActAgent from agentscope.model import OpenAIChatModel from agentscope.memory import InMemoryMemory class ChatUseCase: def __init__( self, long_term_repo: MemoryRepository, working_memory: RedisWorkingMemory, agent: ReActAgent ): self.long_term_repo long_term_repo self.working_memory working_memory self.agent agent async def stream_chat(self, session_id: str, user_input: str): # 1. 检索长期记忆 relevant_memories await self.long_term_repo.retrieve(user_input, top_k5) memory_context \n.join([m.content for m in relevant_memories]) # 2. 组装增强后的输入 enhanced_input user_input if memory_context: enhanced_input f[相关记忆]\n{memory_context}\n\n[用户输入]\n{user_input} # 3. 流式调用 Agent full_response async for chunk in self.agent.stream(enhanced_input): full_response chunk yield chunk # 4. 异步写入长期记忆不阻塞响应 await self._save_memory_async(user_input, full_response) async def _save_memory_async(self, user_input: str, response: str): # 用 LLM 判断是否值得记忆 importance await self._evaluate_importance(user_input, response) if importance 0.6: await self.long_term_repo.save(MemoryItem( contentf用户: {user_input}\n助手: {response}, timestampdatetime.now(), importanceimportance ))4.5 接口层SSE 流式接口的实现接口层用 FastAPI sse-starlettefrom fastapi import FastAPI, Request from sse_starlette.sse import EventSourceResponse import asyncio app FastAPI() app.post(/chat/stream) async def chat_stream(request: Request, body: ChatRequest): async def event_generator(): try: async for chunk in chat_use_case.stream_chat(body.session_id, body.message): if await request.is_disconnected(): break yield {event: message, data: chunk} yield {event: done, data: [DONE]} except Exception as e: yield {event: error, data: str(e)} return EventSourceResponse( event_generator(), ping15 # 每15秒发送心跳防止 idle timeout )sse-starlette的ping参数会自动发送心跳注释行解决了前面提到的 idle timeout 问题。这个参数我踩过坑默认值是 15 秒如果你的 Agent 推理时间经常超过 30 秒建议改成 10 秒。4.6 MCP 工具集成让 Agent 拥有外部能力AgentScope 集成 MCP 的方式很简洁from agentscope.mcp import McpClient from agentscope.tool import Toolkit async def setup_mcp_tools(): toolkit Toolkit() # 接入文件系统 MCP Server fs_client McpClient( namefilesystem, commandnpx, args[-y, modelcontextprotocol/server-filesystem, /tmp/agent_workspace] ) await fs_client.connect() toolkit.register_mcp_client(fs_client) # 接入 Playwright MCP Server浏览器自动化 playwright_client McpClient( nameplaywright, commandnpx, args[-y, playwright/mcplatest] ) await playwright_client.connect() toolkit.register_mcp_client(playwright_client) return toolkitMCP Server 的启动方式通常是npx或uvx需要本地有 Node.js 或 Python 环境。生产环境建议把 MCP Server 部署成独立服务通过 SSE 或 stdio 通信避免和主进程耦合。5. 常见问题与排查技巧实录5.1 SSE 连接频繁断开怎么办这是最高频的问题。排查顺序检查心跳配置确认ping参数小于中间层Nginx、负载均衡器的 idle timeout。Nginx 默认 60 秒建议 ping 设为 15 秒。检查 Nginx 配置需要加proxy_buffering off;和proxy_cache off;否则 Nginx 会缓冲 SSE 流导致客户端收不到实时数据。检查响应头确保有X-Accel-Buffering: no这是给 Nginx 的信号告诉它不要缓冲这个响应。检查客户端浏览器 EventSource 默认会自动重连但如果服务端返回非 200 状态码重连会失败。确保异常时也返回 200通过 event 类型区分错误。5.2 记忆检索召回质量差怎么优化我整理了一个排查表症状可能原因解决方案召回内容完全不相关embedding 模型不匹配语言中文场景换 bge-small-zh召回内容相关但不够精确纯向量检索精度不足引入 BM25 混合检索召回内容重复写入时未去重写入前做相似度检测召回内容过时未做时间衰减引入时间衰减因子召回数量太少top_k 设置过小增大 top_k 后精排5.3 Agent 推理超时怎么处理生产环境必须设置超时。我的配置是单次 LLM 调用超时 30 秒整个 Agent 推理链路超时 120 秒。超时后返回部分结果 提示信息而不是直接报错。实现方式是用asyncio.wait_for包裹推理调用try: result await asyncio.wait_for( self.agent.stream(enhanced_input), timeout120 ) except asyncio.TimeoutError: yield 抱歉处理时间过长请稍后重试或简化您的问题。5.4 MCP Server 连接失败的排查MCP 连接失败通常是三个原因命令不存在检查 npx/uvx 是否安装、参数错误检查 MCP Server 的启动参数、权限不足比如文件系统 MCP 需要目录读写权限。排查时先手动在终端运行 MCP Server 的启动命令确认能正常启动后再集成到代码里。5.5 记忆写入导致响应变慢记忆写入如果同步执行会显著增加响应时间。我的做法是异步写入用asyncio.create_task把写入操作丢到后台不阻塞主响应流。但要注意异步任务需要自己处理异常否则失败会静默丢失。我加了一个写入队列 重试机制失败的任务会重新入队最多重试 3 次。6. 生产部署的关键考量6.1 性能优化从单机到分布式单机部署时Agent 推理和记忆检索在同一个进程里简单但性能有限。上生产后我做了几个优化Agent 推理服务独立部署用 FastAPI 起多个 worker前面挂 Nginx 做负载均衡记忆检索服务独立部署向量库单独部署通过 gRPC 通信减少网络开销Redis 集群工作记忆用 Redis Cluster避免单点故障连接池LLM API 调用、向量库查询都用连接池避免频繁建连实测下来这套架构在 4 核 8G 的机器上能支撑约 50 QPS 的并发对话。6.2 可观测性日志、指标、追踪生产环境必须能定位问题。我接入了三个东西结构化日志用structlog输出 JSON 格式日志包含 session_id、trace_id、耗时等字段Prometheus 指标暴露 QPS、P99 延迟、记忆命中率、工具调用成功率等指标OpenTelemetry 追踪每个请求生成 trace_id串联 Agent 推理、记忆检索、工具调用各环节有了这些线上出问题能快速定位到具体环节。我遇到过一次记忆检索变慢导致整体响应超时通过追踪发现是向量库的 HNSW 索引参数不合理调整ef_search后恢复正常。6.3 成本控制Token 消耗的优化LLM 调用是最大的成本项。我做了几个优化记忆压缩短期记忆超过 10 轮后用便宜的小模型压缩成摘要减少上下文长度缓存相同 query 的 embedding 结果缓存到 Redis避免重复计算模型分级简单任务用便宜模型如 gpt-4o-mini复杂任务才用贵模型流式截断设置 max_tokens 上限避免模型生成过长内容这套组合拳下来token 成本降低了约 60%。6.4 安全与合规数据隔离与访问控制多租户场景下记忆必须严格隔离。我的做法是向量库按租户分 collection每个租户一个 collection物理隔离Redis key 加租户前缀wm:{tenant_id}:{session_id}:{key}API 层鉴权每个请求校验 token提取 tenant_id 后注入上下文敏感信息过滤写入记忆前用正则 LLM 双重检测过滤手机号、身份证号等注意记忆数据涉及用户隐私必须提供删除接口。我实现了遗忘功能用户可以要求删除特定记忆系统会从向量库和 Redis 中彻底清除。7. 我踩过的坑和几条实在建议第一个坑是过度依赖向量检索。我一开始所有记忆都用向量检索结果发现精确匹配场景比如用户问我上次说的那个项目叫什么召回质量很差。后来引入 BM25 混合检索才解决。如果你也在做记忆系统建议一开始就上混合检索别走弯路。第二个坑是记忆写入没有做重要性评估。早期我把所有对话都写入长期记忆结果向量库膨胀到几十万条检索质量断崖式下降。后来加了重要性评分只写入评分高于 0.6 的记忆检索准确率立刻回升。第三个坑是SSE 心跳没配置。上线第一天就遇到大量连接断开排查了半天才发现是 Nginx 的 idle timeout。现在我的标准配置是SSE ping 15 秒Nginxproxy_read_timeout300 秒proxy_buffering off。最后一个建议别急着上多 Agent。我见过很多项目一上来就搞多 Agent 协作结果调试成本爆炸。先把单 Agent 记忆 工具这套跑通稳定运行一段时间后再根据实际需求扩展多 Agent。AgentScope 的多 Agent 能力很强但前提是你的单 Agent 基础扎实。这套架构我目前跑了三个多月支撑了大概几十万次对话整体稳定性不错。后续我打算在记忆检索上引入 rerank 模型进一步优化召回质量也在考虑把 MCP 工具调用做成插件市场让业务方自己接入工具。如果你也在做类似的事情欢迎交流踩坑经验。
返回列表