1. 为什么我要从零手搓一个记忆型 Agent,而不是直接套框架
2026 年这个时间点,市面上能叫得出名字的 AI Agent 框架两只手数不过来,从轻量脚本到企业级平台都有。但我还是花了将近三周时间,从零搭了一个带长期记忆的生产级 Agent,底层用 AgentScope 的思路做编排,交互层用 SSE 做流式渲染,架构上按 DDD 分层。原因很简单:框架能帮你跑通 Demo,但跑不通生产。
我最早也是直接拿现成框架改的,结果遇到三个绕不过去的问题。第一,记忆模块和推理模块耦合太深,想换一套向量检索策略,得把整个 Agent 的调用链翻一遍;第二,流式输出在长会话下会断,前端拿到的是一段一段的碎片,用户看到的是"打字机卡住";第三,人工介入(HITL)没有干净的切入点,想在某个工具调用前插一道人工确认,只能靠 hack 回调。
这三个问题本质上不是框架的锅,而是架构分层没做对。Agent 这个东西,表面上是"大模型 + 工具调用",实际上它是一个有状态、有记忆、有中断恢复需求的分布式系统。你把它当成一个函数来写,它迟早会在生产环境里教你做人。
所以这篇东西,我想把整个搭建过程拆开讲清楚:DDD 怎么切分 Agent 的领域边界、SSE 流式输出怎么做到不丢包不断流、记忆层怎么设计才能既快又准、HITL 怎么插进去才不破坏主流程。适合已经跑通过 Demo、准备往生产推的开发者,也适合想理解 Agent 内部到底怎么运转的人。代码层面我会给关键片段,但重点在为什么这么设计,因为抄代码容易,抄思路难。
2. 用 DDD 给 Agent 划边界:哪些该是领域,哪些只是基础设施
2.1 Agent 的领域模型到底长什么样
很多人做 Agent 开发,脑子里只有一条线:用户输入 → 拼 Prompt → 调模型 → 解析工具调用 → 执行工具 → 再拼 Prompt → 输出。这条线在 Demo 阶段没问题,但一旦你要加记忆、加人工审核、加多轮工具编排,这条线就会变成一团意大利面。
DDD 的核心价值在这里体现得很明显:它逼你先想清楚"什么是业务概念",再想"怎么实现"。我把整个 Agent 拆成四个限界上下文:
- 会话上下文(Conversation Context):管理一次对话的生命周期,包括消息历史、会话状态、超时策略。这里的聚合根是
Session,它持有Message列表和SessionStatus。 - 记忆上下文(Memory Context):负责长期记忆的写入、检索、衰减。聚合根是
MemoryEntry,它有自己的重要性评分和访问频次。 - 推理上下文(Reasoning Context):封装大模型调用、工具选择、思维链编排。这里的核心是
ReasoningStep,每一步都是一次"思考-行动-观察"的循环。 - 人工介入上下文(HITL Context):处理需要人工确认的节点,包括审批流、超时降级、结果回填。
这四个上下文之间通过领域事件通信,而不是直接方法调用。比如推理上下文决定要调用一个高危工具时,它不直接问 HITL "要不要批准",而是发一个ToolExecutionRequested事件,HITL 上下文订阅后决定是自动放行还是挂起等人工。
这么切的好处是:记忆策略换了,推理层完全无感;HITL 从"同步阻塞"改成"事件驱动",主流程不会被卡死。
2.2 分层落地:应用层、领域层、基础设施层各放什么
DDD 讲分层,但 Agent 场景下有个特殊点:大模型调用本身既是基础设施,又深度参与领域逻辑。我的处理方式是:
- 领域层:放
Session、MemoryEntry、ReasoningStep这些实体和值对象,以及领域服务如MemoryRetrievalService(定义"怎么算相关"的规则)。这一层不依赖任何具体的大模型 SDK。 - 应用层:放用例编排,比如
ChatUseCase、MemoryConsolidationUseCase。它负责协调领域对象和基础设施,但不含业务规则。 - 基础设施层:放具体的模型客户端、向量库、SSE 推送器、持久化实现。这一层实现领域层定义的接口(端口)。
关键接口长这样:
// 领域层定义的端口 public interface LanguageModelPort { ReasoningResult reason(ReasoningContext context); } public interface MemoryStorePort { List<MemoryEntry> retrieve(String query, int topK); void persist(MemoryEntry entry); } public interface StreamPushPort { void push(String sessionId, StreamChunk chunk); }基础设施层用 AgentScope 的模型封装、Redis 向量检索、SSE Emitter 分别实现这三个端口。换模型、换向量库、换推送方式,领域层一行不用改。这就是 DDD 在 Agent 项目里最实在的收益。
2.3 一个容易踩的坑:别把 Prompt 当领域对象
我见过不少项目把 Prompt 模板放在领域层,甚至做成实体。这是个陷阱。Prompt 是实现细节,它随模型版本、随业务调优频繁变化,把它放进领域层会导致领域模型极不稳定。
我的做法是把 Prompt 模板放在基础设施层的资源目录,领域层只定义"需要哪些信息"(比如ReasoningContext里包含历史消息、可用工具、记忆片段),由基础设施层负责把这些信息渲染成具体 Prompt。这样调 Prompt 不用动领域代码,领域模型也不会被 Prompt 的频繁变更污染。
3. SSE 流式输出:从"打字机卡顿"到丝滑实时渲染的完整链路
3.1 为什么选 SSE 而不是 WebSocket
热词里有个问题很典型:"react + sse/websocket 轮询文件变化"。选型这事我踩过坑,直接说结论:Agent 的对话流式输出,SSE 是更优解,除非你需要双向实时通信。
对比一下:
| 维度 | SSE | WebSocket |
|---|---|---|
| 通信方向 | 服务端单向推送 | 双向 |
| 协议 | 纯 HTTP | 独立协议,需升级握手 |
| 断线重连 | 浏览器原生支持 | 需自己实现 |
| 代理兼容 | 好 | 部分代理会拦截 |
| 实现复杂度 | 低 | 中高 |
Agent 场景下,客户端主要是"发一次请求,收一串流式响应",天然是单向的。用 WebSocket 属于杀鸡用牛刀,还要处理心跳、重连、协议升级一堆事。SSE 基于 HTTP,天然穿透大部分网络环境,浏览器EventSource自带重连。
但 SSE 有个坑:默认的EventSource不支持 POST,也不支持自定义 Header。而 Agent 请求往往需要带认证 Token 和较长的请求体。解决办法是用fetch+ReadableStream手动解析 SSE 流,而不是用EventSource。
3.2 服务端:SSE 推送的三个关键设计
服务端我用 Spring 的SseEmitter,但直接裸用会出问题。三个关键设计:
第一,分块粒度要合理。大模型返回的是 token 流,但你不能每个 token 推一次,那样网络开销太大。我的做法是按"语义块"推送:遇到标点、换行、或者累积到一定长度才 flush 一次。实测下来,中文场景下按 8-15 个字符一块,前端渲染最顺滑。
第二,心跳不能少。热词里那个stream disconnected before completion: idle timeout waiting for sse就是典型的没做心跳。中间隔太久没数据,网关或代理会主动断开。我的做法是每 15 秒推一个注释行:heartbeat,这个在 SSE 协议里是合法的,客户端会忽略,但能保活连接。
// 心跳任务 ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1); scheduler.scheduleAtFixedRate(() -> { try { emitter.send(SseEmitter.event().comment("heartbeat")); } catch (IOException e) { // 连接已断,清理资源 emitter.complete(); } }, 0, 15, TimeUnit.SECONDS);第三,abort 要能真正中断。用户点了"停止生成",服务端得真的停下来,不能只是前端不显示。我的做法是给每个会话维护一个AtomicBoolean标志,推理循环每步检查一次,同时调用模型客户端的 cancel 接口。
public void abort(String sessionId) { AtomicBoolean flag = abortFlags.get(sessionId); if (flag != null) { flag.set(true); } // 同时通知模型客户端取消 modelClient.cancel(sessionId); }3.3 前端:fetch 流式解析与渲染节流
前端不用EventSource,用fetch拿ReadableStream,然后手动按\n\n切分事件块。核心逻辑:
const response = await fetch('/api/chat/stream', { method: 'POST', headers: { 'Content-Type': 'application/json', 'Authorization': token }, body: JSON.stringify({ sessionId, message }), signal: abortController.signal }); const reader = response.body.getReader(); const decoder = new TextDecoder(); let buffer = ''; while (true) { const { done, value } = await reader.read(); if (done) break; buffer += decoder.decode(value, { stream: true }); const events = buffer.split('\n\n'); buffer = events.pop(); // 最后一段可能不完整,留到下次 for (const event of events) { const dataLine = event.split('\n').find(l => l.startsWith('data:')); if (dataLine) { const chunk = JSON.parse(dataLine.slice(5)); appendToUI(chunk); } } }这里有个细节:渲染要节流。如果每个 chunk 都触发一次 React setState,高频 token 流下会把主线程打满。我的做法是用requestAnimationFrame做批量更新,把一帧内的多个 chunk 合并成一次渲染。
let pending = ''; let rafId = null; function appendToUI(chunk) { pending += chunk.content; if (!rafId) { rafId = requestAnimationFrame(() => { setMessage(prev => prev + pending); pending = ''; rafId = null; }); } }实测下来,这套组合拳能把长会话的渲染帧率稳定在 60fps,用户看到的就是丝滑的打字机效果,而不是一顿一顿的。
3.4 断流恢复:让用户刷新页面也不丢内容
生产环境里,用户网络抖动、切后台、刷新页面都是常态。如果流一断内容就没了,体验极差。我的方案是服务端持久化每个 chunk 的序号,客户端重连时带上最后收到的序号,服务端从那个序号之后继续推。
// 服务端:每个 chunk 带序号 emitter.send(SseEmitter.event() .id(String.valueOf(chunkSeq.incrementAndGet())) .data(chunkJson)); // 客户端重连时带 Last-Event-ID // 服务端从该 ID 之后重放SSE 协议本身支持Last-Event-ID头,浏览器EventSource会自动带上。但我们用fetch手动实现,就得自己管理这个序号。多写几行代码,换来的是用户刷新页面后内容无缝续上,这个投入非常值。
4. 记忆层设计:让 Agent 真的"记得住"而不是"假装记得"
4.1 短期记忆、长期记忆、工作记忆的分工
"记忆型 Agent"这个词被用烂了,但很多所谓记忆就是把历史消息全塞进 Prompt。这不叫记忆,这叫上下文堆砌,token 烧得飞快,效果还差。
我按认知科学的思路分了三层:
- 工作记忆(Working Memory):当前这一轮推理需要的信息,包括最近几轮对话、当前任务相关的记忆片段、可用工具列表。它是有容量上限的,我设的是 8K token。
- 短期记忆(Short-term Memory):本次会话的完整历史,存在 Redis 里,按会话 ID 索引。它不直接进 Prompt,而是作为工作记忆的检索源。
- 长期记忆(Long-term Memory):跨会话的持久记忆,存在向量库里。用户说过的偏好、重要事实、历史决策都在这。
三层之间的流转是:每轮对话结束后,短期记忆里产生新内容,由一个"记忆巩固"过程判断哪些值得写入长期记忆。这个判断用一个小模型或者规则引擎做,不是所有内容都值得长期记。
4.2 记忆写入:什么该记,什么该忘
记忆写入最大的坑是什么都记。用户随口说一句"今天天气不错",你把它写进长期记忆,下次检索出来就是噪音。
我的写入策略是打分制,三个维度:
- 信息密度:是否包含实体、事实、偏好。用简单的 NER + 关键词匹配就能粗筛。
- 情感强度:用户强调、重复、带情绪的内容,权重更高。
- 任务相关性:与当前进行中的任务相关的,优先记。
综合分超过阈值的才写入长期记忆。同时给每条记忆一个衰减因子,随时间推移和访问频次降低而衰减,检索时按相关性 × 衰减因子排序。
public class MemoryEntry { private String content; private float importance; // 初始重要性 0-1 private long createdAt; private int accessCount; public float currentScore() { long ageHours = (System.currentTimeMillis() - createdAt) / 3600000; float decay = (float) Math.exp(-ageHours / 168.0); // 一周半衰期 float accessBoost = 1 + (float) Math.log1p(accessCount) * 0.1f; return importance * decay * accessBoost; } }这个公式不复杂,但效果比"全量塞 Prompt"好太多。实测在 200 轮以上的长会话里,检索准确率能保持在 85% 以上。
4.3 记忆检索:向量 + 关键词的混合召回
纯向量检索有个问题:对精确匹配不敏感。用户问"我上次说的那个项目叫什么来着",向量检索可能召回一堆语义相近但不对的内容。纯关键词检索又抓不住语义。
我的做法是混合召回:向量检索取 Top 20,BM25 关键词检索取 Top 20,然后用 RRF(Reciprocal Rank Fusion)融合排序,取 Top 5 进 Prompt。
public List<MemoryEntry> hybridRetrieve(String query, int topK) { List<MemoryEntry> vectorResults = vectorStore.search(embed(query), 20); List<MemoryEntry> keywordResults = keywordIndex.search(query, 20); Map<String, Double> scores = new HashMap<>(); for (int i = 0; i < vectorResults.size(); i++) { scores.merge(vectorResults.get(i).getId(), 1.0 / (60 + i), Double::sum); } for (int i = 0; i < keywordResults.size(); i++) { scores.merge(keywordResults.get(i).getId(), 1.0 / (60 + i), Double::sum); } return scores.entrySet().stream() .sorted(Map.Entry.<String, Double>comparingByValue().reversed()) .limit(topK) .map(e -> memoryStore.get(e.getKey())) .collect(Collectors.toList()); }RRF 里的常数 60 是经验值,来自信息检索领域的经典论文,不用纠结,直接用就行。
4.4 记忆巩固:会话结束后的异步整理
会话结束后,我跑一个异步任务做记忆巩固:把短期记忆里的内容做摘要、去重、提取事实,然后按写入策略决定哪些进长期记忆。这个任务不阻塞用户,用消息队列异步处理。
巩固过程还会做记忆合并:如果新记忆和已有记忆高度相似,就更新已有记忆的访问时间和重要性,而不是新增一条。这样长期记忆库不会无限膨胀。
5. HITL 人工介入:怎么插进去才不破坏主流程
5.1 哪些节点需要人工确认
HITL 不是到处都要,那样 Agent 就没法自动跑了。我的经验是三类节点需要:
- 高危工具调用:比如删除数据、发送邮件、调用外部支付接口。这类操作一旦执行不可逆,必须人工确认。
- 低置信度决策:模型对某个判断的置信度低于阈值,比如 0.6,就挂起等人工。
- 合规敏感内容:涉及特定业务规则的输出,需要人工审核后才能返回。
5.2 事件驱动的挂起与恢复
前面说了,HITL 用事件驱动,不阻塞主流程。具体实现:
推理上下文在决定调用高危工具时,发一个ToolExecutionRequested事件,然后把当前推理状态序列化存起来,返回一个"等待人工确认"的状态给前端。HITL 上下文收到事件后,生成一个待办任务推给人工审核界面。
人工确认后,HITL 上下文发ToolExecutionApproved或ToolExecutionRejected事件,推理上下文收到后从序列化的状态恢复,继续执行。
// 推理上下文:挂起 public ReasoningResult handleToolCall(ToolCall call) { if (riskAssessor.isHighRisk(call)) { String snapshotId = stateSerializer.snapshot(currentState); eventBus.publish(new ToolExecutionRequested(sessionId, call, snapshotId)); return ReasoningResult.suspended(snapshotId); } return executeTool(call); } // 恢复 public ReasoningResult resume(String snapshotId, boolean approved) { ReasoningState state = stateSerializer.restore(snapshotId); if (!approved) { return state.withObservation("用户拒绝了此操作"); } return executeTool(state.pendingCall()); }这个设计的关键是状态快照。Agent 的推理状态包括消息历史、当前步骤、待执行的工具调用,全部序列化后可以跨请求恢复。这样人工审核可能花几分钟甚至几小时,主流程也不会一直占着资源。
5.3 超时降级:人工不响应怎么办
人工审核不能无限等。我的策略是分级超时:
- 5 分钟无响应:发提醒
- 30 分钟无响应:按预设策略降级(高危操作默认拒绝,低危操作默认放行)
- 2 小时无响应:会话标记为过期,清理快照
降级策略要可配置,不同业务场景不一样。金融场景可能全部默认拒绝,内部工具场景可能默认放行。
6. 生产环境里那些文档不会告诉你的坑
6.1 模型输出的 JSON 解析失败率比你想的高
工具调用需要模型输出结构化 JSON,但实测下来,即使加了严格的格式约束,仍有 3%-5% 的解析失败率。失败原因五花八门:多了个逗号、少了引号、中文标点混入、嵌套层级错误。
我的处理是三层兜底:第一层用宽松解析器(比如 Jackson 的ALLOW_SINGLE_QUOTES、ALLOW_UNQUOTED_FIELD_NAMES);第二层用正则提取 JSON 片段再解析;第三层解析失败就把原始输出回传给模型,让它自己修正。三层下来,失败率能压到 0.5% 以下。
6.2 长会话的 token 成本控制
200 轮以上的会话,如果每轮都把全量历史塞进去,token 成本会爆炸。我的做法是滑动窗口 + 摘要压缩:最近 10 轮保留原文,更早的做摘要,摘要再往前就只保留关键事实。
摘要不是简单截断,而是用模型生成一段 200 字以内的浓缩版。这个摘要生成是异步的,不阻塞主流程。实测能把长会话的 token 消耗降低 60% 以上。
6.3 并发会话下的资源隔离
多个用户同时用,会话之间必须隔离。我用的是会话级线程池 + 信号量限流:每个会话最多占 2 个推理线程,全局最多 50 个并发推理。超过就排队,避免一个用户的复杂任务把整个服务拖垮。
Semaphore globalLimit = new Semaphore(50); Map<String, Semaphore> sessionLimits = new ConcurrentHashMap<>(); public ReasoningResult reason(String sessionId, ReasoningContext ctx) { Semaphore sessionLimit = sessionLimits.computeIfAbsent(sessionId, k -> new Semaphore(2)); if (!sessionLimit.tryAcquire()) { throw new TooManyRequestsException("会话并发超限"); } try { globalLimit.acquire(); try { return doReason(ctx); } finally { globalLimit.release(); } } finally { sessionLimit.release(); } }6.4 可观测性:没有日志的 Agent 就是黑盒
Agent 出问题时,你光看输入输出根本不知道哪一步错了。我的做法是全链路埋点:每次模型调用记录 prompt、response、耗时、token 数;每次工具调用记录参数、结果、耗时;每次记忆检索记录 query、召回结果、得分。
这些日志用结构化格式(JSON)打到专门的日志系统,配合 traceId 串起来。出问题时,一个 traceId 就能还原整个推理链路。这个投入在排查线上问题时回报巨大。
7. 从 Demo 到生产,我总结的几条硬经验
搭完这套东西,回头看,有几个判断我觉得是对的,也分享给准备往生产推的朋友。
第一,架构分层不是过度设计,是保命符。Agent 的复杂度会随着功能增加指数级上升,没有清晰的分层,三个月后你自己都不敢改代码。DDD 那套东西看着重,但它帮你把"什么会变"和"什么不变"分开了,变的部分隔离在基础设施层,核心领域稳定。
第二,流式输出是体验的生命线。用户等 10 秒看到完整回答,和 1 秒开始看到字一个个蹦出来,感受天差地别。SSE 这套东西不难,但细节多:心跳、分块、abort、断流恢复,每一个都得处理到位,少一个用户就会遇到"卡住"。
第三,记忆要做减法而不是加法。什么都记等于什么都没记。写入策略、衰减机制、混合检索,这三样做好,记忆才有价值。我见过太多项目把记忆做成"历史消息全塞",那不是记忆,那是烧钱。
第四,HITL 要异步化。同步阻塞的人工审核会把 Agent 的并发能力锁死。事件驱动 + 状态快照,让审核可以慢慢来,主流程该干嘛干嘛。
最后说个我自己的体会:Agent 这个领域,框架更新太快,今天学的 API 明天可能就变了。但架构思想、流式处理、记忆设计、状态管理这些底层的东西,是跨框架通用的。把精力花在这些上面,比追框架版本划算得多。我搭这套东西的过程中,AgentScope 的编排思路给了我不少启发,但真正让它在生产环境跑稳的,是那些跟具体框架无关的工程决策。