ARTICLE DETAIL

资讯详情

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

DeepAgent实战:SSE流式输出与Agent长期记忆体系设计拆解

DeepAgent实战:SSE流式输出与Agent长期记忆体系设计拆解 DeepAgent 的 SSE 流式输出上线跑了一阵子整体链路算是通了但长期记忆这块我评估下来仍然是个半成品。这篇文章把这次实战的完整过程拆开讲清楚SSE 怎么接、Abort 怎么处理、记忆体系怎么设计、以及为什么说长期记忆还差得远。内容偏工程落地向适合正在做 AI Agent 应用、需要把大模型交互封装成标准化服务的开发者参考。先说结论如果你只是想快速把 LLM 的流式输出接到前端DeepAgent 这套 SSE 方案可以直接抄作业但如果你指望它有成熟的长期记忆能力那还要自己动手补不少东西。1. 项目背景与技术选型思路1.1 DeepAgent 到底是什么场景的产物DeepAgent 本质上是一个面向多智能体协作场景的应用编排层底层对接大模型能力向上提供统一的业务接口。和直接调 OpenAI API 不同它更像一个中间件帮你把 Prompt 管理、工具调用、多轮会话、流式输出这些繁琐的逻辑收敛到一层业务侧只需要关心我要 Agent 做什么而不需要关心Agent 怎么和模型通信。我在这个项目里负责的是交互链路这块核心目标有两点。第一让大模型的回答能像打字机一样逐字出现而不是让用户对着空白页面等三五秒第二让 Agent 的状态能被追踪也就是多轮对话里它记得什么、不记得什么这个数学模型必须清楚。这两点分别指向了 SSE 和记忆体系也是这次标题里两个关键词的由来。技术选型上我们前后对比过几条路。最原始的方案是前端轮询也就是发起请求后每隔 500ms 去查一次结果但这样 WebSocket 的优势体现不出来反而增加服务端压力。第二条路是 WebSocket 全双工通信好处是消息可以双向推但 Agent 场景里用户发消息的频率远低于模型回消息的频率全双工的价值并不大而且 WebSocket 在网关层、鉴权层都要单独做适配成本和复杂度不成比例。最终定了 SSE单向流式推送恰好匹配用户发一次模型流式回多次的交互模型。1.2 为什么 SSE 比 WebSocket 更适合 Agent 场景很多人一听到实时就想到 WebSocket但这里有个误区Agent 交互的实时性需求是单向的。用户在输入框敲完回车那一刻请求已经提交了后续真正需要实时更新的只有模型的输出流。这叫请求-流式响应模型SSE 天生就是干这个的。SSE 基于纯 HTTP意味着它可以复用现有的负载均衡、鉴权中间件、日志链路不需要额外维护一个长连接网关。你可以在 Nginx 层面直接给它配缓冲关闭、超时调优也可以在应用层用最普通的 WSGI/ASGI 服务器跑起来。对部署和运维来说这比 WebSocket 少了一整个维度的麻烦。还有一个工程上的隐性收益SSE 天然支持断线重连。浏览器原生的 EventSource API 自带 reconnection 机制而 WebSocket 挂了之后要自己实现心跳和重连策略。对 Agent 这种动辄十几秒的推理过程来说中间网络抖一下是常有的事SSE 的重连机制能省掉不少脏活。2. SSE 流式输出实战拆解2.1 服务端怎么把 LLM 的流式输出转发给前端核心思路用一个代码块就能说清楚大模型 SDK 返回的是一个生成器我们把这个生成器里的每一段增量文本包装成 SSE 格式的data:事件一次一个 event 推到客户端。from fastapi import FastAPI from fastapi.responses import StreamingResponse from llm_sdk import chat_stream app FastAPI() app.post(/v1/agent/chat) async def agent_chat(request: ChatRequest): # 这里拿到的是一个异步生成器每次 yield 一小段文本 generator chat_stream( messagesrequest.messages, toolsrequest.tools, agent_idrequest.agent_id, ) async def event_stream(): # 先发一个 session 开始事件方便前端做 loading 收敛 yield event: session_start\ndata: {}\n\n async for chunk in generator: # chunk 可能是文本增量也可能是工具调用的元信息 yield fevent: delta\ndata: {chunk.model_dump_json()}\n\n # 流结束后发一个 done 标记 yield event: done\ndata: {}\n\n return StreamingResponse( event_stream(), media_typetext/event-stream, headers{ Cache-Control: no-cache, Connection: keep-alive, X-Accel-Buffering: no, }, )这里有个细节值得注意X-Accel-Buffering: no这个响应头。如果你们的服务前面挂了 Nginx默认情况下 Nginx 会缓冲后端响应攒够一定量才发给客户端这样前端看到的就是卡一下然后一整段出现流式效果直接没了。加上这个头就是为了告诉 Nginx 别缓冲数据到了就往外吐。2.2 前端怎么优雅处理 Abort 和异常中断SSE 接入前端有两种方式。一种是浏览器原生 EventSource但 EventSource 只支持 GET 请求没法带 Authorization header对需要鉴权的业务很不友好。另一种是用 fetch 加 ReadableStream 手动解析这也是我在项目中采用的方式。async function chatWithAgent(messages, { onDelta, onDone, onError }) { const controller new AbortController(); // 把 controller 暴露出去组件卸载或者用户点了“停止生成”时调用 pendingController controller; try { const response await fetch(/v1/agent/chat, { method: POST, headers: { Content-Type: application/json, Authorization: Bearer ${token}, }, body: JSON.stringify({ messages }), signal: controller.signal, }); if (!response.ok) throw new Error(HTTP ${response.status}); const reader response.body.getReader(); const decoder new TextDecoder(utf-8); let buffer ; while (true) { const { done, value } await reader.read(); if (done) break; buffer decoder.decode(value, { stream: true }); // SSE 事件之间以空行分隔按这个规则解析 const events buffer.split(\n\n); buffer events.pop(); for (const rawEvent of events) { const lines rawEvent.split(\n); const eventType lines.find(l l.startsWith(event:))?.slice(6).trim(); const data lines.find(l l.startsWith(data:))?.slice(5).trim(); if (eventType delta) { onDelta(JSON.parse(data)); } else if (eventType done) { onDone(); } } } } catch (err) { if (err.name AbortError) { // 用户主动取消不算错误 onAbort?.(); } else { onError(err); } } }这段代码里最容易踩坑的是TextDecoder的{ stream: true }参数。SSE 数据流是字节级别的一个中文字符可能被拆在两个 chunk 里如果不加 stream 模式解码时就会偶尔出现乱码。加上这个参数之后解码器会把没凑完整的字符留在缓冲区等下一段字节到了再拼起来。关于 Abort 还有个很容易被忽略的点客户端abort()之后服务端其实不一定立刻感知到。Python 的异步生成器里如果模型还在推理StreamingResponse不会因为你断开了连接就自动取消模型调用。这意味着用户取消了生成但后端还在占着模型资源继续算多来几个这样的请求就能把算力池拖垮。解决方案是在生成器里监听客户端断开事件async def event_stream(): try: yield event: session_start\ndata: {}\n\n async for chunk in generator: yield fevent: delta\ndata: {chunk.model_dump_json()}\n\n yield event: done\ndata: {}\n\n except asyncio.CancelledError: # 客户端断开时StreamingResponse 会取消这个生成器 # 在这里手动去 cancel 底层模型调用 await generator.aclose() raise这块我在初版实现里漏了压测时发现前端点了停止后端日志里模型调用还在跑浪费了不少 token。后来加上aclose()才把这个资源泄漏问题解决。2.3 流式输出的超时与空闲断开问题热搜词里有一条 stream disconnected before completion: idle timeout waiting for SSE这个报错我实测中遇到过。它的大致场景是模型在某个工具调用节点上思考了比较久中间有一段时间没有输出任何 SSE 事件结果网关层的 idle timeout 触发了连接被强制断开前端收到一个不完整的流。排查思路是这样先看你的服务链路里有哪些组件有超时配置。通常有三个位置会影响 SSE位置默认超时调整建议负载均衡/网关30-60s调到 300s 以上或按最大工具调用链评估应用服务器通常不限制确认反向代理转发时没有覆盖超时客户端 fetch不超时手动实现超时控制兜底异常场景Nginx 配置里有两个参数最要命proxy_read_timeout和proxy_send_timeout。默认值是 60s如果你的模型工具调用链超过这个时间连接必断。实测中我把它调到了 300slocation /v1/agent/ { proxy_pass http://agent_service; proxy_buffering off; proxy_cache off; proxy_read_timeout 300s; proxy_send_timeout 300s; chunked_transfer_encoding off; }另外一个容易翻车的地方是模型在思考期间真的一个字节都不发。为了让中间链路知道这个连接还活着很多实现会用 SSE 的注释行做心跳。注释行形如: ping\n\n浏览器和代理都会把它当成活跃信号但不会触发表层事件。我建议在流式生成里搞一个后台心跳任务每 15s 发一次注释行async def event_stream(): # 用一个 task 定期发心跳 async def heartbeat(): while True: await asyncio.sleep(15) yield : heartbeat\n\n # 把 heartbeat 和主生成器合并输出这个方案不算复杂但能显著减少因为网络静默导致的误杀。3. 从需求到落地记忆体系的演化路径3.1 Agent 记忆体系的四层结构热搜词里提到agent 记忆体系中短期、长期、永久记忆如何实现这个问题我被问了很多次。说清楚这件事之前先建立一个统一的框架。我把 Agent 记忆按时间尺度和持久化策略分成四层短期记忆当前对话窗口内的上下文直接拼在 Prompt 里送给模型。它的特点是容量小、变化快模型每次调用都重新传一遍。对应到大模型的世界里就是 context window 内的 messages 列表。短期记忆的实现最简单但也是最容易失控的——上下文一旦超过模型的 context window要么截断要么报错。中期记忆跨对话但不过长时间的记忆比如用户在一个工作会话里的历史操作、偏好变化。实践中通常用摘要的方式管理也就是把一段对话压缩成几百字的摘要存下来下次对话时作为系统提示的一部分注入。中期记忆的核心是压缩因为原始的对话内容太多不可能全量传给模型。长期记忆用户跨会话、跨天甚至跨月的稳定偏好、事实信息。比如用户是 Python 开发者、偏好简洁回答、上次讨论过某个项目。这类信息适合结构化存储比如放 SQLite、PostgreSQL 或者向量数据库在对话开始时按需检索注入。长期记忆的难点是什么时候写入和什么时候读取写早了会记录噪音读多了会污染上下文。永久记忆用户画像级别的数据不随对话变化比如账号信息、权限范围、不可修改的基础属性。这种一般直接落在业务数据库跟 Agent 系统解耦需要时通过工具调用查出来。这个四层结构不是拍脑袋定的。你去看 LangGraph、AutoGen 这些框架的记忆模块本质上都在做这四件事只是叫法不同。理解了分层逻辑再去看具体框架的 API 就会豁然开朗。3.2 DeepAgent 当前记忆模块的实现现状回到 DeepAgent 本身。我评估完它的记忆模块给的结论是短期和中期基本可用长期是个半成品。为什么这么说短期记忆方面DeepAgent 维护了一个 messages 数组每次调用模型时完整传入。这个逻辑没毛病问题在于它没有处理上下文超出窗口的情况。我实测发了几十轮对话之后服务端直接报了 context length exceeded而且错误信息里没有任何提示告诉你该滚动了或该压缩了。这说明框架把这个责任完全扔给了上层调用者。中期记忆方面DeepAgent 提供了会话摘要能力基于一个简单的策略对话轮数超过阈值时让模型对已有内容做总结用摘要替换掉最早的一部分消息。我测下来功能是能跑的但触发的时机比较死板。它只按轮数触发不看实际 token 占用导致在长回复场景下摘要还没触发上下文就已经超限了。长期记忆这块标题里说它半成品一点不冤枉。DeepAgent 目前的实现只有一个memory字段挂在一个 agent 实例上进程一重启里面的内容全没了。没有持久化层没有向量检索没有独立的存取服务。如果你要做多实例部署这个 memory 字段连分布式共享都做不到每个 worker 各存各的用户会话被路由到不同实例时记忆完全是错乱的。我给它的定位是提供了一个记忆接口的骨架但工程上能用的长期记忆能力需要自己动手补全。3.3 基于 LangGraph 的交互机制分析热搜词里有一条很有意思deepagent 和 langaph 通过什么方法交互。这里langaph我理解是 LangGraph 的笔误。目前市面上确实有不少团队在用 LangGraph 编排复杂 Agent 流程同时用 DeepAgent 这类框架做应用封装两者的通信机制值得聊一聊。LangGraph 的核心抽象是图节点是处理步骤边是状态流转路径。它本身不关心你用什么协议调用外部服务所以 DeepAgent 和 LangGraph 的集成方式完全取决于你把它放在哪个位置。我实际采用的方式是用 LangGraph 做顶层编排把 DeepAgent 封装成一个 Tool 节点。LangGraph 根据用户的输入决定要不要调用 DeepAgent调用时通过 HTTP 请求把参数传过去DeepAgent 返回的结果作为 Tool 的响应回到 LangGraph 的 state 里。还有一种反向集成DeepAgent 作为入口内部通过 LangGraph 的 API 定义子图。这种场景下通信机制一般是函数调用而非 HTTP因为两个框架在同一个进程内。DeepAgent 的execute_graph接口接收 LangGraph 的 state 对象内部走完图之后把最终 state 返回。不管哪种方式关键点在于状态如何传递。LangGraph 的 state 是线程私有的跨图传递时需要注意可序列化。我的经验是在两个框架的边界处定义一个明确的 DTO数据传输对象只传必要字段不要把整个 state 直接塞过去。曾经踩过坑尝试把 LangGraph 的整个 state 对象 JSON 序列化后传给 DeepAgent结果里面有不可序列化的自定义类直接 500。3.4 长期记忆落地的补全方案既然 DeepAgent 的长期记忆是半成品那就自己动手补。我的思路是把长期记忆拆成两层事实型记忆和语义型记忆。事实型记忆用一个标准的 KV 存储或者关系型表搞定。表结构大概是用户 ID、key、value、更新时间、来源会话 ID。查询时直接按用户 ID 拉取全部记忆注入系统提示。优点是简单、可控、可审计缺点是没法做语义匹配——你没法问用户之前有没有提过和部署相关的事情然后模糊命中。语义型记忆用向量数据库。对话过程中把关键信息切片、embedding、入库新对话开始时用当前问题去检索 top-k 相关的记忆片段。这种做法对用户提过类似问题的场景特别有效但需要额外的 embedding 服务和向量库运维成本。我的建议是预算有限时先用事实型 KV 记忆把语义型记忆作为迭代计划。原因是事实型记忆的效果确定性更强出现问题也好排查向量召回如果质量不好你说不清是 embedding 的问题、切分的问题还是检索参数的问题排查成本会高得多。4. 常见问题与排查技巧实录4.1 SSE 断了、前端不渲染、数据乱码排查表这段时间踩过的 SSE 相关坑整理成一张表希望帮你少走弯路现象直接原因排查/解决前端长时间无输出然后一瞬间全部出现Nginx 缓冲未关闭加X-Accel-Buffering: no响应头或关proxy_buffering流中途断开报 idle timeout网关超时配置过短调大proxy_read_timeout增加心跳注释行中文内容偶尔乱码TextDecoder 未开 stream 模式使用new TextDecoder(utf-8, { stream: true })用户停止后模型还在跑未监听断开取消生成在生成器里捕获asyncio.CancelledError并关闭底层生成器前端收到不完整 JSON事件按\n\n分割时跨包先用缓冲区累积按分割符切分尾段保留到下次多个事件拼接在一起每次 read 返回多个 SSE 事件循环处理当前 buffer 里的所有事件不要只取第一个排查 SSE 问题时有个实用技巧直接用 curl 不带-N参数看报错用curl -N看流式输出。如果 curl 能正常一个 chunk 一个 chunk 地出前端却不行问题大概率在前端解析如果 curl 也卡住或者断了问题就在后端或中间链路。4.2 记忆体系实测中的典型故障记忆相关的坑比 SSE 更难定位因为它的错误往往是逻辑正确但结果不对而不是程序直接报错。举一个我实际遇到的例子。用户问 Agent我上次让你总结的那篇关于 RAG 的文章核心观点是什么 理想情况下Agent 应该从长期记忆里检索出用户上次讨论过 RAG 文章这个事实然后基于记忆内容回答。但实际跑起来Agent 直接回复我没有找到相关记录。排查后发现原因在于记忆写入的时机错了。DeepAgent 的实现是在会话结束时整体写入记忆结果进程崩溃或者会话超时被清理时记忆根本没落库。用户下次来问自然什么都没有。这个问题的根因是写入时机没有保证不是检索逻辑不对。解决思路有两个层面。第一是写入时机下沉把提取记忆并写入这个动作从会话结束挪到每一轮模型调用之后这样即使会话中途挂掉至少已经完成的对话内容会被存下来。第二是加写入确认机制记忆写入后返回一个确认标识在日志里能追踪这条记忆是哪个会话、哪个时间点写入的。另一个高频问题是记忆注入位置。长期记忆如果一股脑全塞进系统提示的最前面模型容易被大量历史事实干扰影响对当前问题的注意力。我实践下来的做法是把记忆按相关度排序只注入 top 5 条事实每条用一句话表达放在系统提示的末尾紧挨着用户消息。这样既给了模型参考又不会让它淹没在陈年旧事里。4.3 生产环境中的稳定性兜底策略SSE 生产环境最大的敌人是所有东西看起来都对但就是偶尔断流。这种随机性问题没有银弹只能用多层次的兜底来对冲。我在服务端做了三个兜底。第一是流式响应的超时熔断单次 SSE 连接超过 5 分钟强制断开防止模型卡死导致连接被无限占用。第二是错误事件注入生成过程中如果模型调用抛异常不直接中断连接而是发一个event: error事件前端收到后可以给用户展示这次回答生成失败请重试而不是让人面对一个毫无反馈的加载转圈。第三是断点续传的简化版前端在收到 done 事件后才把内容写入本地状态在这之前如果断了用户刷新后回看看到的还是上一轮完整回答不会出现一截截残文。前端层面我加了一个指数退避的重试策略。前 5 秒内断开不重试因为大概率是网络抖动用户自己会刷新超过 5 秒断开每 5 秒重试一次最多 3 次。重试时带上当前已渲染内容的长度服务端如果解析到这个参数可以选择跳过已经输出的部分从断点继续生成。这个功能我还没完全实现目前是重试时从头生成但至少用户不用手动刷新页面了。5. 后续演进方向与个人经验总结5.1 记忆模块下一步的两个优先改造项长期记忆要从不合格到能用我的优先级排序是先做持久化再做召回质量最后才是多模态扩展。持久化是最急迫的。目前 memory 字段在内存里进程重启就丢这连半成品都算不上只能算 demo。我的计划是引入 SQLite 做本地持久化把 memory 结构化成表以 user_id agent_id 作为联合主键。这样至少做到进程重启不丢数据多实例部署时再考虑引入 Redis 做共享存储。召回质量方面先用 BM25 这类传统的稀疏检索试试。不要一上来就上向量数据库因为 embedding 的维护成本高且对短期迭代不友好。BM25 基于关键词匹配对大部分用户之前提过 XXX这类事实性召回已经够用。等数据量积累到一定程度确实出现语义匹配需求了再切换到向量检索方案。5.2 从这次实战里沉淀下来的几条经验第一SSE 不是越底层越好。直接用 fetch ReadableStream 是最灵活但也最容易出 bug。如果你们的场景简单不需要自定义 header 和 Abort 控制那原生 EventSource 反而更靠谱省得自己处理断线重连和事件解析。第二记忆体系的设计要奔着可观测去。每一层记忆的写入、读取、命中情况都要有日志否则线上出了问题你根本无从下手。我在项目里给每次记忆查询都加了 trace_id能顺着一条完整链路看到当前问题是什么、检索到了哪些记忆、最终注入了几条、模型有没有引用这些记忆。这个投入大大缩短了排查问题的平均时间。第三框架封装得再好核心逻辑还是得自己吃透。DeepAgent 把 SSE 和会话管理的骨架搭好了但记忆这块基本是留白。接手这样的项目不要想着找个框架一步到位而是先把这个领域的基本原理摸清然后针对自己的场景做定制。框架负责省时间你负责正确性。这次实战最大的感受是AI Agent 的工程化流的实时性只是入门记忆才决定体验的上限。SSE 是管道管道通了只是开始记忆是大脑大脑还没发育好之前Agent 充其量是个带打字效果的聊天机器人。后续我会持续跟进长期记忆的迭代等这块完善了再写一篇更深的拆解。
返回列表