拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

从零手搓生产级记忆型AI Agent:AgentScope+DDD+SSE+MCP实战

从零手搓生产级记忆型AI Agent:AgentScope+DDD+SSE+MCP实战

1. 为什么我要从零手搓一个记忆型 AI Agent

先说结论:市面上大部分号称“AI Agent”的东西,本质上就是给大模型套了个 while 循环,再挂几个工具调用,跑完一轮对话就失忆。你问它昨天聊过什么,它一脸茫然;你让它接着上周的任务继续干,它得让你从头把背景再讲一遍。这种“金鱼记忆”的 Agent,拿来做 demo 演示还行,真要放到生产环境里跑业务,基本活不过三天。

我这次要聊的,就是怎么用AgentScope这套框架,从零搭一个真正带长期记忆、能扛住生产流量、支持SSE 流式输出、还能通过MCP 协议接入外部工具的生产级 AI Agent。关键词里提到的 AgentScope、AI Agent、DDD、SSE、MCP,这五个东西不是随便凑在一起的,它们分别对应了这类系统的五个核心命题:框架选型、智能体行为建模、领域驱动设计、实时通信、工具生态接入。

这篇文章适合谁看?如果你已经写过几个 LangChain 的 demo,但一到“怎么让 Agent 记住用户偏好”“怎么让前端实时看到 Agent 的思考过程”“怎么把公司内部系统接进来当工具”这些问题就卡壳,那这篇就是写给你的。如果你是完全的新手,也没关系,我会把每个概念用生活化的类比讲清楚,保证你能跟着一步步复现。

我踩过的坑不少,比如最开始用轮询文件变化的方式做记忆持久化,结果并发一上来文件锁直接炸了;又比如 SSE 连接空闲超时被网关掐断,前端一直报stream disconnected before completion: idle timeout waiting for SSE,排查了大半天才发现是心跳没做。这些经验我都会在下面掰开揉碎讲。

2. 整体架构设计与技术选型思路

2.1 为什么是 AgentScope 而不是别的框架

选框架这件事,我的判断标准就三条:消息传递模型是否清晰、是否原生支持多智能体协作、记忆和工具的扩展点是否开放。AgentScope 在这三点上做得比较扎实。它的核心抽象是Message(消息)和Agent(智能体),所有交互都通过消息驱动,这跟现实世界里人跟人协作的方式是一致的——你发消息,我回消息,消息里带着内容和元数据。

对比一下,有些框架把 Agent 设计成一个巨大的“链”,每一步都耦合在一起,想加个记忆模块得改一堆代码。AgentScope 的设计更像搭积木,记忆、工具、规划器都是可以插拔的组件。这对生产环境太重要了,因为生产环境的需求是会变的,今天要加个向量检索记忆,明天要换个工具调用协议,框架如果不支持热插拔,你就得推倒重来。

还有一个现实考量:AgentScope 有比较完整的中文文档和社区案例,关键词里“agentscope中文文档”“agentscope教程”出现频率很高,说明中文开发者群体对它接受度不错。遇到问题能搜到答案,这在赶工期的时候能救命。

2.2 DDD 领域驱动设计怎么落到 Agent 上

很多人觉得 DDD 是后端业务系统才用的东西,跟 AI Agent 没关系。我的看法恰恰相反:Agent 越复杂,越需要 DDD 来划清边界。

举个具体例子。一个生产级 Agent 通常要处理这些事:理解用户意图、检索记忆、调用工具、生成回复、记录日志。如果你把这些逻辑全塞在一个Agent.run()方法里,代码会迅速膨胀到几千行,改一处崩三处。用 DDD 的思路,我会这样划分:

  • 对话域(Conversation Domain):负责消息的接收、路由、上下文组装,核心实体是Session和Message。
  • 记忆域(Memory Domain):负责短期上下文和长期记忆的存取,核心实体是MemoryItem,包含内容、时间戳、重要度、向量表示。
  • 工具域(Tool Domain):负责工具的注册、发现、调用和结果解析,核心实体是ToolSpec。
  • 编排域(Orchestration Domain):负责决定下一步做什么,是调工具还是直接回复,核心是规划逻辑。

