SSE流式输出与Agent长期记忆:从管道优化到记忆闭环
发布时间:2026/9/24 23:57:44 锦皓数字建站

1. 上线两周后SSE 带来的流畅感与记忆功能的短板DeepAgent 的 SSE 流式输出上线两周了体验上确实是一次明显的跨越。之前大模型回答要等全部生成完后一次性返回哪怕问一句今天天气怎么样也要先看三五秒的加载圈用客户的话说是像在等一个慢性子的人说话。切到 SSE 之后第一个 token 基本在 1 秒内就能出现在屏幕上后面每个字都是流出来的。首字延迟从 3~5 秒降到了 0.8~1.5 秒整个产品的交互手感完全不一样。但作为一个每天都泡在 code review 和日志里的人我对这个版本的评价只能说半场胜利。原因在于SSE 解决的回答能不能实时渲染只是管道层问题而一个 Agent 产品真正值钱的长期记忆能力到现在还停在半成品状态——能写入一些用户偏好但召回质量、更新机制、使用效果都不达预期。也正是这种一头通一头堵的现状让我决定把这一段的实战经验完整记录下来。这篇文章我记得比较散干脆按两条线整理先讲 SSE 流式输出的完整落地过程包括为什么选 SSE、怎么和 LangGraph 交互、abort 取消怎么做以及上线后遇到的 idle timeout 断流问题后一半重新梳理 Agent 记忆体系的分层设计最后复盘长期记忆为什么还没做完。适合正在做 Agent 服务、或者刚把 LLM 应用从一次性返回改成流式对话的开发者参考。1.1 SSE 上线带来的直接变化从数据上看SSE 上线后最明显的变化有三个首 token 延迟中位数在 1.2 秒左右此前是完整响应等待平均要 8 秒以上用户主动中断生成的比例明显增加因为大家知道回答可以随时停下来重新问前端正在输入的展示不再靠定时器猜状态而是直接渲染真实的 token 流这些变化指向同一个结论流式输出不是一个炫技功能而是 Agent 对话产品的基本盘。没有它再强的模型能力也会被等待拖垮。用户不会关心你的 Agent 内部是不是有复杂的工具调用链他只会觉得回答得慢或者回答得快。1.2 半成品长期记忆的具体表现再说长期记忆。目前仓库里长期记忆模块的状态我给自己打 4 分10 分制写入链路已经通了每轮对话结束后会有一条抽取任务跑出来把用户偏好摘要写进 SQLite召回链路只有最粗糙的版本新会话开始前把该用户最近 20 条记忆全量拼进 system prompt没有去重、没有时效衰减、没有冲突处理、没有用户反馈入口所以与其说上线了长期记忆不如说打通了记忆写入的管道但记忆读取和使用的循环还没闭合。这种状态在 demo 里看着像那么回事真正放进生产环境就露馅了。2. SSE 流式输出落地的完整链路从选型到 LangGraph 封装2.1 为什么没有用 WebSocket对话场景看起来双向通信很合理我最初也倾向于 WebSocket但仔细评估之后放弃了。我们实际的通信模式是用户请求一次服务器持续推送响应上行只有一次下行是持续的流。WebSocket 的双向能力在这个场景里基本用不上反而带来一堆额外成本。SSE 这边的好处是很实在的SSE 基于 HTTP可以很自然地放进现有网关、负载均衡和监控体系WebSocket 要额外处理连接升级、粘性会话、心跳保活、跨域、连接池管理等问题浏览器原生 EventSource 支持断线重连对长对话场景很友好。虽然我实际没用 EventSource后面会讲原因但这条协议能力本身是加分项服务端用 FastAPI 一个 StreamingResponse 就能搞定不需要维护独立的连接管理器维度SSEWebSocket通信方向服务端单向推送双向全双工协议基础HTTP独立升级协议断线重连EventSource 原生支持需自己实现鉴权方式HTTP headers / cookie子协议或首包处理服务端实现成本StreamingResponse 即可需要连接管理器适用场景模型回答流式输出实时白板、音视频信令WebSocket 真正适合的场景是上行也要实时流式比如用户持续发语音、画布协作、多人协同编辑这类。纯聊天对话SSE 是更轻的选择。2.2 与 LangGraph 的交互astream_events 是最实用的入口DeepAgent 的编排层是基于 LangGraph 实现的两者的交互走的是 LangGraph 编译后的图对象上的流式接口。LangGraph 提供了 stream / astream / astream_events 等几个方法我最终选了astream_events并且指定versionv2。原因很简单astream_events会把你精心编排的图里每个节点的事件全部平铺出来包括模型的 token 流、工具调用开始/结束、节点完成等。这样我在下游做 SSE 转发时既能推送 token 级内容也能推送工具调用的状态给前端展示比如正在搜索资料这类过程反馈。import json from langchain_core.messages import HumanMessage class DeepAgentService: def __init__(self, compiled_graph, memory_serviceNone): self.graph compiled_graph self.memory memory_service async def stream_events(self, thread_id: str, message: str): config {configurable: {thread_id: thread_id}} inputs {messages: [HumanMessage(contentmessage)]} async for event in self.graph.astream_events( inputs, configconfig, versionv2 ): etype event[event] if etype on_chat_model_stream: chunk event[data][chunk] if chunk.content: yield {type: token, content: chunk.content} elif etype on_tool_start: yield {type: tool, name: event[name], phase: start} elif etype on_tool_end: yield {type: tool, name: event[name], phase: end}这里有一个容易踩的细节不同模型供应商的流式 chunk 字段可能不一样。我实际测试过有些模型在on_chat_model_stream事件里除了主文本内容字段还会把 reasoning 内容、工具调用参数等放在其他字段。如果你只关注content可能会漏消息或者把不该展示的内部推理内容一起渲染到前端。所以封装层这里需要做一层字段白名单过滤确定哪些字段需要透出给 SSE。2.3 FastAPI 封装 SSE 端点心跳、断连检测与取消服务端我用 FastAPI 的 StreamingResponse 作为 SSE 出口。最开始的实现很天真直接在 async generator 里遍历agent.stream_events()每拿到一个事件就 yield 一行 SSE。上线后才发现这个实现有个致命问题当 agent 长时间不产出事件时连接就僵在那里了。比如模型在做推理、工具在等外部返回这段时间可能 20 秒甚至 1 分钟没有任何 token但 HTTP 连接还挂着。如果代理层或浏览器做了空闲超时连接会被静默断开。后面我会专门讲这个问题。先给出一个更健壮的封装核心思路是用 asyncio.Queue 把事件生产和SSE 输出解耦然后用 wait_for 的空闲超时被动触发心跳from fastapi import FastAPI, Request from fastapi.responses import StreamingResponse import asyncio, json, uuid app FastAPI() agent DeepAgentService(compiled_graph) app.post(/api/agent/stream) async def agent_stream(request: Request): payload await request.json() thread_id payload.get(thread_id) or ft_{uuid.uuid4().hex} message payload.get(message, ) async def gen(): queue: asyncio.Queue asyncio.Queue() async def producer(): try: async for event in agent.stream_events(thread_id, message): await queue.put(event) except Exception as exc: await queue.put({type: error, message: str(exc)}) finally: await queue.put(None) producer_task asyncio.create_task(producer()) try: while True: try: event await asyncio.wait_for(queue.get(), timeout15) except asyncio.TimeoutError: yield : ping\n\n continue if event is None: break if await request.is_disconnected(): break yield fdata: {json.dumps(event, ensure_asciiFalse)}\n\n finally: producer_task.cancel() await agent.cancel_current_run() return StreamingResponse( gen(), media_typetext/event-stream, headers{ Cache-Control: no-cache, Connection: keep-alive, X-Accel-Buffering: no, }, )这个实现里有三处容易漏但每一处都直接影响线上稳定性第一X-Accel-Buffering: no必须加上。如果不加nginx 默认会缓冲响应前端看到的不是一个字一个字地出现而是一坨一坨地出现体感跟非流式没什么区别。第二Queue wait_for 做心跳的方式比在事件循环里主动 sleep 更符合实际。因为真正的问题是事件源长时间没有产出你不能靠每 N 秒主动发一个心跳这种插队方式解决那会打乱事件顺序。wait_for 超时意味着当前这一刻确实没有新事件此时发一个注释行正好。第三request.is_disconnected()检测客户端断连。配合前端的 abort一旦用户停止生成这个函数返回 True循环 break我们就不再往下游要资源。2.4 前端消费为什么用 fetch ReadableStream 而不是 EventSourceSSE 协议本身很成熟但浏览器的 EventSource API 有一个让我绕不开的限制只支持 GET 请求且不能自定义请求头。我们的接口需要带 Authorization headertoken 放 query 里不仅丑而且不安全所以我直接用 fetch ReadableStream 手写了一段 SSE 解析。const controller new AbortController(); async function sendMessage(message: string, onEvent: (e: any) void) { const res await fetch(/api/agent/stream, { method: POST, headers: { Content-Type: application/json, Authorization: Bearer ${getToken()}, }, body: JSON.stringify({ thread_id, message }), signal: controller.signal, }); if (!res.ok || !res.body) { throw new Error(stream request failed: ${res.status}); } const reader res.body.getReader(); const decoder new TextDecoder(); let buffer ; while (true) { const { value, done } await reader.read(); if (done) break; buffer decoder.decode(value, { stream: true }); const blocks buffer.split(\n\n); buffer blocks.pop() ?? ; for (const block of blocks) { for (const line of block.split(\n)) { if (!line.startsWith(data:)) continue; const raw line.slice(5).trim(); if (!raw) continue; onEvent(JSON.parse(raw)); } } } }abort 的核心就一个 AbortController用户点击停止生成按钮时调用controller.abort()浏览器会中断 fetchreader.read()返回或者抛异常服务端request.is_disconnected()检测到连接断开会主动停止 LangGraph 的生成我实测过如果前端只 abort、后端不做取消下游 LLM API 的 tokens 还会继续烧因为语言模型的生成任务没人叫停。所以前端 abort 通知用户和后端取消止损必须双向配合缺一环都是浪费。3. 断流重连的排查记idle timeout 是怎么被干掉的3.1 现场现象stream disconnected before completion上线之后我们陆续收到反馈在 agent 执行工具调用的场景里尤其是搜索工具跑到第 4、5 步的时候页面会突然没有任何反馈过一会儿前端提示连接断开。去看服务端日志发现异常信息是stream disconnected before completion: idle timeout waiting for sse这句话直译过来是SSE 流在完成之前被断开了原因是空闲超时——在设定的超时时间内连接上没有任何数据流动。这个报错第一次出现在日志里的时候我以为是某个第三方库的问题。但持续观察后确认这是整个链路里某一层判定连接空闲时间过长后主动断开的表现。3.2 排查链路从网络层到业务层我的排查过程完整记录如下希望能帮后面的人少走弯路。第一步本地复现。我写了一个最小脚本直接调 FastAPI 端点发现本地很快就能看到 token不会断。所以问题大概率不在业务代码本身而在部署链路。第二步看 nginx 配置。当时proxy_read_timeout用的是默认 60 秒。Agent 场景下如果工具链路很长某一步的 LLM 调用到第一个 token 回来可能就要 20~40 秒再叠加多个工具调用之间的间隙很容易超过 60 秒没有任何数据推送。我把 nginx 的proxy_read_timeout先调到 300 秒情况缓解了一部分但没有根治。第三步抓事件流时间线。我在流式接口里把每个 data 事件的耗时都打进日志发现一个规律断开瞬间往往出现在一次长时间的无事件之后。比如 agent 在调用搜索工具工具执行完但还没把结果喂给下一轮模型之前中间有 30 秒到 1 分钟是完全安静的状态。这段时间里连接上没有任何数据流动所以不管超时阈值调到多高只要用户侧网关或浏览器自己的空闲策略更保守依然会断。第四步确认根因。核心问题不是代理超时时间不够大而是SSE 应用层没有保活机制。SSE 协议本身并不要求服务器定期发数据但现实中的通道持有者nginx、云负载均衡、浏览器往往会基于一段时间无流量去回收连接。3.3 修复方案与验证修复分两层。第一层是应用层心跳。让 SSE 连接在空闲时也能向通道发送连接还活着的信号。具体做法就是前面代码里那个 15 秒的注释行心跳每 15 秒如果没有业务事件就发送一个: ping\n\nSSE 解析逻辑会把它当作注释行忽略对浏览器、nginx 来说这代表连接上有数据流动不会判定为空闲第二层是代理层配合。nginx 需要同时调整几项配置location /api/agent/stream { proxy_pass http://upstream; proxy_http_version 1.1; proxy_set_header Connection ; proxy_buffering off; proxy_cache off; proxy_read_timeout 300s; proxy_send_timeout 300s; chunked_transfer_encoding on; }几个关键点proxy_buffering off关掉缓冲否则 nginx 会攒一批数据才发给客户端proxy_set_header Connection 对 HTTP/1.1 保持长连接避免每次转发都重新建立上游连接proxy_read_timeout 300s虽然有心跳机制但仍然留一个宽裕的业务超时作为兜底验证的时候我和前端同学配合做了个实验手动让 agent 进入只思考不出 token的状态保持 2 分钟无业务事件。修复前连接会在 60 秒左右被断掉修复后连接一直保持等到模型恢复输出后token 流继续正常走到了 done。这个坑填完之后线上再没有出现过 idle timeout 断流。4. Agent 记忆的分层设计短期、长期、永久各管什么4.1 先区分消息记录不等同于记忆很多做 Agent 的朋友会把把对话历史存进数据库说成我有记忆功能这是常见的混淆。LangGraph 的 checkpointer 确实可以保存每个 thread 的完整 state包括消息列表这保证对话中断后能恢复上下文。但它不能回答一个跨会话的问题用户三天前明确说过我偏好用表格形式呈现结果今天的新会话怎么知道这件事这类信息不在当前会话的消息列表里但又是真实影响回答质量的长期事实。所以记忆体系的本质是从对话数据中沉淀出可跨会话复用的知识。它不能只靠存下消息实现需要一套独立的抽取、存储、检索和更新机制。4.2 短期记忆会话窗口与上下文压缩短期记忆承载的是当前这轮对话进行中的信息。我在 DeepAgent 里的实现方式比较朴素但有效原始消息列表保持全量交给 LangGraph 的 state 管理当 token 总量接近阈值时触发上下文压缩把较早的旧消息交给一个摘要模型生成一段对话摘要并替换掉原文最近的几条消息保留原文这样既保留当前讨论重点又不至于让 token 无限膨胀这里分享一个参数经验压缩触发的阈值不能只看模型窗口的绝对上限要留出足够余量给 agent 的工具返回结果和系统提示词。比如模型窗口是 128ksystem prompt 加工具定义可能占 20k工具返回结果占 20k 以上那消息历史最多留到 60k 左右就该考虑压缩否则某个大工具结果刚回来就爆窗。另外短期记忆并不需要做得特别智能。很多人一上来就设计复杂的摘要策略结果摘要比原文还绕。我的建议是小对话直接滑动窗口长对话才启用摘要压缩先用最简单可靠的方式跑起来。4.3 长期记忆跨会话的知识抽取与召回长期记忆是我认为 Agent 产品最值得投入的部分。它的职责是把对话中出现的用户偏好、个人事实、项目信息、承诺与待办等沉淀到独立存储在新会话中按需召回并注入上下文。设计上分为四个环节记忆抽取对话过程中或结束后用 LLM 判断这一轮有没有值得存的信息输出结构化记忆条目记忆存储结构化字段进关系型数据库需要语义检索的部分进向量库记忆注入新会话开始时检索与该用户当前问题相关的记忆择优注入 system prompt记忆更新用户纠正或新信息覆盖旧信息时做合并、作废或版本更新抽取这一步我踩过的坑是不要把所有对话内容都当作记忆。记忆应该是有损的压缩而不是无损的备份。我用 JSON Schema 约束抽取结果只保留四类用户偏好、个人事实、项目信息、承诺与待办。{ memory_items: [ { category: preference, content: 用户希望回答尽量用列表呈现, confidence: 0.9, source_message_id: msg_12345 } ] }这四类的区分不是为了分表而是为了后续更新和召回时能按类型处理。比如承诺与待办可能在下次对话开始时优先召回而用户偏好要在回答语气上持续生效。4.4 永久记忆不是固定存储而是闭环终点永久记忆在语义上很像人脑中的长期记忆固化。我理解的永久记忆不是永远不会改的数据而是经过验证、可靠性较高、有版本与时间戳可追溯的记忆。举一个实际例子用户说我喜欢简洁回答 → 这是一个偏好过了几天用户说算了你还是给我详细一点我要写报告 → 旧偏好应被更新或标记失效永久记忆的关键不是存下来而是可演进。我在存储层为每条记忆维护了创建时间、更新时间、来源消息 ID、置信度等字段后续做冲突消解时会依赖这些元数据。三类记忆放在一起看应该是这样的分层层级数据范围生命周期存储载体当前状态短期记忆当前会话消息一次会话LangGraph checkpointer / Redis可用长期记忆跨会话的用户事实与偏好数周至数月SQLite/Postgres 向量库半成品永久记忆经过校验的稳定事实长期可演进事实表 版本字段设计阶段5. 长期记忆的半成品复盘写入通了召回还差得远5.1 目前已经做到的部分先客观盘点。当前长期记忆模块做到的程度写入对话结束后会触发一个记忆抽取任务用 LLM 判断本段对话中是否有值得记录的条目有则按 JSON 结构写入 SQLite基础召回新会话开始时把该用户最近 20 条记忆全部拼进 system prompt存储一张 memory_items 表字段包括 uid、thread_id、category、content、created_at看起来最核心的链路都有了那为什么还说半成品因为它只是能跑离能用差得很远。demo 模式下一个用户聊两轮你给他看我还记得你喜欢简洁回答效果很好但放到连续使用一周的真实用户身上问题就全出来了。5.2 为什么说它只是半成品问题集中暴露在召回侧具体有四点。第一全量注入导致 token 膨胀和注意力漂移。把最近 20 条记忆全部塞进 system prompt假设每条平均 80 字20 条就是 1600 字。这些不一定和当前问题相关的上下文会分散模型注意力甚至让模型优先回忆记忆里的内容而不是认真回答用户当前的问题。第二没有去重和合并。我们的库里已经出现了同一个用户 17 条喜欢简洁回答的记录来源是 17 次对话里用户反复表达同一个偏好。每次抽取任务都像一个第一次见面的人一样重新记录完全不知道库里已经有同样语义的条目。第三没有更新和失效机制。用户说以后不用给我表格了这条新信息本应该覆盖之前偏好表格的旧条目。但当前实现里新旧条目会同时存在互相矛盾的记忆一起进 prompt模型根本不知道该信哪个。第四没有用户反馈闭环。用户可以对回答点不对但系统不会把这个信号接回记忆存储。也就是说记忆系统在写入—使用之间少了一个校验—修正—淘汰的循环。这四点合起来就是它撑不起长期记忆这四个字的原因。用工程上的话说输入侧通了输出侧通了但缺少一个持续运转的质量控制回路。5.3 后续迭代计划针对上述问题我目前的迭代优先级是这样的召回层改造把最近 20 条全量注入改成按相关性检索 相关度阈值过滤只有和当前 query 相关的记忆才进入 system prompt记忆合并任务定期扫描同一 uid 下相似度高的记忆条目做去重与融合比如把 N 条喜欢简洁回答合并成一条并累计置信度冲突消解与失效策略当同类别新旧记忆冲突时以更新时间戳为准把旧条目标记为 superseded反馈回路用户对回答点纠正时生成一条记忆修订事件由抽取任务决定是否更新或作废某个记忆条目效果评估建立一个小型评估集对比是否注入记忆、注入哪些记忆时的回答质量差异这个迭代计划看起来还有不少工作量但每一步都不需要特别花哨的技术难的是把闭环里每个环节都做得够扎实。从这段实践里我最大的体会是SSE 再复杂本质上是协议 传输问题花一周时间总能踩平长期记忆则完全不同它没有一个固定正确的实现写入、召回、融合、淘汰每一步都直接影响最终质量。如果你也在做 Agent 的记忆模块建议不要急着把写入做得花哨先把召回—验证—更新这条闭环里的每一环都跑通哪怕每一步都朴素一点也比一个看起来很智能但不可控的写入器强。
锦
锦皓数字建站
深耕本土企业品牌数字化升级,专注原创端正雅致商务官网,从视觉设计到稳定运维全程保驾护航。