ARTICLE DETAIL

资讯详情

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

DeepAgent 实战:SSE 流式输出已上线,长期记忆为何仍是半成品

DeepAgent 实战:SSE 流式输出已上线,长期记忆为何仍是半成品 1. 从SSE 已上线说起DeepAgent 的实时交互链路到底怎么搭DeepAgent 这个项目我断断续续跟了两周最直观的感受是流式输出这条链路已经能跑通但长期记忆那块还处在能存不能取、能取不准的半成品状态。这个判断不是拍脑袋来的而是把 SSE 通道、LangChain 的 Agent 调度、FastAPI 的后端封装、Vue 3 的前端渲染这几层拆开逐个验证之后得出的结论。先说清楚 DeepAgent 是什么。它是一个基于 LangChain 构建的智能体应用前端用 Vue 3后端用 FastAPI模型交互通过 SSEServer-Sent Events做流式输出让大模型的回答能一个字一个字地实时渲染到页面上。这个组合在当下的 AI 应用开发里非常典型——LangChain 负责 Agent 的编排和工具调用FastAPI 负责把 Agent 的能力封装成 HTTP 接口SSE 负责把生成过程推给前端Vue 3 负责把流式数据渲染成用户能看的界面。为什么是 SSE 而不是 WebSocket这是很多人第一个会问的问题。我实测下来的结论是对于一问一答、服务端单向推送这种场景SSE 比 WebSocket 更省事。WebSocket 是全双工需要处理连接握手、心跳、断线重连、消息分帧而 SSE 本质就是一个长连接的 HTTP 响应Content-Type: text/event-stream服务端持续往这个响应里写数据就行。DeepAgent 的交互模式是用户发一次请求、模型流式返回一次回答中途用户可能点停止生成abort这个场景 SSE 完全够用而且浏览器原生EventSource就支持前端代码量能少一大截。但 SSE 也有它的坑。第一个坑是它默认只支持 GET 请求而 DeepAgent 需要把用户的 prompt、会话 ID、历史上下文这些信息传给后端用 GET 把长文本塞进 URL 既不优雅也不安全。所以实际项目里通常不用原生EventSource而是用fetchReadableStream手动解析 SSE 格式。这一点后面会详细讲。第二个坑是abort 的处理。用户点了停止生成前端要能中断请求后端也要能感知到连接断开并停止调用模型否则就是在白白烧 token。这个链路涉及前端的AbortController、FastAPI 的Request.is_disconnected()、以及 LangChain 回调里的中断检查三处要配合好缺一处就会出现前端停了、后端还在跑的浪费。第三个坑也是这个项目标题里点名的——长期记忆。SSE 让当前这一轮对话体验很流畅但跨会话的记忆、用户偏好的沉淀、历史事实的召回这些属于长期记忆的范畴DeepAgent 目前只做到了把对话存进数据库但在合适的时候把合适的记忆取出来喂给模型这一步还很粗糙。这就是标题说的半成品。下面我会把这几层逐个拆开讲清楚每一层的技术选型逻辑、实操细节、以及我踩过的坑。不管你是刚接触 LangChain 的新手还是已经在做类似 Agent 项目的开发者应该都能从里面找到能直接抄作业的部分。2. SSE 流式输出的完整链路从模型 token 到页面渲染2.1 为什么放弃原生 EventSource 改用 fetch 流原生EventSource用起来确实简单const es new EventSource(/api/chat?promptxxx); es.onmessage (e) { /* 处理数据 */ };但它有三个硬伤在 DeepAgent 这种场景下直接劝退只能 GETprompt 稍微长一点 URL 就爆了而且中文需要 encode历史上下文根本没法传。不能自定义请求头如果后端需要Authorization做鉴权EventSource没法加 header。abort 能力弱只能es.close()但关闭时机和连接状态不好精确控制。所以 DeepAgent 用的是fetchReadableStream手动解析。核心代码大概长这样const controller new AbortController(); const response await fetch(/api/chat/stream, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ prompt, sessionId }), signal: controller.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 }); // SSE 以 \n\n 分隔事件 const parts buffer.split(\n\n); buffer parts.pop(); // 最后一段可能不完整留到下次 for (const part of parts) { if (part.startsWith(data: )) { const data part.slice(6); if (data [DONE]) return; // 解析并渲染 appendToMessage(JSON.parse(data).content); } } }这里有个非常容易踩的坑buffer的处理。SSE 的数据是按\n\n分隔的事件块但网络传输是分片的一次read()拿到的value很可能把一个事件块切成两半。如果你每次都直接split(\n\n)然后全部解析遇到半截的 JSON 就会JSON.parse报错。正确做法是保留最后一段不完整的 buffer等下一次 read 拼上再解析。我第一次写的时候没注意这个本地测试因为数据快没暴露一上生产环境网络抖动就疯狂报错。2.2 后端 FastAPI 怎么把 LangChain 的输出变成 SSE后端这一层FastAPI 用StreamingResponse来返回 SSE。关键是要把 LangChain Agent 的流式回调接进来。LangChain 的 Agent 执行时可以通过astream_events拿到细粒度的事件流包括 LLM 开始生成、每个 token 产出、工具调用开始/结束等。DeepAgent 里我主要关注on_chat_model_stream这个事件它对应模型吐出的每一个 tokenfrom fastapi import FastAPI, Request from fastapi.responses import StreamingResponse import json app FastAPI() async def event_generator(request: Request, prompt: str, session_id: str): try: async for event in agent.astream_events( {input: prompt}, versionv2 ): # 客户端断开就停止 if await request.is_disconnected(): break kind event[event] if kind on_chat_model_stream: chunk event[data][chunk] if chunk.content: payload json.dumps( {content: chunk.content}, ensure_asciiFalse, ) yield fdata: {payload}\n\n yield data: [DONE]\n\n except Exception as e: yield fdata: {json.dumps({error: str(e)})}\n\n app.post(/api/chat/stream) async def chat_stream(request: Request): body await request.json() return StreamingResponse( event_generator(request, body[prompt], body[sessionId]), media_typetext/event-stream, headers{ Cache-Control: no-cache, X-Accel-Buffering: no, # 关键禁用 Nginx 缓冲 }, )这里有两个细节值得单独说。第一个是X-Accel-Buffering: no。如果你的服务前面挂了 NginxNginx 默认会缓冲响应导致 SSE 的数据被攒成一大块才发出去前端看起来就是卡半天然后一次性全出来完全失去了流式的意义。加上这个 header 告诉 Nginx 不要缓冲。如果用的是其他反向代理也要检查有没有类似的缓冲配置。第二个是request.is_disconnected()。这是实现 abort 的后端侧关键。用户在前端点了停止AbortController.abort()会中断 fetch连接断开FastAPI 这边is_disconnected()就会返回 True循环 break不再继续调用模型。但要注意这个检查要放在每次事件循环里不能只在开头检查一次否则模型已经跑了一半你才发现断开token 已经烧了。2.3 abort 链路的三处配合少一处都白搭abort 这个功能看起来简单实际上要三处配合位置做什么不做的后果前端AbortController.abort()中断 fetch用户点了没反应界面卡住后端request.is_disconnected()检测并 break前端停了后端还在烧 tokenAgent 层回调里检查中断标志停止工具调用工具已经执行了副作用无法撤销第三处是最容易被忽略的。假设 Agent 正在调用一个写数据库的工具用户这时候点了停止如果工具已经执行到一半你 break 了事件循环但工具本身的副作用已经产生了。所以对于有副作用的工具要么设计成幂等可回滚要么在工具执行前就检查中断标志。DeepAgent 目前的处理是只读类工具允许执行完写类工具在执行前检查一次中断标志。这个策略不是完美的但比什么都不做强。2.4 前端 Vue 3 的流式渲染别用 v-html 直接怼前端拿到流式 token 之后渲染方式也有讲究。最朴素的做法是每来一个 token 就message.content token然后模板里{{ message.content }}。这样能跑但有两个问题Markdown 渲染模型输出的是 Markdown如果等全部生成完再渲染用户看到的是原始符号如果每个 token 都重新渲染 Markdown性能会很差Markdown 解析是 O(n)每个 token 都跑一遍就是 O(n²)。代码块高亮代码块在流式过程中是未闭合的高亮库可能会报错或渲染错乱。DeepAgent 的做法是节流渲染token 累积到一个 buffer用requestAnimationFrame或者 50ms 的节流批量更新一次 DOM。Markdown 渲染用marked配合DOMPurify做 XSS 过滤代码高亮用highlight.js但设置ignoreUnescapedHTML遇到未闭合的代码块就按纯文本渲染等闭合了再高亮。提示流式渲染时千万不要用v-html直接插入未经净化的模型输出模型完全可能输出包含script的内容这是实打实的安全风险。DOMPurify这一步不能省。3. LangChain Agent 的调度逻辑与工具设计3.1 LangChain 和 LangGraph 到底怎么选这是被问得最多的问题之一。我的理解是LangChain 是组件库 简单链LangGraph 是状态机 复杂编排。LangChain 的 Agent比如create_react_agent或者AgentExecutor适合一次思考、调用几个工具、给出答案这种线性流程。它的执行逻辑是隐式的你给它工具和 prompt它自己决定调哪个工具、调几次。LangGraph 则是把整个流程显式地建模成一张图节点是处理步骤边是流转条件状态在节点之间传递。它适合需要循环、分支、人工介入human in the loop、多 Agent 协作的场景。DeepAgent 目前用的是 LangChain 的 Agent因为核心场景还是用户提问 → Agent 决定调工具 → 返回答案。但我已经在规划往 LangGraph 迁移原因是长期记忆的读写需要更精细的控制——什么时候读记忆、什么时候写记忆、记忆冲突怎么处理这些用 LangChain 的隐式 Agent 很难精确控制用 LangGraph 画成图就清晰多了。如果你正在选型我的建议是只是做个问答 几个工具调用LangChain 够了别过度设计。需要多轮循环、条件分支、人工审核节点直接上 LangGraph别在 LangChain 上硬凑。两者不是替代关系LangGraph 里可以调用 LangChain 的组件实际项目里经常混用。3.2 工具设计的三个原则Agent 的能力边界由工具决定。DeepAgent 里我设计了几个工具知识库检索、网页抓取、计算器、时间查询。设计过程中总结出三条原则原则一工具描述要写给模型看不是写给人看。工具的description直接决定模型会不会在正确的时机调用它。我一开始写查询知识库模型经常在该调用的时候不调用。改成当用户询问产品文档、技术规范、历史资料相关的问题时使用此工具检索内部知识库输入应该是精炼的查询关键词之后调用准确率明显提升。原则二工具参数要少而精。参数越多模型填错的概率越大。如果一个工具需要 5 个参数考虑拆成两个工具或者把一些参数设成有默认值的可选参数。原则三工具要有明确的失败返回。工具执行失败时不要抛异常让整个 Agent 崩掉而是返回一个结构化的错误信息让模型知道这个工具没成功它可以选择重试或者换一个工具。比如检索工具返回空结果时返回{result: 未找到相关文档, hint: 尝试更换关键词}模型看到 hint 往往会调整查询再试一次。3.3 astream_events 的事件类型与过滤astream_events会吐出很多类型的事件如果不过滤处理起来会很乱。常用的几类事件类型含义用途on_chat_model_start模型开始生成前端显示思考中on_chat_model_stream模型吐出 token流式渲染正文on_chat_model_end模型生成结束记录 token 用量on_tool_start工具开始执行前端显示正在检索...on_tool_end工具执行结束前端隐藏工具状态on_chain_end整个链结束触发记忆写入DeepAgent 前端会根据这些事件显示不同的状态提示用户体验上比一直转圈好很多。比如检索知识库时显示正在查阅资料用户就知道系统在干活而不是卡住了。注意astream_events的version参数要显式指定v2不同版本的 LangChain 事件结构有差异不指定可能拿到旧格式。4. 长期记忆为什么是半成品存储、召回与冲突4.1 当前实现存得进去取不精准DeepAgent 现在的记忆实现是这样的每轮对话结束后把用户输入和模型回答存进数据库PostgreSQL SQLAlchemy同时把对话内容做 embedding 存进向量库Chroma。下次对话时用当前问题去向量库检索 top-k 相似的历史片段拼进 prompt。这套流程能跑但效果不稳定。问题出在几个地方问题一检索的是相似不是相关。向量相似度高不代表内容对当前问题有用。比如用户问上次那个方案改好了吗向量检索可能召回一堆包含方案这个词的历史对话但真正相关的是哪个方案——这需要实体识别和指代消解纯向量检索做不到。问题二没有记忆的时效性和重要性区分。三个月前的一句闲聊和昨天确认的一个关键决策在向量库里权重是一样的。理想情况下应该有衰减机制和重要性打分。问题三记忆冲突没有处理。用户上周说我用 Python这周说我改用 Go 了两条记忆都在库里检索时可能同时召回模型就懵了。需要有一个新记忆覆盖旧记忆或者标记冲突的机制。4.2 记忆分层的思路要解决上面的问题我倾向于把记忆分成三层这也是业界比较常见的做法短期记忆会话内当前会话的完整对话历史直接放在 context 里不需要检索。会话结束就丢弃或归档。长期事实记忆用户的偏好、身份、关键决策这类结构化信息存成 key-value 或者三元组检索时精确匹配。长期情景记忆历史对话的片段存向量库用于语义召回。DeepAgent 目前只做了第三层前两层是缺失的。这就是半成品的核心原因——只有情景记忆没有事实记忆导致 Agent 记不住你是谁你要什么只能模糊地回忆你好像说过类似的话。4.3 事实记忆的抽取与写入事实记忆的关键是从对话里抽取结构化信息。这一步可以用 LLM 来做每轮对话结束后让模型判断这轮对话里有没有值得长期记住的事实如果有抽成{subject, predicate, object}的形式。比如用户说我下个月要做一个基于 FastAPI 的项目可以抽成{ subject: user, predicate: upcoming_project_tech, object: FastAPI, valid_from: 2024-06, confidence: 0.8 }写入时先查有没有同 subject predicate 的旧记录有的话根据时间戳和 confidence 决定是覆盖还是标记冲突。这样下次用户问我那个项目用什么框架来着就能精确召回FastAPI而不是从一堆对话片段里猜。这一步的难点在于抽取的准确率和成本。每轮对话都调一次 LLM 抽取token 成本不低。DeepAgent 现在的策略是异步抽取对话结束后把任务丢进队列后台慢慢处理不阻塞用户交互。同时加了一个简单的规则过滤——太短的对话比如好的谢谢直接跳过不浪费抽取调用。4.4 记忆召回的时机与预算控制记忆召回不是越多越好。context 窗口是有限的塞太多历史反而会稀释当前问题的注意力。DeepAgent 现在的召回策略是先用当前问题做向量检索拿 top-20 候选。用一个小模型或者规则对候选做重排序选出 top-5。事实记忆单独查一次精确匹配的优先。把选出的记忆拼成一段背景信息放在 system prompt 里而不是混在对话历史里。预算上记忆部分占用的 token 控制在总 context 的 20% 以内。超过就截断优先保留事实记忆和最近的情景记忆。提示记忆召回一定要做重排序纯向量 top-k 的效果在真实场景里往往不够看。重排序可以用 cross-encoder 模型也可以用 LLM 打分成本换效果。5. FastAPI 后端的工程化细节5.1 项目目录结构FastAPI 项目最容易写乱一开始不分层后面越堆越乱。DeepAgent 用的是这样的结构app/ ├── main.py # 入口注册路由和中间件 ├── api/ │ ├── chat.py # 对话相关接口 │ └── memory.py # 记忆管理接口 ├── core/ │ ├── config.py # 配置pydantic-settings │ └── agent.py # Agent 初始化 ├── models/ │ └── conversation.py # SQLAlchemy 模型 ├── schemas/ │ └── chat.py # Pydantic 请求/响应模型 ├── services/ │ ├── memory.py # 记忆读写逻辑 │ └── retrieval.py # 检索逻辑 └── db/ └── session.py # 数据库会话管理分层的核心原则是api 层只做参数校验和响应封装业务逻辑放 services数据访问放 models/db。这样测试的时候可以单独测 service不用起 HTTP 服务。5.2 SQLAlchemy 的异步会话管理FastAPI 是异步的数据库访问也要用异步否则会阻塞事件循环。SQLAlchemy 2.0 支持 async配置大概是这样from sqlalchemy.ext.asyncio import ( create_async_engine, async_sessionmaker, AsyncSession ) engine create_async_engine( postgresqlasyncpg://user:passlocalhost/deepagent, pool_size10, max_overflow20, pool_pre_pingTrue, ) AsyncSessionLocal async_sessionmaker( engine, class_AsyncSession, expire_on_commitFalse ) async def get_db(): async with AsyncSessionLocal() as session: yield session几个参数值得说pool_pre_pingTrue会在每次取连接前 ping 一下避免拿到已经断开的连接数据库重启后特别有用expire_on_commitFalse让 commit 后对象属性还能访问不然会触发额外的查询。5.3 SSE 接口的超时与异常处理SSE 是长连接如果模型生成很慢或者卡住连接会一直挂着。要设置合理的超时import asyncio async def event_generator(request, prompt, session_id): try: async with asyncio.timeout(120): # 整体超时 120 秒 async for event in agent.astream_events(...): ... except asyncio.TimeoutError: yield fdata: {json.dumps({error: 生成超时})}\n\n except Exception as e: logger.exception(stream error) yield fdata: {json.dumps({error: 服务异常})}\n\n异常处理里有个细节不要把原始异常信息直接返回给前端可能包含数据库连接串、文件路径等敏感信息。返回一个通用的错误提示详细信息记到日志里。6. 踩坑实录那些文档里不会写的问题6.1 中文乱码ensure_ascii 的坑json.dumps默认ensure_asciiTrue中文会被转义成\uXXXX。虽然前端JSON.parse能正确还原但如果你在浏览器 Network 面板看 SSE 流会看到一堆转义字符调试起来很痛苦。加上ensure_asciiFalse就能直接看到中文。6.2 代理缓冲导致流式失效前面提过X-Accel-Buffering: no但如果你用的是云厂商的负载均衡可能还有一层缓冲。排查方法是直接访问后端端口看流式是否正常。如果直连正常、走代理不正常那就是代理层的缓冲问题需要查代理的配置文档。6.3 向量库的持久化Chroma 默认是内存模式重启就丢数据。生产环境一定要用持久化模式import chromadb client chromadb.PersistentClient(path./chroma_data)而且要注意多个进程同时访问同一个 Chroma 持久化目录会有锁冲突。如果 FastAPI 用多 worker 启动要么每个 worker 用独立的目录要么把向量库访问收敛到单独的服务。6.4 记忆写入的并发问题异步抽取记忆的任务如果并发写同一条记录可能产生冲突。DeepAgent 用的是数据库的唯一约束 upsert 来处理(user_id, subject, predicate)建唯一索引写入时用ON CONFLICT DO UPDATE保证同一时刻只有一条有效记录。6.5 token 用量统计的准确性LangChain 的on_chat_model_end事件里能拿到 token 用量但流式模式下有些模型不返回准确的 usage。DeepAgent 的做法是流式时用估算按字符数除以系数非流式时用真实 usage两者都记对账的时候以真实值为准。7. 后续可以怎么补从半成品到可用长期记忆这块我接下来的计划是分三步走第一步把事实记忆的抽取和写入做扎实。现在抽取是异步的但抽取质量还没系统评估过。打算建一个小测试集人工标注哪些对话应该抽出什么事实然后跑抽取看准确率和召回率迭代 prompt。第二步引入记忆的时效衰减和冲突解决。给每条记忆加last_accessed和access_count检索时综合考虑相似度、时效、访问频率。冲突解决用时间戳优先 显式确认的策略新记忆和旧记忆冲突时如果新记忆时间更近且 confidence 够高就覆盖否则标记为待确认下次对话时主动问用户。第三步把整个记忆流程用 LangGraph 重构成显式的图。读记忆、判断是否需要记忆、抽取、写入、冲突处理每个环节做成一个节点状态在节点间流转。这样调试的时候能看到每一步的输入输出比现在藏在 Agent 内部的黑盒清晰得多。SSE 这条链路目前是稳的abort 也验证过没问题。真正花时间的还是记忆——流式输出解决的是体验记忆解决的是智能。一个记不住东西的 Agent体验再流畅也只是个复读机。这也是为什么我把标题写成SSE 已上线长期记忆还是半成品——前半句是完成时后半句是进行时这个项目真正的硬骨头在后面。
返回列表