这样划分之后,每个域可以独立测试、独立演进。比如我想把记忆从“全量塞进上下文”改成“向量检索 Top-K”,只需要改记忆域,其他域完全不受影响。这就是 DDD 带来的隔离价值。

2.3 SSE 与 WebSocket 的取舍

实时通信这块,关键词里同时出现了 SSE 和 WebSocket,还有“react + sse/websocket 轮询文件变化”这种组合。我的选择是:Agent 的输出流用 SSE,双向控制通道用 WebSocket。

为什么这么分?因为 Agent 生成回复是一个单向的、服务器推送到客户端的过程,SSE 天然适合这种场景,它基于 HTTP,实现简单,浏览器原生支持EventSource,断线重连也有标准机制。而 WebSocket 是全双工的,适合客户端频繁给服务器发指令的场景,比如用户中途打断 Agent、切换工具、调整参数。

如果你只用 WebSocket 做所有事,当然也行,但复杂度会上去。如果你只用 SSE,那客户端就没法主动发消息给服务器(SSE 是单向的),得另开 HTTP 接口。所以我的建议是:流式输出走 SSE,控制指令走普通 HTTP POST 或 WebSocket,各司其职。

这里有个坑必须提前说:SSE 连接很容易被中间的网关、负载均衡器因为空闲超时掐断,报错就是stream disconnected before completion: idle timeout waiting for SSE。解决办法是服务端定期发送心跳注释行(以冒号开头的行),比如每 15 秒发一个: heartbeat\n\n,保持连接活跃。

2.4 MCP 协议:工具生态的“USB 接口”

MCP 是什么?你可以把它理解成 AI 世界的 USB 接口。以前每个工具都要写一套自己的接入代码,A 工具的调用方式和 B 工具完全不一样,Agent 想用新工具就得改代码。MCP 定义了一套标准协议,只要工具实现了 MCP Server,任何支持 MCP 的 Agent 都能直接调用,不用改一行代码。

关键词里出现了大量 MCP 相关的词:playwright mcp、burpsuite mcp、blender mcp、unity mcp、chrome devtools mcp、同花顺 mcp……这说明 MCP 生态正在快速膨胀,从浏览器自动化到 3D 建模到金融数据,各种工具都在往 MCP 上靠。对 Agent 开发者来说,这是巨大的红利——你不需要自己写工具适配层,接上 MCP 就能用现成的。

我的 Agent 里,工具域就是围绕 MCP 设计的:启动时连接配置好的 MCP Server 列表,拉取每个 Server 提供的工具清单,注册到工具注册表里。Agent 需要调工具时,通过 MCP 协议发请求,拿到结果再继续推理。

3. 核心模块拆解与实操要点

3.1 记忆模块:短期上下文与长期记忆的分层设计

记忆是这类 Agent 的灵魂,也是最容易做砸的地方。我的设计分两层:

短期记忆(Short-term Memory)就是当前会话的上下文窗口。它不需要持久化到数据库,放在内存里就行,但要控制长度。我的做法是维护一个滑动窗口,保留最近 N 轮对话,N 根据模型上下文长度动态计算。比如模型支持 128K token,我预留 20K 给系统提示和工具定义,剩下 108K 按平均每轮 500 token 算,大概能放 200 轮。但实际不会放这么多,因为太长的上下文会让模型注意力分散,我一般控制在 30 到 50 轮。

长期记忆(Long-term Memory)才是关键。它要解决的是“跨会话记住用户”的问题。我的实现方案是:

  1. 每轮对话结束后,异步提取值得记住的信息(用户偏好、事实性知识、任务状态),生成MemoryItem。
  2. 对MemoryItem的内容做向量化,存进向量数据库。
  3. 下一轮对话开始时,用当前用户输入去检索相关记忆,Top-K 条注入到上下文里。

