1. 为什么我要从零手搓一个记忆型 AI Agent
市面上开箱即用的 Agent 框架已经多到挑花眼,但我还是决定从零构建一个生产级记忆型 AI Agent。原因很直接:大部分教程里的 Agent 只能叫“演示级”,一旦放到真实业务里,多轮对话记不住上下文、并发一上来就雪崩、流式输出断断续续、人工介入没有入口,这些问题不解决,Agent 永远只是个玩具。
这次我选的技术底座是AgentScope,配合DDD 架构做工程分层,用SSE做流式推送,把HITL(Human-in-the-Loop)作为一等公民设计进去。整套东西的目标很明确:让 AI Agent 真的能“下地干活”,而不是停留在 notebook 里自嗨。
这篇文章适合三类人看。第一类是想从 0 到 1 搭建 AI Agent 的开发者,你不需要有很深的框架经验,但至少要写过 Python 或 Java 的 Web 服务。第二类是被并发和流式输出折磨过的后端,SSE 断连、idle timeout、消息乱序这些坑我会一个个拆开讲。第三类是正在做 AI Agent 中台或二次开发的同学,DDD 分层和记忆模块的设计思路可以直接抄作业。
先把结论摆出来:一个生产级记忆型 Agent 的核心不是模型多强,而是记忆怎么存、上下文怎么拼、流怎么稳、人怎么插手。这四件事做好了,模型换个便宜的也能跑得不错;这四件事做不好,用再贵的模型也是一地鸡毛。
2. 整体架构设计与技术选型拆解
2.1 为什么是 AgentScope 而不是自己造轮子
我一开始也想过纯手写,毕竟 Agent 的逻辑看起来不复杂:接收消息、拼上下文、调模型、返回结果。但真动手就会发现,消息抽象、工具调用、多 Agent 协作、记忆管理这些如果全自己写,光是把不同模型的返回格式统一就是个大坑。
AgentScope 吸引我的点在于它的消息抽象足够干净,Msg对象把 role、content、工具调用结果都统一了,换模型的时候业务代码基本不用动。另外它对多 Agent 对话的原生支持,让我后面扩展“规划 Agent + 执行 Agent”这种模式时省了很多事。
不过 AgentScope 不是银弹。它的默认记忆实现偏简单,直接拿来做生产级记忆是不够的,所以我在它上面套了一层自己的记忆管理层。这也是我建议的做法:框架负责消息流转和模型适配,记忆和业务状态自己管,边界清晰,后面出问题也好排查。
2.2 DDD 分层:别把 Agent 写成一个大泥球
很多人搭 Agent 的习惯是全部塞进一个agent.py,几百行下去,改一个 prompt 都要小心翼翼。我用 DDD 的思路把它拆成四层,这里说的 DDD 不是要你搞一堆领域事件、聚合根那么重,而是借用它的分层思想。
| 层级 | 职责 | 典型内容 |
|---|---|---|
| 接口层 | 对外暴露 HTTP/SSE 接口 | Controller、SSE Emitter、鉴权 |
| 应用层 | 编排用例,协调领域对象 | AgentService、会话管理、HITL 调度 |
| 领域层 | 核心业务逻辑 | 记忆策略、上下文组装、工具定义 |
| 基础设施层 | 技术实现 | 模型客户端、向量库、Redis、DB |
这么分的好处是,当我想把记忆从“滑动窗口”换成“向量召回 + 摘要”时,只动领域层,接口层和应用层完全不用改。反过来,当我要把 SSE 换成 WebSocket 时,也只动接口层。变化被隔离在单层内,这是生产级代码和演示级代码最大的区别。
2.3 SSE 还是 WebSocket:流式输出的选型逻辑
流式输出这块我纠结过 SSE 和 WebSocket。最后选 SSE,理由有三条。
第一,Agent 的输出本质是单向流,服务端推、客户端收,SSE 天然契合,WebSocket 的双向能力用不上还增加复杂度。第二,SSE 基于 HTTP,走标准端口,网关、负载均衡、鉴权中间件都能直接复用,WebSocket 经常要在网关层做额外配置。第三,SSE 的自动重连机制浏览器原生支持,客户端代码简单。
但 SSE 有个绕不开的坑:idle timeout。就是那个经典的stream disconnected before completion: idle timeout waiting for SSE。原因是 Agent 在思考或者调工具的时候,可能几十秒不吐一个字,中间的反向代理或者网关就认为连接死了,直接掐断。解决办法后面会详细讲,核心就是心跳保活 + 服务端超时配置对齐。
2.4 HITL:让 AI 在关键节点停下来等人
HITL 是我认为最容易被忽略但最重要的设计。纯自动的 Agent 在演示时很酷,但生产环境里,涉及资金、删除、对外发送这类操作,你敢让它全自动吗?我不敢。
所以我在架构里把 HITL 做成一个可插拔的拦截点。Agent 在执行工具调用前,先判断这个工具是否标记为“需要人工确认”。如果是,就暂停执行,通过 SSE 推一个确认请求给前端,等用户点了确认或拒绝,再继续。这个暂停不是阻塞线程,而是把会话状态存下来,等确认回来再恢复。这样即使确认要等几分钟,也不会占着连接和线程。
3. 记忆模块的核心设计与实操要点
3.1 记忆不是简单的对话历史堆叠
新手最容易犯的错,就是把所有历史消息一股脑塞进 context。短对话没问题,一旦聊到几十轮,token 直接爆炸,而且模型还会被无关信息干扰,回答质量反而下降。
我的记忆模块分三层,这也是目前比较主流的生产级做法。
短期记忆是最近 N 轮对话,保证连贯性,直接放 context。长期记忆是把历史对话做摘要或者向量化存起来,需要的时候召回。工作记忆是当前任务相关的临时状态,比如正在处理的订单号、用户刚提到的偏好,任务结束就清掉。
三层各司其职,context 里永远只放“当前最相关”的内容,而不是“所有内容”。
3.2 上下文组装:token 预算怎么算
上下文组装的核心是token 预算管理。我一般这么分配:系统提示词占 10%,长期记忆召回占 20%,短期记忆占 40%,当前用户输入和工具结果占 30%。这个比例不是死的,但要有意识地去控制。
具体操作上,我会先算当前模型的最大 context 长度,比如 128k,然后按比例切分。短期记忆从最近往远取,取到预算用完为止。长期记忆用向量检索,取 top-k 相关片段。如果加起来超了,就触发摘要压缩,把更早的对话压成一段摘要。
注意:token 计算不要用字符数除以 2 这种土办法,不同模型的分词器差异很大。用对应模型的 tokenizer 算,误差能控制在几个 token 内。
3.3 记忆持久化:Redis 加向量库的组合拳
短期记忆我放 Redis,因为读写快,而且可以设 TTL 自动过期。长期记忆放向量库,我用的是轻量级的方案,存摘要和对应的向量。工作记忆也放 Redis,但用单独的 key 前缀,任务结束主动删。
这里有个实操心得:Redis 的 key 设计要带会话 ID 和用户 ID,比如mem:short:{user_id}:{session_id}。这样查的时候一次命中,也方便做多租户隔离。我见过有人把所有会话塞一个 list 里,查的时候全量扫,并发一上来直接拖垮 Redis。
3.4 记忆召回的相关性判断
向量召回不是召回了就完事,还要做相关性过滤。我的做法是设一个相似度阈值,低于阈值的直接丢掉,宁可少召回也不要召回噪音。另外召回的内容要带上时间戳,因为用户上周说的偏好和今天说的可能冲突,时间新的优先。
还有个细节:召回的记忆要重新格式化再放进 context,不能直接把数据库里的原始记录塞进去。我会把它包装成“根据历史对话,用户曾经提到……”这种自然语言形式,模型理解起来更顺。
4. SSE 流式接口的完整实现与踩坑记录
4.1 SSE 服务端实现的关键参数
SSE 服务端的核心是保持连接、按格式推数据、正确处理断开。以 Python 的 FastAPI 为例,返回一个StreamingResponse,media_type 设成text/event-stream。
from fastapi import FastAPI from fastapi.responses import StreamingResponse import asyncio app = FastAPI() async def event_generator(session_id: str): # 先推一个心跳,告诉客户端连接建立 yield "event: connected\ndata: {\"status\":\"ok\"}\n\n" try: async for chunk in agent_stream(session_id): # 每条消息按 SSE 格式封装 yield f"event: message\ndata: {chunk}\n\n" # 每推一条就检查是否需要心跳 except asyncio.CancelledError: # 客户端断开,清理资源 cleanup(session_id) raise @app.get("/agent/stream") async def stream(session_id: str): return StreamingResponse( event_generator(session_id), media_type="text/event-stream", headers={ "Cache-Control": "no-cache", "Connection": "keep-alive", "X-Accel-Buffering": "no", # 关键:禁用 Nginx 缓冲 }, )这里有几个参数是血泪教训换来的。X-Accel-Buffering: no必须加,否则 Nginx 会缓冲你的流,客户端要等一大坨才收到,流式就失去意义了。Cache-Control: no-cache防止中间层缓存。Connection: keep-alive保持长连接。
4.2 心跳保活:解决 idle timeout 的正解
stream disconnected before completion: idle timeout waiting for SSE这个报错,本质是连接空闲太久被中间层掐了。解决办法是定期发心跳。
我的做法是起一个后台任务,每 15 秒检查一次,如果距离上次推数据超过 15 秒,就推一个注释行: heartbeat\n\n。SSE 规范里以冒号开头的行是注释,客户端会忽略,但能保持连接活跃。
async def keepalive(emitter, interval=15): while True: await asyncio.sleep(interval) if emitter.idle_seconds() >= interval: await emitter.send_comment("heartbeat")同时,服务端、网关、负载均衡的超时时间要对齐。我一般设成:服务端 300 秒,网关 300 秒,负载均衡 300 秒。如果网关是 60 秒而服务端是 300 秒,那 60 秒一到照样断。这个对齐工作不做,光加心跳也没用。
4.3 前端消费 SSE 的正确姿势
前端用EventSource消费 SSE,但要注意它只支持 GET 请求,而且不能自定义 header。如果鉴权需要 token,要么放 query 参数,要么用 fetch + ReadableStream 自己解析。
const es = new EventSource(`/agent/stream?session_id=${sid}&token=${token}`); es.addEventListener('message', (e) => { const data = JSON.parse(e.data); appendToChat(data.content); }); es.addEventListener('error', (e) => { // EventSource 会自动重连,但要处理重连后的状态同步 console.warn('SSE error, will reconnect', e); // 重连后需要拉取断连期间的消息 syncMissedMessages(sid); }); es.addEventListener('done', () => { es.close(); });这里有个坑:EventSource断线重连后,断连期间的消息会丢。所以我在服务端给每条消息带一个递增的 seq,前端重连时带上最后收到的 seq,服务端把之后的消息补推。这个机制不做,用户就会遇到“回答突然少了一段”的诡异问题。
4.4 并发场景下的 SSE 连接管理
“AI Agent 怎么扛并发”是热词里高频出现的问题。SSE 是长连接,一个用户一条连接,1 万用户就是 1 万条连接。这时候连接管理就很重要。
我的做法是:连接和 Agent 执行解耦。Agent 的执行放到独立的 worker 里,结果写到一个消息队列或者 Redis 的 stream 里。SSE 连接只负责从队列里读消息推给客户端。这样即使客户端断开,Agent 的执行也不受影响,重连后还能接着推。
另外要限制单用户的连接数,防止有人开一堆标签页把连接占满。我一般限制单用户最多 3 条 SSE 连接,超了就踢掉最老的。
5. HITL 人工介入的落地实现
5.1 什么操作该触发人工确认
不是所有操作都要人工确认,那样用户体验会很差。我的判断标准是不可逆 + 高影响。比如删除数据、发起支付、对外发送消息,这些必须确认。查询类、计算类操作直接放行。
实现上,我在工具定义里加一个requires_confirmation标记,Agent 调用工具前先检查这个标记。这样新增工具时只要打个标,不用改调度逻辑。
5.2 暂停与恢复的状态管理
HITL 最难的是状态管理。Agent 执行到一半要暂停,这时候上下文、已执行的工具结果、待确认的操作,都要存下来。我用一个PendingAction对象存这些状态,序列化后放 Redis,key 是hitl:{session_id}:{action_id}。
用户确认后,根据 action_id 把状态捞出来,继续执行。这里要注意幂等性,用户可能重复点确认,或者网络重试导致确认请求发两次。我的做法是确认操作带一个唯一 ID,服务端处理前先检查这个 ID 是否已处理过。
5.3 确认请求的推送与超时处理
确认请求通过 SSE 推给前端,前端弹窗让用户选择。这里要设超时,比如 5 分钟没确认就自动取消,并推一条消息告诉用户“操作已超时取消”。
超时时间不能太长,否则状态一直挂着占资源;也不能太短,用户可能正在看别的。5 分钟是我实测下来比较平衡的值。
提示:HITL 的确认请求要带上足够的上下文,让用户知道自己在确认什么。只显示“是否确认执行 delete_user?”是不够的,要显示“即将删除用户 张三(ID: 12345),该操作不可恢复”。
6. 常见问题排查与避坑速查
6.1 SSE 相关高频问题
| 问题现象 | 根本原因 | 解决办法 |
|---|---|---|
| 流式输出一次性全出来 | Nginx 缓冲 | 加X-Accel-Buffering: no |
| 几十秒后连接断开 | idle timeout | 心跳保活 + 超时对齐 |
| 重连后消息丢失 | 无 seq 机制 | 消息带 seq,重连补推 |
| 中文乱码 | 编码未指定 | 响应头加charset=utf-8 |
| 连接数暴涨 | 未限制单用户连接 | 限制单用户连接数 |
6.2 记忆模块常见坑
第一个坑是记忆无限增长。短期记忆如果不设上限,聊得越久 context 越大,最后直接超模型限制。一定要设滑动窗口大小,比如最近 20 轮。
第二个坑是摘要丢失关键信息。做长期记忆摘要时,如果摘要太粗,关键信息就丢了。我的做法是摘要时保留实体(人名、订单号、金额),这些是后续召回的关键。
第三个坑是多会话串味。用户开了两个会话,结果 A 会话的记忆跑到 B 会话里去了。这是 key 设计问题,一定要用user_id + session_id做隔离。
6.3 并发与性能问题
Agent 执行是 IO 密集型,等模型返回的时候线程是空闲的。所以要用异步,别用同步阻塞。Python 里用asyncio,Java 里用CompletableFuture或者响应式。
模型调用要做超时和重试。模型服务偶尔抽风很正常,设个 30 秒超时,失败重试 2 次,还失败就降级返回一个兜底回复。别让一个模型调用卡死整个会话。
还有个容易被忽略的点:工具调用的并发。如果 Agent 一次要调多个独立工具,可以并发调,别串行等。我实测下来,3 个工具并发调用比串行快 2 倍多。
6.4 我踩过的三个真实坑
第一个坑:早期我没做 SSE 的 seq 机制,测试时网络一抖,用户就反馈“回答少了一段”。排查了半天才定位到是重连丢消息。加上 seq 后彻底解决。
第二个坑:HITL 的确认状态我一开始放内存,单机测试没问题,一上多实例就出问题——确认请求打到 A 实例,但状态存在 B 实例。后来改成 Redis 共享状态才解决。任何要跨请求的状态,都别放进程内存。
第三个坑:记忆召回我一开始没做相关性阈值,结果召回一堆无关内容,模型被带偏,回答质量反而比不召回还差。加上阈值过滤后,召回质量明显提升。召回不是越多越好,是越准越好。
7. 从练手项目到生产级的扩展思路
如果你只是想练手,把前面说的记忆三层、SSE 流式、HITL 拦截跑通,就已经比 90% 的教程项目完整了。但要做成生产级,还有几件事要补。
可观测性。Agent 的执行链路要能追踪,每次模型调用、工具调用、记忆召回都要打点。出了问题能快速定位是哪一环。我用的是 OpenTelemetry 那套,trace 一拉,整个链路清清楚楚。
成本控制。模型调用是花钱的,要做 token 用量统计和限额。单用户单日 token 超了就限流,防止有人恶意刷。这个不做,账单会教你做人。
多模型路由。简单问题用便宜模型,复杂问题用强模型。我做了个简单的路由规则,根据输入长度和关键词判断,能省不少成本。
灰度与回滚。Prompt 改了、记忆策略改了,不能直接全量上。要有灰度机制,先放 5% 流量,观察指标没问题再全量。出问题能一键回滚到上个版本。
这套东西搭下来,你会发现 Agent 的难点从来不在模型,而在工程。记忆怎么管、流怎么稳、人怎么插手、并发怎么扛,这些才是决定一个 Agent 能不能上生产的关键。AgentScope 给了你一个好的起点,但剩下的路,还得自己一步步走。