从零构建生产级记忆型AI Agent:DDD分层架构与SSE流式输出实战

发布时间:2026/9/28 13:17:46
从零构建生产级记忆型AI Agent:DDD分层架构与SSE流式输出实战 1. 为什么我要从零手搓一个记忆型 AI Agent2025 年年底那阵子我手上同时压着三个跟 AI Agent 相关的需求一个内部知识库问答、一个自动化运维助手、还有一个给业务部门做的流程编排工具。最开始我想得很简单直接拿现成的框架套一套把大模型的接口接上再挂几个工具函数不就完事了结果真动手才发现能跑起来的 Demo 和能上生产的 Agent中间隔着一整条鸿沟。最典型的坑就是“记忆”。你问它上一轮说了什么它记得你隔了三天再问同一件事它一脸茫然。你让它记住用户的偏好它转头就忘。你让它基于历史对话做决策它把上下文窗口塞爆了之后开始胡言乱语。这不是模型不行是架构层面就没有把记忆当成一等公民来设计。所以我决定从零构建一个生产级的记忆型 AI Agent。这里的“记忆型”不是指简单地把对话历史拼进 prompt而是要有分层记忆结构、记忆的写入与召回策略、记忆的衰减与压缩机制还要能通过 SSE 流式输出把大模型的回答实时渲染到前端配合 abort 机制让用户随时打断。技术栈上我选了 AgentScope 作为核心框架用 DDD 的思路做分层架构用 MCP 协议对接外部工具前端用 React SSE 做流式渲染。这篇文章我会把整个项目的设计思路、核心实现、踩过的坑全部摊开讲。适合谁看如果你已经会用大模型 API 写个聊天机器人但不知道怎么把它做成一个有记忆、能流式输出、可扩展工具、能上生产的 Agent那这篇就是写给你的。如果你还在纠结“AI Agent 有哪些”“AI Agent 开发从哪入手”也可以先跟着走一遍至少能建立起完整的工程认知。提示本文涉及的所有代码和配置都是基于实际项目脱敏后的版本你可以直接参考复现。涉及具体业务逻辑的部分我会用伪代码或简化实现代替但架构和关键机制是完整的。2. 整体架构设计与技术选型背后的取舍2.1 为什么是 AgentScope 而不是自己从零写市面上做 AI Agent 的框架不少LangChain、AutoGen、AgentScope 各有各的路子。我最后选 AgentScope核心原因有三个。第一AgentScope 对多 Agent 通信的原生支持更干净。它的消息传递机制是基于消息队列的Agent 之间通过msg对象通信而不是靠全局状态或者回调地狱。这意味着我后面要扩展成多 Agent 协作时不需要重构通信层。第二它的记忆模块是可插拔的。AgentScope 把 memory 抽象成了一个独立的组件你可以自己实现Memory接口替换掉默认的短期记忆。这一点对我这种要做“记忆型 Agent”的需求来说太关键了。LangChain 的 memory 虽然也灵活但它的抽象层级偏高很多细节被封装掉了反而不容易做精细控制。第三AgentScope 2.0 开始对 RAG as a Service 和 MCP 的支持更完善。MCP 协议在 2025 年已经成为工具对接的事实标准之一AgentScope 对 MCP server 的接入做了封装省了我不少事。当然AgentScope 也不是没有缺点。它的中文文档虽然有了但部分高级用法的示例还是偏少有些地方得翻源码。另外它的 Java 版本AgentScope Java 2.0在企业级实战中的资料相对 Python 版更少如果你团队是 Java 技术栈前期调研成本会高一些。2.2 DDD 分层架构怎么落到 Agent 项目上DDD领域驱动设计在传统业务系统里很常见但用在 AI Agent 项目上很多人会觉得“是不是过度设计”。我的经验是如果你的 Agent 只是玩具那确实不需要但如果要上生产DDD 的分层能帮你把“模型交互逻辑”和“业务逻辑”彻底隔离开。我的分层是这样的接口层Interface负责对外暴露 SSE 接口、接收用户请求、处理 abort 信号。这一层不包含任何业务逻辑只做协议转换和流式输出。应用层Application编排 Agent 的执行流程协调记忆模块、工具调用、模型推理。这一层是“导演”不干具体活。领域层Domain定义 Agent、Memory、Tool 这些核心领域对象和领域服务。记忆的写入策略、召回算法、衰减规则都在这一层。基础设施层Infrastructure对接具体的大模型 API、向量数据库、MCP server、持久化存储。这一层是“可替换的零件”。这么分的好处是当我从 OpenAI 的模型切换到国内某个模型时只需要改基础设施层的适配器领域层和应用层完全不动。当我要把短期记忆从内存换成 Redis 时也只需要换一个 Memory 实现类。2.3 SSE 流式输出与 abort 机制的配合大模型回答的实时渲染目前主流方案就是 SSEServer-Sent Events。相比 WebSocketSSE 的优势是单向推送、基于 HTTP、实现简单、自动重连。对于“用户发一个问题模型流式返回答案”这个场景SSE 完全够用而且前端用EventSource就能接不需要引入额外的 WebSocket 库。但 SSE 有个坑用户想中途打断怎么办比如模型开始胡言乱语了用户想点“停止生成”。这时候需要前端发一个 abort 信号后端收到后中断对模型的请求并关闭 SSE 连接。我的做法是前端用AbortController配合fetch来发 SSE 请求而不是用EventSource因为EventSource不支持自定义 abort。后端在 SSE 流中监听一个 abort 端点或者通过请求的signal来判断。具体实现我在第 4 章会详细讲。2.4 MCP 协议在工具对接中的角色MCPModel Context Protocol本质上是一个标准化的工具描述和调用协议。你可以把它理解成“AI Agent 世界的 USB 接口”——只要工具实现了 MCP server任何支持 MCP 的 Agent 都能直接调用它不需要为每个工具写适配代码。我项目里对接了三个 MCP server一个是内部知识库的检索工具一个是 Playwright MCP用于浏览器自动化还有一个是蓝湖 MCP用于读取设计稿信息。有了 MCP我只需要在配置里声明 server 地址和 tokenAgent 就能自动发现这些工具并调用。注意MCP server 的 token 一定要通过环境变量注入不要硬编码在代码或配置文件里。我见过有人把 token 直接写在application.yml里然后提交到了仓库这是大忌。3. 记忆系统的核心设计与实现细节3.1 三层记忆结构短期、长期、工作记忆记忆型 Agent 的核心不是“记住所有对话”而是在该记住的时候记住在该忘记的时候忘记。我设计了三层记忆结构记忆类型存储内容存储介质生命周期召回方式短期记忆当前会话的最近 N 轮对话内存会话结束即销毁直接拼入 prompt工作记忆当前任务相关的关键信息内存 Redis任务结束或超时按任务 ID 检索长期记忆用户偏好、历史事实、知识片段向量数据库持久化语义相似度召回短期记忆就是最常见的“对话历史”但我做了一个优化不是把所有历史都拼进去而是只保留最近 N 轮并且对每一轮做摘要压缩。比如用户和 Agent 聊了 20 轮我不会把 20 轮原文都塞进 prompt而是把前 15 轮压缩成一段摘要后 5 轮保留原文。这样既保留了上下文又控制了 token 消耗。工作记忆是任务级别的。比如用户让 Agent 帮忙订机票工作记忆里会存“出发地、目的地、日期、预算”这些槽位信息。任务完成后工作记忆可以选择性地写入长期记忆比如“用户偏好靠窗座位”。长期记忆用向量数据库存储我选的是 Milvus 的轻量版也可以用 Chroma 或 Qdrant。每次用户发消息时先用 embedding 模型把消息向量化然后在长期记忆里做相似度检索把最相关的 K 条记忆召回拼入 prompt。3.2 记忆写入策略什么时候该记什么时候不该记这是最容易踩坑的地方。我一开始的做法是“每轮对话都写入长期记忆”结果向量数据库迅速膨胀而且召回质量越来越差——因为大量无意义的对话“你好”“谢谢”“好的”也被存进去了。后来我改成基于重要性的写入策略用户明确表达的偏好“我喜欢”“我习惯”“以后都”→ 高优先级写入包含事实性信息的陈述“我的项目地址是”“我们团队有 5 个人”→ 中优先级写入任务相关的槽位信息 → 写入工作记忆任务结束后按需转长期寒暄、确认、无信息量的对话 → 不写入重要性的判断可以用一个轻量级的分类模型也可以用规则 关键词匹配。我用的是规则引擎 小模型打分的方式成本低且可控。3.3 记忆召回语义检索 时间衰减 重要性加权召回不是简单的“取相似度最高的 K 条”。我设计了一个综合打分公式score α * similarity β * importance γ * recency其中similarity是向量相似度余弦相似度范围 0~1importance是记忆写入时的重要性评分范围 0~1recency是时间衰减因子recency exp(-λ * days_since_created)λ 取 0.05 左右α、β、γ 是权重我实测下来 α0.6、β0.25、γ0.15 效果比较均衡这个公式的好处是一条很久以前但非常重要的记忆比如用户的过敏史不会因为时间久远就被完全淹没而一条最近但不重要的记忆也不会因为“新鲜”就挤掉真正有用的信息。3.4 记忆压缩与衰减别让上下文窗口爆掉上下文窗口是有限资源。我的策略是短期记忆超过 10 轮时触发压缩把最旧的 5 轮对话用一个小模型做摘要替换掉原文长期记忆超过 90 天且重要性低于阈值的做归档或删除工作记忆在任务结束后 24 小时自动过期压缩用的摘要模型不需要太强我用的是一个 7B 级别的本地模型成本几乎可以忽略。摘要的 prompt 也很简单“请用一句话概括以下对话的核心信息保留事实和偏好去掉寒暄。”实操心得记忆压缩一定要做但不要做得太激进。我试过把压缩阈值调到 5 轮结果 Agent 经常“忘记”前面刚说过的话用户体验很差。10 轮是个比较平衡的值。4. SSE 流式输出与 abort 机制的完整实现4.1 后端用 Spring WebFlux 做 SSE 流式推送后端我用的是 Spring Boot WebFlux。WebFlux 的Flux天然适合做流式输出配合ServerSentEvent可以很方便地把大模型的流式响应转发给前端。核心代码结构是这样的GetMapping(value /agent/stream, produces MediaType.TEXT_EVENT_STREAM_VALUE) public FluxServerSentEventString streamAgentResponse( RequestParam String sessionId, RequestParam String message, ServerHttpRequest request) { // 创建 abort 信号 Sinks.ManyString abortSink Sinks.many().multicast().onBackpressureBuffer(); // 监听客户端断开 request.getBody().subscribe( data - abortSink.tryEmitNext(abort), error - abortSink.tryEmitError(error) ); return agentService.execute(sessionId, message, abortSink.asFlux()) .map(chunk - ServerSentEvent.builder(chunk).build()) .onErrorResume(e - Flux.just( ServerSentEvent.builder([ERROR] e.getMessage()).build() )); }这里的关键是abortSink。当客户端断开连接或者主动发 abort 信号时abortSink会发出一个事件Agent 服务在流式处理过程中会监听这个信号一旦收到就中断对大模型的请求。4.2 前端React fetch 实现流式渲染与中断前端我没有用EventSource因为它不支持自定义请求头和 abort。我用的是fetchReadableStreamconst abortController new AbortController(); async function streamChat(message) { const response await fetch(/api/agent/stream, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ sessionId, message }), signal: abortController.signal }); const reader response.body.getReader(); const decoder new TextDecoder(); while (true) { const { done, value } await reader.read(); if (done) break; const chunk decoder.decode(value, { stream: true }); // 解析 SSE 格式并更新 UI appendToChat(chunk); } } function abortStream() { abortController.abort(); }这样用户点“停止生成”时调用abortController.abort()fetch 请求会被中断后端也会收到连接断开的信号从而停止对大模型的请求。4.3 踩坑记录SSE 空闲超时与断线重连我遇到的最头疼的问题是SSE 空闲超时。有些网关或代理会在连接空闲 60 秒后自动断开导致前端报stream disconnected before completion: idle timeout waiting for SSE。解决方案有两个后端定期发送心跳每隔 15 秒发一个:heartbeat\n\n注释行保持连接活跃。前端做断线重连检测到连接断开后自动重新发起请求并带上最后收到的事件 ID让后端从断点继续。我两个都做了。心跳用Flux.interval实现FluxServerSentEventString heartbeat Flux.interval(Duration.ofSeconds(15)) .map(i - ServerSentEvent.builder().comment(heartbeat).build()); return Flux.merge(agentStream, heartbeat);注意心跳的 comment 行不会被前端EventSource触发onmessage所以不会干扰正常的数据流。但如果你用的是 fetch ReadableStream需要自己过滤掉以:开头的行。5. MCP 工具对接与多 Agent 协作的工程实践5.1 MCP server 的接入配置与工具发现AgentScope 对 MCP 的接入方式是在配置文件里声明 server 列表mcp: servers: - name: knowledge-base url: http://localhost:8081/mcp token: ${MCP_KB_TOKEN} - name: playwright url: http://localhost:8082/mcp token: ${MCP_PW_TOKEN} - name: lanhu url: http://localhost:8083/mcp token: ${MCP_LANHU_TOKEN}Agent 启动时会自动连接这些 server拉取工具列表并生成对应的工具描述注入到 prompt 中。当模型决定调用某个工具时AgentScope 会把调用请求转发给对应的 MCP server拿到结果后再继续推理。这里有个细节工具描述的质量直接影响模型调用工具的准确率。MCP server 返回的工具描述如果太简略模型可能不知道该什么时候用。我一般会在 MCP server 端把工具描述写详细包括参数说明、使用场景、返回值格式。5.2 多 Agent 协作什么时候需要怎么拆分单 Agent 能搞定的事不要上多 Agent。我见过太多项目为了“看起来高级”硬拆成多 Agent结果通信开销比收益还大。我的判断标准是当一个 Agent 的职责超过 3 个明显不同的领域时才考虑拆分。比如我的运维助手一开始是一个 Agent 同时负责“查日志、执行命令、发通知、做报表”后来发现 prompt 越来越长模型经常搞混。拆成“日志分析 Agent”“命令执行 Agent”“通知 Agent”之后每个 Agent 的 prompt 都很聚焦准确率明显提升。AgentScope 的多 Agent 通信是通过消息队列做的。每个 Agent 有自己的inbox发送消息时指定目标 Agent 的 ID。我一般用一个“协调者 Agent”来做路由它根据用户请求判断该交给哪个专业 Agent 处理。5.3 工具调用的错误处理与重试策略工具调用失败是常态。网络抖动、MCP server 重启、参数格式错误都会导致调用失败。我的处理策略是参数错误不重试直接把错误信息返回给模型让模型修正参数后重新调用网络错误重试 2 次间隔 1 秒仍失败则返回错误超时设置 30 秒超时超时后中断并返回错误MCP server 不可用降级到备用方案比如知识库检索失败时直接用模型自身知识回答实操心得工具调用的错误信息一定要结构化包含错误码、错误描述、可能的修正建议。这样模型才能根据错误信息做出正确的下一步决策。我试过只返回“调用失败”模型完全不知道该怎么办只能反复重试同一个错误调用。6. 常见问题排查与生产环境避坑指南6.1 记忆召回不准怎么办这是最高频的问题。排查思路检查 embedding 模型是否匹配写入和召回必须用同一个 embedding 模型否则向量空间不一致相似度计算完全没意义。检查相似度阈值阈值太高会召回不到太低会召回一堆无关的。我一般从 0.7 开始调。检查记忆内容质量如果写入的记忆本身就是一堆废话召回再准也没用。回到第 3.2 节优化写入策略。检查时间衰减参数λ 太大老记忆会被过度惩罚λ 太小新记忆没有优势。0.05 是个不错的起点。6.2 SSE 流式输出卡顿或断流现象可能原因解决方案首字节延迟高模型首 token 生成慢检查模型负载考虑换更快的模型或加缓存输出到一半断了网关空闲超时加心跳前端做重连输出乱码编码问题确保前后端都用 UTF-8SSE 的Content-Type带charsetUTF-8abort 不生效后端没有监听 abort 信号检查abortSink是否正确传递到模型调用层6.3 MCP 工具调用超时或返回异常MCP server 本身也是服务也会挂。我的做法是每个 MCP server 配置独立的超时时间不要用全局超时加健康检查server 不可用时自动从工具列表中移除关键工具做降级方案非关键工具直接报错让模型换一种方式6.4 生产环境部署的注意事项记忆存储要持久化短期记忆可以放内存但长期记忆必须落盘。我用的是 Milvus PostgreSQLMilvus 存向量PostgreSQL 存元数据。SSE 连接数要限制每个 SSE 连接都会占用一个线程或一个响应式流并发高了会耗尽资源。我一般限制单实例 500 个并发 SSE 连接超了就走负载均衡。日志要脱敏Agent 的日志里会包含用户输入和模型输出可能涉及敏感信息。我在日志输出前做了一层脱敏过滤把手机号、邮箱、身份证号等自动替换成占位符。监控要到位我监控了四个核心指标——首 token 延迟、每秒输出 token 数、工具调用成功率、记忆召回命中率。任何一个异常都能第一时间发现。6.5 常见问题速查表问题排查方向快速修复Agent 不记得上一轮对话短期记忆是否被正确加载检查 sessionId 是否一致Agent 记不住用户偏好长期记忆写入策略是否触发检查重要性判断规则流式输出没有实时渲染前端是否正确解析 SSE检查Content-Type和解析逻辑工具调用返回 401MCP token 是否过期刷新 token 并重启 Agent多 Agent 消息丢失消息队列是否满增大队列容量或加背压上下文窗口爆了记忆压缩是否生效调低压缩阈值检查摘要模型7. 我在这套架构上踩过的三个大坑第一个坑是过度依赖向量检索做记忆召回。我一开始把所有记忆都向量化靠相似度召回。结果发现有些记忆是“事实型”的比如“用户的名字是张三”这种用关键词精确匹配比向量检索更靠谱。后来我改成了混合检索事实型记忆走关键词匹配经验型记忆走向量检索。第二个坑是SSE 的 abort 信号没有传递到模型调用层。前端点了停止后端也收到了 abort但模型请求还在跑白白消耗 token。后来我在模型调用的Flux上加了takeUntil(abortSink)才真正做到了中断。第三个坑是MCP server 的工具描述写得太随意。有个工具叫query描述只写了“查询数据”。模型完全不知道这个工具是查什么的经常在不该调用的时候调用。后来我把描述改成“根据用户 ID 查询订单信息返回订单号、金额、状态”调用准确率立刻上来了。这套东西我前后迭代了大概两个月从最开始只能跑 Demo到现在能稳定支撑日均几千次对话。记忆系统的效果最明显——用户反馈“它终于能记住我了”。如果你也在做类似的东西我的建议是先把记忆做扎实再考虑多 Agent 和复杂工具链。记忆是 Agent 的地基地基不稳上面盖什么都会塌。