这里有个细节:不是所有对话都值得记。如果用户说“今天天气不错”,这没必要记。如果用户说“我习惯用 Python 而不是 Java”,这就值得记。我用一个轻量的分类逻辑来判断,可以是一个小模型,也可以是一组规则加关键词匹配。生产环境里我倾向于用规则打底加模型兜底,因为规则快且可控,模型处理边界情况。

注意:向量检索的 Top-K 不是越大越好。K 太大,注入的无关记忆会干扰模型;K 太小,可能漏掉关键信息。我的经验值是 K=5 到 8,配合一个相似度阈值(比如 0.75),低于阈值的直接丢弃。

3.2 消息流转与 SSE 流式输出实现

Agent 的消息流转是这样的:用户输入 -> 组装上下文(系统提示 + 短期记忆 + 检索到的长期记忆 + 工具定义)-> 调用模型 -> 模型可能返回工具调用请求 -> 执行工具 -> 把工具结果再喂给模型 -> 模型生成最终回复。

这个过程中,前端最关心的是“Agent 现在在干嘛”。如果等全部跑完再返回,用户会盯着空白屏幕等好几秒,体验很差。所以要用 SSE 把中间过程流式推出去。

我的 SSE 事件设计是这样的:

事件类型数据内容触发时机
thinkingAgent 的思考文本模型开始生成推理内容时
tool_call工具名和参数Agent 决定调用工具时
tool_result工具返回结果摘要工具执行完成时
message最终回复的增量文本模型生成回复时逐 token 推送
done结束标记整个流程完成时
error错误信息任何环节出错时

服务端用 Python 的sse-starlette或者自己基于StreamingResponse实现都行。关键点是每个事件之间要 flush,不能攒着一起发,否则就失去流式的意义了。

async def event_generator(request): async for event in agent.stream_run(user_input): yield { "event": event.type, "data": json.dumps(event.payload, ensure_ascii=False) } await asyncio.sleep(0) # 让出控制权,确保及时 flush

前端用EventSource接收:

const es = new EventSource('/api/agent/stream?session_id=xxx'); es.addEventListener('message', (e) => { const data = JSON.parse(e.data); appendToChat(data.text); }); es.addEventListener('tool_call', (e) => { showToolIndicator(JSON.parse(e.data)); }); es.addEventListener('done', () => es.close());

实操心得:SSE 的EventSource默认只支持 GET 请求,如果你的用户输入很长,URL 长度可能超限。解决办法是把用户输入先 POST 到一个接口存起来,返回一个 token,然后 SSE 连接带上这个 token 去取。或者直接用fetch加ReadableStream手动解析 SSE 格式,这样就能用 POST 了。

3.3 MCP 工具接入的完整流程

接入 MCP 工具分三步:发现、注册、调用。

发现阶段,Agent 启动时读取配置文件里的 MCP Server 列表,每个 Server 有地址和认证信息。通过 MCP 协议的tools/list方法拉取工具清单。每个工具包含名称、描述、参数 schema。

注册阶段,把拉到的工具清单转换成 Agent 内部统一的ToolSpec格式,存进工具注册表。这里要注意工具名的冲突处理,不同 Server 可能有同名工具,我的做法是加前缀,比如playwright__click和burp__click。

调用阶段,Agent 决定用某个工具时,通过 MCP 的tools/call方法发请求,带上工具名和参数,拿到结果后解析成文本喂回模型。

class MCPToolAdapter: def __init__(self, server_config): self.server_url = server_config['url'] self.token = server_config['token'] self.tools = {} async def discover(self): resp = await self._rpc('tools/list', {}) for tool in resp['tools']: self.tools[tool['name']] = tool async def call(self, name, arguments): return await self._rpc('tools/call', { 'name': name, 'arguments': arguments })

注意:MCP Server 的日志管理是个容易被忽视的点。生产环境里,你需要把每个工具调用的入参、出参、耗时都记下来,方便排查问题。我建议在 Adapter 层统一加日志中间件,而不是在每个工具里单独加。

