SSE流式输出实战:从协议原理到LangChain与前端接入全解
1. 一场线上事故SSE 连接被掐断的前因后果先从我上个月踩的一个坑讲起。公司内部做了一个基于 LangChain 的工业知识库问答系统前端用 Vue后端是 FastAPI。上线第二天运维就甩过来一条报错stream disconnected before completion: idle timeout waiting for sse。当时我的第一反应是网络问题连夜查 Nginx 配置调了各种超时参数结果故障还是偶发。后来把整个链路捋了一遍才意识到问题根本不在网络层而是我自己压根没搞懂 SSE 的传输模型。这里给还不熟悉的朋友补个背景。SSEServer-Sent Events是 HTML5 标准里定义的服务器推送方案它建立在普通 HTTP 之上服务端通过Content-Type: text/event-stream持续往客户端写数据客户端用EventSource这个原生 API 就能接收。整个过程是单向的服务器可以持续往客户端推数据但客户端没法通过同一条连接往服务器发业务数据。注意建立在普通 HTTP 之上这句话后面所有坑都从这里来。很多做 AI 应用的朋友会犯一个认知错误觉得 SSE 就是把接口返回改成一段一段往外吐就行了。实际上普通流式响应比如 FastAPI 直接返回文本流和真正的 SSE 是两种完全不同的东西——前者只是分块传输后者是事件协议。这篇文章我就把这件事一次讲透从 SSE 协议本身到前端 EventSource 的正确姿势再到后端 LangChain 流式链路怎么对接最后附上我在实际项目中踩过的坑和排查清单。适合谁看所有在做聊天式 AI 应用、实时推送、大模型流式输出的后端、前端同学。就算你暂时没碰到问题把这套逻辑理顺了以后遇事能少走很多弯路。2. SSE 协议拆解它不是流式响应是事件流2.1 协议格式每一帧都在说什么SSE 的传输格式看起来很像纯文本但它有严格的行协议。服务端发送的一条事件由若干字段行和一个空行组成每个字段的格式是字段名: 值。我直接给一段最小示例event: message data: {type: token, content: 你好} data: 这是一条没有事件名的消息 id: 42 retry: 3000 data: 带 ID 的消息逐行解释一下event:定义事件类型客户端可以通过addEventListener监听对应事件名。如果省略默认就是message事件。data:是事件的实际内容可以有多行多行data会被拼接成一个完整数据换行符保留。id:给事件一个标识符。客户端断开重连后会通过请求头Last-Event-ID把它带回去服务端据此决定补发哪些数据。retry:指定断开后自动重连的等待毫秒数默认是 3000ms。以冒号开头的行是注释比如: ping用于心跳保持连接。这里有个关键的规范细节每条事件的结束必须是两个连续换行也就是事件之间用一个空行分隔。很多初级做 SSE 接口的人会在每行只加一个\n结果前端收到的数据黏成一团解析全部错位。正确做法是每条事件末尾输出\n\n如果数据是多行也要在最后补上完整的空行才广播出去。2.2 普通流式响应和 SSE 到底差在哪很多人的困惑是我用StreamingResponse直接返回一段不停变化的文本前端用fetch配合ReadableStream去读不也能做到打字机效果吗为什么非要用 SSE技术上确实都能出效果但两者定位不同普通流式响应身体是一大块未知内容客户端拿到的是字节流要自己按字符或字节切分。没有事件语义不知道这段数据代表开始还是结束断线后也没有重试机制更谈不上回放未读事件。SSE身体是一组结构化事件每个事件带着类型、编码好的 JSON、可选 ID。客户端可以针对不同事件分别处理连接意外断开后EventSource 默认自动重连还能带上Last-Event-ID拉取补偿数据。简单类比普通流式像是往纸箱里不停塞货收货人得自己开箱清点SSE 是每一件货都贴着标签、写明品类收货人按标签分类入库即可。所以如果你只是自己内部调试普通流式能凑合但一旦面对生产环境、多端对接、断线恢复这些真实场景老老实实按 SSE 协议去做才是省心的选择。2.3 SSE 和 WebSocket不是替代关系这个话题在热搜里常年出现。很多初学者容易把 SSE 当成低配版 WebSocket这是误解。简单说两个东西解决的问题不一样对比维度SSEWebSocket方向服务器到客户端单向全双工双向底层协议普通 HTTP独立的 WS 协议重连机制内置自动重连需要自己实现消息格式纯文本可自定义字段文本或二进制帧浏览器限制HTTP/1.1 下同域约 6 条连接无此限制实现成本极高几行就能跑需要握手、心跳、帧解析适用场景推送通知、LLM 流式输出聊天、在线协作、实时游戏AI 对话场景我几乎都推荐 SSE原因很实际大模型输出本质是服务端生成、客户端消费的单向流客户端不需要往服务端持续推消息最多发个停止请求断开连接就够了。用 WebSocket 不仅复杂度高还要单独维护心跳和重连属于典型的杀鸡用牛刀。3. 前端接入EventSource 的坑比想象中多3.1 EventSource 的 GET 限制与 fetch 替代方案原生 EventSource 是最简单的 SSE 客户端const es new EventSource(/api/chat/stream?id123); es.onmessage (event) { console.log(JSON.parse(event.data)); }; es.addEventListener(error, () { // 出错时 EventSource 会自动重连 });但它在实际 AI 项目里有一个绕不开的限制EventSource 只支持 GET 请求没法自定义请求头也没法传 POST body。而大模型聊天往往需要把用户的 Prompt、历史消息、参数配置都放在请求体里用 GET 走 query 参数既不安全长度也有限。生产中我一般直接用 fetch ReadableStream 来实现 SSE 客户端。以 Vue 3 聊天场景为例async function streamChat(messages, onToken, onDone) { const response await fetch(/api/chat/stream, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ messages }) }); if (!response.ok || !response.body) { 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 }); // 按空行拆分成完整事件 const events buffer.split(\n\n); buffer events.pop(); // 最后一段可能是不完整的事件留在缓冲区 for (const rawEvent of events) { const event parseSSEEvent(rawEvent); if (!event) continue; if (event.event token) { onToken(JSON.parse(event.data)); } else if (event.event done) { onDone(); } else if (event.event error) { console.error(服务端错误, event.data); } } } } function parseSSEEvent(raw) { const lines raw.split(\n); let event message; let data ; for (const line of lines) { if (line.startsWith(event:)) { event line.slice(6).trim(); } else if (line.startsWith(data:)) { data line.slice(5).trimStart(); } } if (!data) return null; return { event, data }; }这段代码看起来简单但有两个细节决定生死必须保留不完整事件。reader.read()返回的字节边界和数据边界不一定是重合的一个事件可能被拆到两次读取里。如果每次都直接处理buffer你会把半个事件当完整事件去解析。缓冲区设计是我反复提醒团队的重点。decoder.decode(value, { stream: true })要一直开着。否则中文或其他多字节字符恰好被截断时会产生乱码。这个参数告诉解码器后面还有数据先把当前字节缓存住。3.2 标签返回未完整怎么处理标签返回未完整怎么处理这个热搜词几乎每个做 AI 流式前端的人都遇到过。大模型生成 Markdown 或 HTML 时常常是这个样子流式到达**这是一个加粗的 标题**等等还没完。如果前端实时把这段 Markdown 渲染成 HTML就会出现**开头但没闭合的瞬间闪烁或者更糟——直接把断开的 HTML 标签插入 DOM页面结构直接乱掉。处理思路有三种按项目复杂度选第一种暴力简单型在渲染前做延迟缓冲比如 500ms 内不更新视图等 500ms 后再把最新获取的内容一次性渲染。缺点非常明显流式输出的即时感被牺牲了体验大打折扣。只适合内容不长、对实时性要求低的场景。第二种增量修补型前端保留上一次的完整渲染结果每次拿到新片段时从完整内容而不是增量片段重新渲染。也就是说后端每次发来的数据不一定是新增 token而是截至当前的完整文本。这样渲染器永远面对一个相对完整的文本进行 Markdown 解析残缺标签的渲染闪烁会少很多但副作用是内容越长重复解析性能越差。第三种状态化渲染型把 Markdown 解析器放到 Web Worker 里每次新 token 到达都发送全量文本由 Worker 解析后只更新干净的 HTML 片段。这种方式我已经在生产环境用了很久配合diff算法只更新变化部分几乎感觉不到性能损耗也彻底规避了标签残缺问题。我个人建议除非你只是做个 Demo否则别用第一种。第二种适合轻量场景第三种工作量稍大但体验和维护成本最可控。4. 后端链路LangChain 流式输出到 SSE 的正确姿势4.1 为什么 LangChain 的 stream 不能直接往前端丢现在主流的大模型应用尤其是知识库问答和智能体后端基本都套了一层 LangChain。很多同学把chain.stream()拿到 token 后直接yield token当成响应返回前端也确实能看到打字机效果——但问题在上面已经说过没有事件结构没有结束标志断线无法恢复。正确做法是在 LangChain 拿到 token 后把 token 包装成标准 SSE 事件再发送。框架层面FastAPI 配合StreamingResponse写起来非常直接import json from fastapi import FastAPI from fastapi.responses import StreamingResponse from langchain_openai import ChatOpenAI from langchain_core.prompts import ChatPromptTemplate app FastAPI() llm ChatOpenAI(modelgpt-4o, temperature0.7) prompt ChatPromptTemplate.from_template(你是工业设备运维专家请回答{question}) chain prompt | llm async def event_generator(question: str): # 发送起始事件让前端可以展示正在思考 yield fevent: start\ndata: {json.dumps({status: thinking})}\n\n # 流式获取 LLM token async for chunk in chain.astream({question: question}): content chunk.content if content: # 每得到一个 token包装为一个 token 事件 yield fevent: token\ndata: {json.dumps({content: content})}\n\n # 发送结束事件 yield fevent: done\ndata: {json.dumps({status: ok})}\n\n app.post(/api/chat/stream) async def chat_stream(request: dict): question request.get(question, ) return StreamingResponse( event_generator(question), media_typetext/event-stream, headers{ Cache-Control: no-cache, X-Accel-Buffering: no, # 关键关闭 Nginx 缓冲 } )这里有三个点值得单独拿出来说。第一media_typetext/event-stream是 SSE 协议的标志少了它浏览器不会按事件流解析。第二event: start、event: done这些自定义事件可以给前端明确的阶段信号而不是让前端在超时后靠猜来判断是不是结束了。第三X-Accel-Buffering: no这个响应头是给 Nginx 这类反向代理看的告诉它这个响应别给我缓冲实时转发。我在生产环境见过太多人漏掉这一行结果 Nginx 把 SSE 数据攒到一大坨才吐出去前端看着像假流式。4.2 用 astream_events 处理复杂链路的流式输出上面的例子是最简单的单模型链。但真实项目里链路上经常有检索、工具调用、多个模型节点比如 RAG 里的召回 → 重排 → 生成。这时候astream已经不够用了因为它只输出最终结果。用astream_events可以拿到链路上每一个节点的完整事件流from langchain_core.callbacks import AsyncCallbackHandler from langchain_core.callbacks import adispatch_custom_event async def event_generator(question: str): # 通过配置 LLM开启 verbose 级别事件 async for event in chain.astream_events( {question: question}, versionv2, ): kind event[event] if kind on_chat_model_stream: chunk event[data][chunk] token getattr(chunk, content, ) if token: yield fevent: token\ndata: {json.dumps({content: token})}\n\n elif kind on_retriever_end: # 检索完成时前端可以展示已找到相关文档 yield fevent: status\ndata: {json.dumps({message: 检索完成开始生成})}\n\n elif kind on_tool_end: # 工具执行完成 yield fevent: status\ndata: {json.dumps({message: 工具调用完成})}\n\n elif kind on_chain_end: yield fevent: status\ndata: {json.dumps({message: 任务结束})}\n\nastream_events的事件类型前缀是on_模型、链、检索器、工具都会有对应的事件。前端拿到这些事件后可以非常细腻地渲染当前进行到哪一步——这在工业智能体、复杂工作流类产品里是体验加分项用户不再对着空白页面干等大模型出第一个字。但要注意开启astream_events后事件量会成倍增长如果链路里有大量并行节点后端 CPU 和网络开销都不小。生产环境建议在迭代到稳定阶段后只保留必要的on_chat_model_stream事件把其他节点事件降级为 debug 级别避免给前端刷屏。4.3 JDBC 查询流式输出一个容易被忽略的相似场景热搜词里有jdbc查询流式输出虽然它不是 SSE 的直接范畴但在架构思路上高度相关。很多报表系统在查大表时习惯用Statement.executeQuery()一次性拉全量结果数据量大时 JVM 直接 OOM。JDBC 本身支持流式读取方法因数据库而异MySQLstatement.setFetchSize(Integer.MIN_VALUE)配合statement.setFetchDirection(FETCH_FORWARD)让 ResultSet 一行一行从服务端拉取而不是一次性全部载入内存。PostgreSQLconnection.setAutoCommit(false)statement.setFetchSize(1000)。Oraclestatement.setFetchSize(100)即可本来就按游标方式取数。这跟 SSE 共享同一个哲学能流式传输的数据不要整包缓冲。数据从数据库到后端、从后端到前端的每一跳都应该保持流式形态否则在任一跳上做全量缓冲都会把前面的流式优化抵消掉。我见过一个系统JDBC 这边已经改好流式了结果后端组装 JSON 时用StringBuilder堆了全量数据再返回前功尽弃内存照样爆。做完整链路流式化时一定要检查每个环节是否都在流式地处理。5. 再进一步LangGraph 与 Human-in-the-loop 下的流式设计5.1 LangChain 和 LangGraph 到底怎么选这个也是热搜高频问题。简单说LangChain 是链式思维LangGraph 是图式思维。LangChain面向线性的任务编排A → B → C一条路走到底。适合绝大多数 RAG、问答、简单工具调用场景。LangGraph面向有状态、有分支、需要回环的 Agent 场景。比如先规划、再执行、发现结果不对要回退重新规划甚至需要暂停等待人工审批。在流式输出这件事上两者也有明显差异。LangChain 的流式输出相对简单token 按顺序往外吐就行。LangGraph 因为要支持节点之间的状态流转、条件分支和回环它的流式设计需要考虑每个节点的输出以及节点之间的状态迁移早期版本在这方面确实拉胯很多智能体项目都得自己拼装节点事件。新版本已经内置了stream_mode[updates, messages]这类模式可以同时输出节点状态更新和 token 流但整体心智负担比 LangChain 重不少。我的选型建议很直接如果业务只是输入 → 检索 → 生成用 LangChain 就好别为图结构过度设计如果业务有明确的暂停/恢复人工确认多轮工具调用后重新规划等控制流需求上 LangGraph 是合理的但要留足调试时间。5.2 Human-in-the-loop 与流式输出的组合拳Human-in-the-loop人在回路是智能体落地时绕不开的一环。以工业场景为例Agent 自动生成设备维护工单时不能直接下发给维修系统得先推到人工审核台确认无误再执行。这时候流式输出和人工介入怎么共存我这边的一个参考实现是这样的from langgraph.graph import StateGraph, END from langgraph.checkpoint.memory import MemorySaver # 定义状态 class AgentState(TypedDict): question: str draft: str human_approved: bool # 生成草稿节点这里开启流式输出到前端 def generate_draft(state: AgentState): # 实际场景中用 LLM 生成并把 token 通过 SSE 推给前端预览 return {draft: 根据巡检数据建议更换传感器模块...} # 人工审核节点 def human_review(state: AgentState): # 暂停在图节点等待外部调用 resum 恢复 return {human_approved: True} graph StateGraph(AgentState) graph.add_node(generate, generate_draft) graph.add_node(review, human_review) graph.add_edge(generate, review) graph.add_edge(review, END) app graph.compile(checkpointerMemorySaver())LangGraph 的interrupt()机制可以将执行暂停在某个节点等待前端收到 SSE 事件后展示草稿人工点击通过/驳回再通过graph.resume()恢复流程。这个模式的完整链路是Agent 执行到生成草稿节点通过 SSE 把草稿内容流式推给前端用户能实时看到 AI 在写什么。草稿生成完毕Agent 进入interrupt()暂停点此时图的状态会被检查点机制持久化。前端收到 SSE 的 waiting human review 事件展示通过/驳回按钮。用户点击后前端调用后端接口后端拿到用户决策调用graph.resume(thread_id, {human_approved: True})恢复执行。这里面最容易被忽略的是检查点持久化。MemorySaver只在内存里保存状态进程一重启所有暂停中的流程就丢了。生产环境要找专门的外部存储检查点实现比如数据库或 Redis 版本否则一台后端实例重启所有待审核工单全部断掉这种事故我在早期 demo 版本里差点踩过。5.3 工业智能体场景的流式注意事项说一个我们实际落地工业知识库时踩到的教训。现场环境网络不稳定车间无线网络经常断SSE 连接说断就断。这时候 EventSource 自动重连 Last-Event-ID的机制帮了大忙但前提是后端必须正确解析Last-Event-ID并做补发。实现思路是前端每次发送的 SSE 事件都带上id字段ID 用一个单调递增的序号后端从request.headers.get(Last-Event-ID)拿到断点 ID如果发现客户端漏了序号从断点后的第一条事件重新推起。这套机制最怕的是事件 ID 不连续所以我建议 ID 由后端统一生成用一个自增序列就好别用时间戳——毫秒级时间戳在并发下可能一样事件去重会出问题。另外工业场景的现场监控大屏经常同时开着好几条 SSE 连接设备状态、告警、AI 助手注意 HTTP/1.1 下同域连接数上限的问题。能合并的通道尽量合并比如把所有推送事件都塞进同一条 SSE 流的自定义事件里前端按event类型分流这个方案可以顶住很长一段时间。6. 踩坑实录常见问题与排查技巧6.1 在线事故速查表把我在项目里碰到过的典型问题整理成了表格方便大家直接对照现象常见原因解决方案idle timeout waiting for sse中间层空闲超时Nginx/网关/负载均衡或后端长时间没吐数据缩短心跳间隔每 1530 秒发一条注释事件调大proxy_read_timeout前端收到的 token一坨一坨地出现Nginx 开启了缓冲加proxy_buffering off或响应头加X-Accel-Buffering: no事件数据黏连JSON.parse报错事件之间没有用空行分隔或只用了单个\n严格规范每条事件以\n\n结尾中文乱码TextDecoder未用{stream: true}或多字节字符被截断保持 stream 模式解码缓冲区保留不完整字节连接数达到上限新连接被挂起HTTP/1.1 同域 6 连接限制合并 SSE 通道或升级 HTTP/2或换子域名断线重连后数据重复未使用Last-Event-ID或 ID 不连续后端生成单调递增 ID解析请求头的Last-Event-ID流式响应中间报ChunkedEncodingError后端在流式过程中抛了异常响应被强制截断在生成器外层捕获异常把错误信息包装成 SSE 错误事件前端偶尔收不到最后一条done代理层把 EOF 当成连接断开丢弃了尾部数据在done事件后再发一条注释事件作为flush 标记6.2 排查流程从浏览器到网关再到后端遇到 SSE 问题我习惯按这条链路从上到下排查先在浏览器DevTools → Network里找到这条 SSE 请求看它的响应内容是不是真的在一条一条地蹦出来。如果 Network 面板能看到数据持续追加但页面没反应问题在前端解析。如果 Network 面板里数据是攒了一大批才出现那就是中间层缓冲的问题重点查 Nginx 和网关。接着看网关。除了上面说的proxy_buffering还要确认网关有没有对连接做空闲超时。我遇到过一个诡异案例SSE 已经正常推流了但只要某次模型思考时间超过 60 秒第一个 token 迟迟没来网关就把连接掐了。解决方案是后端要尽早发前缀事件比如event: start和 正在检索 状态事件让连接始终有流量经过这比单纯调大超时参数更可靠。最后看后端。无论是 LangChain 还是原生 LLM 调用streamTrue参数都要确保传对另外确认是否在生成器内部做了不必要的await asyncio.sleep()或串行阻塞操作。FastAPI 的StreamingResponse是异步的生成器不能写成同步版本否则事件会被阻塞住表现为前端一顿一顿。6.3 两个值得长期保留的调试技巧第一个是写一个极简的 SSE 测试端点。排查问题时先用它排除业务代码干扰确认链路是通的app.get(/api/sse-test) async def sse_test(): async def gen(): for i in range(10): yield fid: {i}\ndata: {json.dumps({seq: i})}\n\n await asyncio.sleep(0.5) yield event: done\ndata: {}\n\n return StreamingResponse(gen(), media_typetext/event-stream)前端用任意一个在线 SSE 测试工具或者直接 EventSource 连一下如果它正常就说明基础设施没问题业务代码的问题可以缩小范围去查。如果它也不正常那基本是网络或代理层的问题。第二个是用curl -N直接从命令行观察原始事件流curl -N -X POST http://localhost:8000/api/chat/stream \ -H Content-Type: application/json \ -d {question: 你好}-N参数会禁用 curl 的缓冲看到什么就输出什么。这样能确认后端吐出来的协议格式是否符合规范注意观察事件之间是否有空行、event:和data:字段是否完整。6.4 最后再说一个规范层面的教训我在复盘第一次线上事故时发现自己最大的问题还不是配置而是把 SSE 当成了普通流式来处理——没有事件的边界没有协议语义断线重连全凭运气。现在团队里我定了两个硬性规范所有流式接口必须返回text/event-stream所有 SSE 事件必须带id和event字段缺一不可。这看起来只是加两个字段的小事但在多端对接、前端换人维护、第三方系统接入时节省的沟通成本是巨大的。我个人在实际操作中的体会是SSE 的门槛真的不高但它是一个协议不是功能。协议就意味着有边界、有格式、有约定照着规范做就能避免绝大多数问题。如果你现在正被各种流式问题折磨先别急着改配置回头把协议层面的东西捋一遍往往问题自己就浮出来了。
上一篇/下一篇内容由系统自动关联
返回资讯列表 →