ARTICLE DETAIL

资讯详情

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

从零手搓生产级记忆型Agent:DDD分层、SSE流式与HITL实战

从零手搓生产级记忆型Agent:DDD分层、SSE流式与HITL实战 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 { ListMemoryEntry 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 是更优解除非你需要双向实时通信。对比一下维度SSEWebSocket通信方向服务端单向推送双向协议纯 HTTP独立协议需升级握手断线重连浏览器原生支持需自己实现代理兼容好部分代理会拦截实现复杂度低中高Agent 场景下客户端主要是发一次请求收一串流式响应天然是单向的。用 WebSocket 属于杀鸡用牛刀还要处理心跳、重连、协议升级一堆事。SSE 基于 HTTP天然穿透大部分网络环境浏览器EventSource自带重连。但 SSE 有个坑默认的EventSource不支持 POST也不支持自定义 Header。而 Agent 请求往往需要带认证 Token 和较长的请求体。解决办法是用fetchReadableStream手动解析 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 20BM25 关键词检索取 Top 20然后用 RRFReciprocal Rank Fusion融合排序取 Top 5 进 Prompt。public ListMemoryEntry hybridRetrieve(String query, int topK) { ListMemoryEntry vectorResults vectorStore.search(embed(query), 20); ListMemoryEntry keywordResults keywordIndex.search(query, 20); MapString, 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, DoublecomparingByValue().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); MapString, 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 的编排思路给了我不少启发但真正让它在生产环境跑稳的是那些跟具体框架无关的工程决策。
返回列表