3.4 会话状态管理与并发控制

生产级 Agent 必须处理并发。同一个用户可能开多个标签页,同一个会话可能同时有多个请求进来。如果不做控制,会出现消息错乱、记忆覆盖的问题。

我的方案是会话级锁 + 消息队列。每个session_id对应一个异步锁,同一会话的请求串行处理,不同会话并行。这样既保证了会话内的一致性,又不影响整体吞吐。

会话状态存 Redis,包含:当前上下文窗口、最近检索到的记忆 ID、正在执行的任务状态。Redis 的过期时间设成 24 小时,超过就清理,避免内存无限增长。

async def handle_message(session_id, user_input): lock = await get_session_lock(session_id) async with lock: state = await load_session_state(session_id) context = build_context(state, user_input) async for event in agent.stream_run(context): yield event await save_session_state(session_id, state)

4. 从零搭建的完整实操流程

4.1 环境准备与依赖安装

我用的技术栈是 Python 3.11 + AgentScope + FastAPI + Redis + PostgreSQL(带 pgvector 扩展)。Python 3.11 是因为它在异步性能和类型提示上比 3.10 有改进,而且 AgentScope 对 3.11 支持最好。

python -m venv venv source venv/bin/activate pip install agentscope fastapi uvicorn redis psycopg2-binary pgvector sse-starlette httpx

数据库初始化:

CREATE EXTENSION IF NOT EXISTS vector; CREATE TABLE memory_items ( id SERIAL PRIMARY KEY, session_id VARCHAR(64), content TEXT, embedding vector(1536), importance FLOAT, created_at TIMESTAMP DEFAULT NOW() ); CREATE INDEX ON memory_items USING ivfflat (embedding vector_cosine_ops);

实操心得:pgvector 的索引类型选ivfflat还是hnsw,取决于数据量。数据量小于 100 万条用ivfflat够用,构建快;超过 100 万条建议hnsw,查询更快但构建慢、占内存多。我一开始用ivfflat,后来数据涨到 200 万条查询明显变慢,换成hnsw后恢复。

4.2 Agent 核心类的编写

Agent 的核心是一个状态机,我用 AgentScope 的AgentBase派生:

class MemoryAgent(AgentBase): def __init__(self, model, memory, tools, config): super().__init__() self.model = model self.memory = memory self.tools = tools self.config = config async def reply(self, msg): # 1. 检索长期记忆 relevant = await self.memory.retrieve(msg.content, top_k=5) # 2. 组装上下文 context = self._build_context(msg, relevant) # 3. 推理循环 while True: response = await self.model(context) if response.has_tool_call: yield ToolCallEvent(response.tool_call) result = await self.tools.call(response.tool_call) yield ToolResultEvent(result) context.append(result) else: yield MessageEvent(response.text) break # 4. 异步写入长期记忆 asyncio.create_task(self.memory.extract_and_store(msg, response))

这个循环就是 Agent 的“思考-行动”循环。模型如果决定调工具,就执行工具把结果加回上下文继续想;如果决定直接回复,就输出并结束。

4.3 SSE 接口的完整实现

FastAPI 里实现 SSE 接口:

from fastapi import FastAPI, Request from sse_starlette.sse import EventSourceResponse app = FastAPI() @app.get("/api/agent/stream") async def stream(request: Request, session_id: str, message: str): async def generator(): try: async for event in agent.stream_run(session_id, message): if await request.is_disconnected(): break yield { "event": event.type, "data": json.dumps(event.payload, ensure_ascii=False) } except Exception as e: yield {"event": "error", "data": str(e)} finally: yield {"event": "done", "data": "{}"} return EventSourceResponse( generator(), ping=15 # 每15秒发心跳,防止空闲超时 )

ping=15这个参数就是解决idle timeout waiting for SSE的关键,它会自动发送心跳注释行。

4.4 MCP Server 的配置与联调

配置文件mcp_servers.yaml:

servers: - name: playwright url: http://localhost:3001/mcp token: ${PLAYWRIGHT_MCP_TOKEN} - name: filesystem url: http://localhost:3002/mcp token: ${FS_MCP_TOKEN}

启动时加载配置,逐个连接并发现工具。联调阶段我建议先用一个简单的 MCP Server 测试,比如官方的 filesystem server,确认发现和调用都通了,再接复杂的。

注意:MCP Server 的连接要加超时和重试。生产环境里 Server 可能临时不可用,不能让一个工具挂掉导致整个 Agent 卡死。我的做法是每个工具调用设 30 秒超时,失败后返回错误信息给模型,让模型决定是重试还是换方案。

5. 常见问题排查与避坑经验

5.1 SSE 连接频繁断开

这是最高频的问题,报错通常是stream disconnected before completion: idle timeout waiting for SSE。原因有三个:一是没发心跳,二是网关超时设置太短,三是服务端处理太慢导致连接被判定为空闲。

排查顺序:先确认服务端有没有发心跳(用 curl 直接连 SSE 接口看有没有: ping行),再看网关的超时配置(Nginx 的proxy_read_timeout默认 60 秒,要调大),最后看 Agent 处理逻辑有没有阻塞事件循环的地方。

5.2 记忆检索不准确

表现是 Agent 答非所问,或者重复问已经告诉过它的信息。原因可能是向量模型不适合中文、相似度阈值设得不对、或者记忆提取环节漏掉了关键信息。

我的排查方法:把检索到的记忆和用户输入打印出来,人工判断相关性。如果检索结果明显不相关,先换向量模型试试;如果检索结果相关但 Agent 没用上,那是上下文组装的问题,检查记忆有没有正确注入到 prompt 里。

5.3 工具调用参数错误

模型生成的工具参数经常不符合 schema,比如该传数字传了字符串,该传数组传了单个值。解决办法是在工具定义里把参数描述写清楚,给出示例。另外可以在调用前加一层参数校验和自动修正,比如字符串数字自动转数字。

5.4 并发下的状态错乱

多个请求同时改同一个会话状态,导致消息顺序错乱。这个前面说了,用会话级锁解决。但要注意锁的粒度,锁太大会影响吞吐,锁太小会失效。我的经验是按session_id加锁,锁的持有时间尽量短,只包住状态读写,不包住模型调用。

问题现象可能原因排查方法解决方案
SSE 频繁断开无心跳/网关超时curl 看心跳行加 ping=15,调大网关超时
记忆检索不准向量模型不匹配打印检索结果换模型,调阈值
工具参数错误schema 描述不清看调用日志完善描述,加校验层
并发状态错乱无会话锁压测复现按 session_id 加锁
上下文超长记忆注入过多统计 token 数限制 Top-K,加摘要

5.5 模型输出不稳定的兜底

生产环境不能假设模型每次都输出正确格式。我的做法是加一层解析兜底:如果模型输出的工具调用 JSON 解析失败,就重试一次,重试还失败就降级为纯文本回复,并记录日志。这样至少不会让用户看到报错。

6. 一些关于扩展方向的个人想法

这套架构搭起来之后,扩展性其实很好。比如想加 RAG,就在记忆域里加一个文档检索模块,把检索结果和记忆一起注入上下文。关键词里提到的“agentscope 2.0 rag as service”就是这个思路。想加多智能体协作,就在编排域里引入多个 Agent 实例,让它们通过消息互相调用。

MCP 生态的膨胀也带来了很多可能性。以前接一个新工具要写适配代码,现在只要那个工具提供了 MCP Server,配置一下就能用。我最近在试把浏览器自动化、数据库查询、文件操作都通过 MCP 接进来,Agent 的能力边界一下就打开了。

最后分享一个小技巧:调试 Agent 的时候,把每一轮的完整上下文(包括系统提示、记忆、工具定义、模型输出)都落盘成 JSON 文件,按 session 和时间戳命名。出问题的时候直接翻文件,比看日志快得多。这个习惯帮我省了无数排查时间。

返回列表