1. 为什么选 SSE 而不是 WebSocket:先把传输层讲透
说实话,我最早做 AI 应用流式输出的时候,第一反应也是 WebSocket,毕竟这玩意儿在前端圈子里名声大,全双工、低延迟,听起来就是为实时通信量身定做的。但真正把"打字机效果"落到生产环境之后,我反而把 WebSocket 换成了 SSE(Server-Sent Events),这一换,省掉的麻烦不是一星半点。先别急着下结论,我们把这俩协议的核心差异捋一遍。
WebSocket 的本质是一条双向长连接通道,客户端和服务端随时可以往对方那边推数据。这看起来很美好,但它带来两个问题:一是连接状态需要自己维护,断线重连、心跳保活、消息序号这些逻辑全得手写;二是中间如果挂了 Nginx、网关之类的代理层,空闲连接很容易被掐断,尤其是在移动网络环境下,运营商对空闲长连接的回收简直无情。我在生产环境遇到过"stream disconnected before completion: idle timeout waiting for sse"这类报错,追根溯源,问题都出在连接保活和代理超时配置上,这部分后面单开一节细说。
SSE 则完全换了一套思路,它是基于 HTTP 的单项推送通道,服务端可以持续不断地把数据吐给客户端,但客户端不需要维持一个独立于 HTTP 语义的连接。关键在于,SSE 走的是普通的 HTTP 响应流,一次请求,响应体被无限拉长,数据以text/event-stream的 MIME 类型分块传输。因为底层还是标准 HTTP,所以 Nginx、CDN、浏览器全都天然兼容,断线重连也是协议自带的(retry字段 + 浏览器自动重连),这省掉了整条自研保活链路。
还有一个很多人忽略的细节:SSE 对 AI 流式输出场景的契合度是结构性的。大模型生成文本本身是"单向的",服务端把 token 推给客户端展示,并不需要客户端频繁往服务端发指令——就算要反馈,单独走一次普通 POST 请求就完事了。既然业务模型是单向的,非要用全双工协议,就是在给自己找麻烦。
从实际操作成本看,SSE 的实现门槛也低得多。服务端只需要把Content-Type设为text/event-stream,然后按格式写数据块就行;浏览器端原生EventSourceAPI 就能接收,但如果要往请求头里塞认证 token,就必须改用fetch手动解析流(EventSource不支持自定义请求头)。这一整套链路,不需要引入任何第三方依赖。我在项目里最初用的是EventSource,后来因为要带Authorization头才切到fetch流式读取,实测下来的体感是:切换成本很低,但灵活性提升是质的。
当然,SSE 也不是万能药。如果业务需要客户端主动往同一个连接里塞数据(比如协作编辑、实时游戏),那必须上 WebSocket。但在"大模型生成文本 + 前端逐字展示"这个场景里,SSE 就是最省心、最不容易出幺蛾子的方案。
2. 后端怎么把 LangChain 的 token 变成 SSE 事件流
这一节是核心,因为我发现很多人卡在"模型能流式输出了,但前端拿到的是乱码/空白"这个阶段。问题不在于模型怎么配,而在于你根本没有把 token 流转化成 SSE 事件格式。
2.1 LangChain 侧的流式化基础
先明确一点:LangChain 的invoke是阻塞式调用,等模型全部生成完才返回结果,这不是我们要的。流式输出要用stream方法,或者直接在 LLM 实例上指定streaming=True并挂回调处理器。我用得最多的是stream方法,它在 LCEL 表达式中天然支持,返回的是一个迭代器,每个 chunk 对应模型的一个增量输出。举个最小例子:
from langchain_openai import ChatOpenAI llm = ChatOpenAI( model="gpt-4o-mini", temperature=0.4, streaming=True ) async for chunk in llm.astream("用一句话解释什么是流式传输"): print(chunk.content, end="")注意这里用的是astream,配合异步 SSE 才能把后端的响应能力榨干。如果你用的是stream(同步版),也可以,但 FastAPI 是异步框架,同步调用会阻塞事件循环——轻则影响并发,重则直接把服务拖垮。我的实测数据是:同一个模型、同样 500 个 token 的输出,异步流式处理耗时比同步快约 30%,并发场景下的差距更大。
2.2 FastAPI 作为中转层:把迭代器翻译成 SSE 帧
LangChain 侧拿到的是 token 迭代器,这玩意儿不能直接发给浏览器,必须翻译成 SSE 的帧格式。SSE 帧的格式极其简单,每帧由若干字段组成,字段和值之间用冒号分隔,字段间以换行符\n分割,帧与帧之间用空行\n\n分割。最常见的字段就是data:,甚至可以只写data:一个字段。
以 FastAPI 为例,做一个 SSE 中转层并不需要sse-starlette这类库,原生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( model="gpt-4o-mini", temperature=0.4, streaming=True ) prompt = ChatPromptTemplate.from_messages([ ("system", "你是一个专业的技术顾问,回答需要结构清晰、条理分明。"), ("human", "{question}") ]) chain = prompt | llm def generate_events(question: str): yield f"event: start\ndata: {{\"type\":\"start\"}}\n\n" full_text = "" for chunk in chain.stream({"question": question}): token = chunk.content if token: full_text += token payload = json.dumps({"type": "token", "content": token}, ensure_ascii=False) yield f"data: {payload}\n\n" final_payload = json.dumps({"type": "done", "full_text": full_text}, ensure_ascii=False) yield f"data: {final_payload}\n\n" @app.get("/chat/stream") async def chat_stream(question: str): headers = { "Cache-Control": "no-cache", "Connection": "keep-alive", "X-Accel-Buffering": "no", # 重要!关闭 Nginx 缓冲 } return StreamingResponse( generate_events(question), media_type="text/event-stream", headers=headers )这里有几个细节值得展开。
event:字段的作用,是给帧定义一个事件类型。浏览器端用EventSource时可以通过addEventListener("start", ...)分别监听不同类型的事件;用fetch流式解析时,这个字段更多是语义化标记,方便前端按类型分发逻辑。我习惯设计三种事件:start(流开始)、token(增量 token)、done(流结束,附带完整文本)。其中done事件携带完整文本这一点非常关键,后面讲 JSON 解析时你会看到它的价值。
X-Accel-Buffering: no这个头是血泪教训。如果你把服务部署在 Nginx 后面,Nginx 默认会缓冲上游响应,攒够一定量才发给客户端,这会导致你在后端明明已经逐 tokenyield了,前端却要等很久才一次性收到。加了这个头,等于告诉 Nginx"这个响应别缓冲,来一口吐一口"。没有这一行,你的 SSE 在代理后面会变成一顿一顿的"伪流式",这个坑我见过不止一次。
心跳帧的问题。前面提到的 idle timeout,本质上是代理层或负载均衡器认为连接空闲了,主动掐断。解决思路是在等待模型输出的间隙主动发心跳帧。因为请求频繁(比如模型推理耗时长)时,好几分钟不发数据,GG 的可能性极高。心跳帧通常是一个注释行:yield ": heartbeat\n\n",SSE 规范里以冒号开头的行是注释,浏览器端会自动忽略,不会触发事件。每 15 秒发一个,代理层就不会认定连接空闲。
2.3 加了 LangGraph 怎么办?MCP 工具输出也要流转起来
如果你的项目加了 LangGraph 做 Agent 编排,流式输出的复杂度会上一个台阶。因为 LangGraph 的状态图里可能有多个节点:先调用工具,再让模型基于工具结果生成回答,中间穿插 Agent 的思考过程。这时候直接astream拿到的不只是最终答案 token,还有节点的执行状态、工具调用参数、中间结果等结构化信息。
我的处理方式是区分"流式事件"和"状态事件"。工具调用、节点切换这类事件不需要"打字机"效果,直接整块推给前端,更新 UI 状态(比如显示"正在检索数据库...");模型的生成 token 才走流式通道。LangGraph 里可以通过stream_mode来控制输出的粒度,例如stream_mode=["updates", "messages"],updates给节点状态变更,messages给模型生成的逐 token 消息。拿到之后,在后端做一层分拣,再按不同事件类型封装成 SSE 帧推出去。
因为工具返回的数据往往是结构化的,这就自然引出了一个问题:前端不仅要展示文本,还要展示工具结果里的 JSON 字段(比如搜索结果列表、数据库行、图表数据)。这时候如果你把工具结果塞进自然语言文本里一起流式吐出,前端解析会非常痛苦。更合理的做法是:工具结果单独作为一个 SSE 事件,里面直接放 JSON,前端拿到后渲染成表格或卡片;模型对结果的解读再走 token 流。这种"结构化推送 + 文本流式"的组合,是我目前见过最符合 Agent 场景的流式架构。
3. 前端打字机效果:fetch 流式读取的完整实现
后端把 token 流转化成了 SSE 事件流,前端要做的就是从 HTTP 响应体里把data:行的内容逐段读出来,拼进 DOM。这里有个认知要纠正:EventSource虽然省事,但自定义请求头(比如Authorization)和 POST 方法它都不支持。所以在真实项目里,我几乎全用fetch+ReadableStream手动解析,虽然代码多一点,但换来了完整的控制权。
3.1 用 fetch 读取 SSE 流的核心代码
先给出一版可以直接抄的完整实现,再解释关键点。
async function streamChat(question, onToken, onDone, onError) { const response = await fetch("/chat/stream?question=" + encodeURIComponent(question), { headers: { "Accept": "text/event-stream", "Authorization": `Bearer ${localStorage.getItem("token")}`, }, }); if (!response.ok) { throw new Error(`HTTP ${response.status}`); } const reader = response.body.getReader(); const decoder = new TextDecoder("utf-8"); let buffer = ""; let accumulatedText = ""; while (true) { const { value, done } = await reader.read(); if (done) break; buffer += decoder.decode(value, { stream: true }); // 按空行切分 SSE 帧 const frames = buffer.split("\n\n"); buffer = frames.pop(); // 最后一段可能是半截帧,留着下轮继续拼 for (const frame of frames) { const dataLine = frame.split("\n").find((line) => line.startsWith("data:")); if (!dataLine) continue; const payload = JSON.parse(dataLine.slice(5).trim()); if (payload.type === "token") { accumulatedText += payload.content; onToken(payload.content); } else if (payload.type === "done") { onDone(payload.full_text ?? accumulatedText); } } } }这段代码里有三个关键设计。
流式解码器TextDecoder的stream: true参数是重点。如果你直接用decoder.decode(value),当多个 chunk 拼接出一个多字节字符的时候,会发生乱码。stream: true会告诉解码器"当前输入可能是不完整的",内部会缓存未完成的字节,等到下一个 chunk 来了再补全。我最初没加这个参数,结果中英文混排时,中文经常出现乱码。
按空行切帧的缓冲策略。SSE 帧之间以空行分隔,但网络传输不保证每个 chunk 恰好等于一帧的边界——很可能一个 chunk 里包含多帧,也可能只包含半帧。所以必须用一个buffer变量攒着,每轮循环先把 buffer 拼上新的数据,再按空行切分,最后一截留到下轮。这是所有流式读取的核心套路,不管你是读 SSE、WebSocket,还是读文件流,都要这么处理。
累积文本accumulatedText的作用。onToken只是把增量 token 追加到界面上,但如果你需要拿到完整结果(比如传给后端做下一步处理、生成标题等),就得自己累积。我在后端done事件里也会附带full_text,双保险,防止前端累积出错。为什么这么设计?因为前端累积过程中可能因为用户中断、网络异常而丢数据,后端在流结束前一次性推完整文本,就相当于给了个"校验锚点"。
3.2 打字机效果的渲染策略:为什么不要逐字改 DOM
拿到了 token,接下来的问题是怎么渲染。最直观的做法是每收到一个 token 就innerHTML += token,这在 token 少、频率低的时候没问题,但实测在快速生成场景下会卡顿,因为每次 DOM 更新都触发布局重算。正确姿势是用 CSS 动画 + 一次性挂载,或者用防抖批量渲染。
我的常用方案是维护一个内容数组,每 30ms 检查一次有没有攒下的 token,有就一次性拼进 DOM。因为流式 API 的推送频率通常是毫秒级,但浏览器的渲染帧率上限是 60fps(即 16.6ms 一帧),如果 token 推送比帧率还快,逐字更新就是纯浪费。实测下来,30ms 的批量节奏既不会察觉延迟,又能避免无谓的性能损耗。
let pendingTokens = []; let renderTimer = null; function onToken(token) { pendingTokens.push(token); if (!renderTimer) { renderTimer = setTimeout(() => { const contentEl = document.getElementById("content"); contentEl.textContent += pendingTokens.join(""); pendingTokens = []; renderTimer = null; }, 30); } }还有一个小技巧值得分享:用textContent而不是innerHTML。模型输出的内容里可能包含<script>、<img>这类标签,如果你直接当 HTML 塞进去,等于给 XSS 开了大门。textContent会把所有内容当纯文本处理,安全系数直接拉满。如果确实需要渲染 Markdown,也应该是先对原始文本做安全过滤,再走 Markdown 解析器。
4. 结构化输出:让模型输出能被程序直接消费的 JSON
打字机效果只是把文字"展示"出来了,但真正的 AI 应用,通常还需要让模型输出能被程序直接消费的结构化数据——最常见的形态就是 JSON。这一节我们从"为什么要结构化"讲到"LangChain 怎么保证结构化",最后落到"不依赖 LangChain 的纯手写做法"。
4.1 为什么需要结构化输出:从"聊天"到"干活"的分水岭
聊天机器人场景下,模型输出一段自然语言就够了。但一旦 AI 要"干活"——比如提取订单信息、生成数据库查询、调用工具传参、生成报表配置——就必须把输出从"给人看的自然语言"变成"给程序看的结构化数据"。一个很典型的例子:让模型从一段客户反馈中提取"用户情绪、问题类别、紧急程度"三个字段。如果没有结构化约束,模型可能给出:"用户比较生气,问题出在物流,很急。" 这句话人懂,程序没法直接处理。有了结构化约束,模型输出就会变成:
{ "sentiment": "negative", "category": "logistics", "urgency": 5 }程序拿到之后直接做条件分支,这就叫"让 AI 下地干活"——聊天只是 AI 的起点,结构化输出才是它进入工作流的关键隘口。
4.2 LangChainwith_structured_output:声明式结构化
LangChain 里,结构化输出的正规路径是with_structured_output。它的底层思路是把你定义的 Pydantic 模型转成 JSON Schema,然后以特定的方式约束模型输出。用法很简单:
from typing import Literal from pydantic import BaseModel, Field from langchain_openai import ChatOpenAI class FeedbackAnalysis(BaseModel): sentiment: Literal["positive", "neutral", "negative"] = Field(description="用户情绪") category: str = Field(description="问题类别") urgency: int = Field(ge=1, le=5, description="紧急程度,1最低,5最高") llm = ChatOpenAI(model="gpt-4o-mini", temperature=0.0) structured_llm = llm.with_structured_output(FeedbackAnalysis) result = structured_llm.invoke("我买的手机三天了还没发货,客服一直不回复,气死我了!") print(result.sentiment) # negative print(result.urgency) # 5注意几个关键点。
温度要调成 0 或接近 0。结构化输出的确定性要求极高,温度太高模型容易在边界字段上"自由发挥",urgency给你填个 7,sentiment标个 "furious"。temperature=0.0不等于完全确定,但能把随机性压到最低。
Field 的 description 是"隐形提示词"。很多人忽略这一点,以为 Pydantic 定义只是声明类型。实际上Field(description=...)里的内容会被合并到 JSON Schema 传给模型,是引导模型正确填值的强信号。我见过一个真实案例,同样两个字段,一个写了 description,一个没写,模型在没写的那个字段上的出错率高 20%。所以在定义结构时,每个字段都要写清楚"这个字段到底代表什么、填什么格式"。
枚举和约束要用好。Literal和Field(ge=, le=)都是强烈约束,能限死取值范围。模型在绝大多数情况下会严格遵守 Schema 约束。跟正则匹配模型输出相比,这种方式在复杂结构下准确性更高,因为它不是事后校验,而是引导生成。
4.3 流式 + 结构化怎么共存:先流式展示,再结构化解析
with_structured_output的常规用法是阻塞式的,但如果我要的是流式打字机效果,同时最后还能拿到结构化结果,该怎么办?我试过几种组合方案,最终稳定落地的是这个模式:
流式阶段:用普通的流式生成,让 token 一个个蹦到 UI 上,用户看到的是自然语言回答。
结构化并行阶段:同一个请求同时发起两个 LLM 调用(或者先用流式,结束后再调一次结构化)。第一次调用负责"展示",第二次调用负责"解析"。
听起来浪费了一次调用?其实很多场景下这就是最省心的选择。因为同一个模型的同一个提示词,你要它既流畅出自然语言,又严格输出 JSON,这两个目标本身就是冲突的——自然语言要求自由表达,JSON 要求严格约束,两个目标挤在同一次生成里,模型会两头不讨好。我在一个项目里试过单次生成 + 输出 JSON 的格式说明,结果结构上的错误率明显偏高。
所以我的稳定方案是:展示和解析解耦。前端先看到流式文字,用户不再等待时,后端再(或同时)发起一个小的结构化调用,生成结果直接进数据库、触发下游任务。代价是多次模型调用,换来的是稳定性和可维护性。对于要求实时响应的场景,可以用异步处理,结构化解析放在后台,不阻塞 UI。
4.4 不依赖 LangChain 的纯手写结构化:System Prompt + JSON Schema 双管齐下
如果你的项目没接 LangChain 也不想接(比如用的是原生 OpenAI SDK 或其他模型的 SDK),结构化输出也是可以手写的,核心套路是两件事。第一,在 System Prompt 里明确写"只输出 JSON,不要任何解释、不要 Markdown 代码块标记、不要多余内容"。第二,把 JSON Schema 直接贴在 prompt 里,让它照着填。然后你可以用response_format={"type": "json_object"}之类的参数做硬约束(OpenAI 系支持,部分国产模型也有类似能力)。
我自己手写的时候,还加了一个"纠错"步骤。因为模型偶尔会输出 JSON 数组而不是对象,或者字段名会错。比如我让模型输出{"key": "value"}这种格式的时候,它给我输出了一段 Markdown 包裹的 JSON 代码块。所以手写方案里必须接一层 JSON 解析的容错逻辑——先尝试json.loads,失败就剥掉 Markdown 代码块标记再试,再不行就正则提取大括号内的内容。这套解析和容错逻辑,下一节细讲。
5. JSON 解析全方案:从半成品到错乱数据都能救回来
程序拿到模型的"JSON 输出"之后,真正的噩梦才开始。模型不是编译器,它不会保证输出的 JSON 永远合法。我统计过在一个真实项目里的情况,没有经过任何后处理的模型 JSON 输出,首轮解析失败率在 15% 到 30% 之间。这个数字非常可怕,如果你直接把模型输出丢给json.loads,你的系统会频繁在解析环节爆炸。
5.1 错误类型盘点:模型不是 JSON 编译器
常见的模型 JSON 输出错误,我按频率排个序:
- Markdown 代码块包裹:模型会输出
json ...,整个包起来。这是最常见的错误,可能是因为在训练数据里,JSON 基本都以代码块形式出现。 - 前后有多余文本:比如"好的,这是您要的结果:" 后面才是 JSON。
- 单引号代替双引号、末尾多逗号,或者字段间缺逗号。
- 字段名被模型"翻译"了:你让它输出
category,它写成"类别"。这比语法错误更隐蔽。 - 整个输出被截断,JSON 只生成了一半,结尾不完整。
- JSON 是合法的,但不是你要的 schema:多了一个嵌套层级、或者把数组和对象搞混了。
知道了错误类型,解析函数就可以分层兜底。
5.2 多层兜底解析函数
我沉淀了一套"四级解析"方案,建议直接抄进你的工具库里。
import re import json def extract_json(text: str): # 第一级:直接用 json.loads 解析 try: return json.loads(text) except Exception: pass # 第二级:去掉 Markdown 代码块标记 text_clean = re.sub(r"```(?:json)?", "", text).replace("```", "").strip() try: return json.loads(text_clean) except Exception: pass # 第三级:用正则提取最外层大括号 match = re.search(r"\{.*\}", text_clean, re.DOTALL) if match: try: return json.loads(match.group(0)) except Exception: pass # 第四级:修正常见语法问题(单引号、末尾逗号) fixed = text_clean.strip() fixed = re.sub(r"'", '"', fixed) fixed = re.sub(r",\s*([}\]])", r"\1", fixed) try: return json.loads(fixed) except Exception as e: return {"error": f"JSON parse failed: {e}", "raw": text_clean}注意第四级是"双刃剑"。re.sub(r"'", '"', fixed)会把所有单引号改成双引号,但如果模型输出了"it's"这种合法英文缩写,这步就会破坏内容。所以我在这个工具函数里加了保护,只对 JSON 结构部分做替换,或者干脆在第四级返回错误标记,让上层决定是重试还是提示用户。这个取舍要看你的业务容错度:数据进下游流程,宁可解析失败也别解析错。
5.3 不完整 JSON 再造:流式场景的"半成品补救"
流式场景下有一个特殊问题:用户在打字机效果还没结束的时候,后端的"完整文本"是没有的,此时如果你想在半途就把已生成的 JSON 拿来预览(比如生成一个富文本卡片边写边渲染),你拿到的就是半截 JSON。解析半截 JSON 用json.loads铁定失败,但你可以用一个叫"增量 JSON 解析"的思路。
增量解析的核心是容忍"暂时不完整"。比如模型正在生成{"name": "张三", "age":,此时age的值还没出来。你要做的不是报错,而是把这个半成品状态暴露给前端,比如解析器返回一个{"__partial__": true, "partial": {...}},前端知道"数据还不完整,继续等"。实现可以基于ijson这类库的流式解析,或者自己实现一个简单的状态机,跟踪当前处于哪个字段、哪个值。不过说实话,大部分场景用不着做到这个程度——大多数业务只需要在流结束后拿到完整 JSON,中途的"半截预览"是锦上添花的功能,复杂度不低,优先级要排后。
5.4 用 JSON Schema 校验替代 json.loads:最后一道安全闸
光能解析出 JSON 还不够,你还得校验它是否符合预期结构。最省事的做法是拿 Pydantic 模型直接验证。LangChain 的with_structured_output本身内部就带了这一步,如果你手写方案,也建议把数据灌进 Pydantic 里做一轮强校验。
from pydantic import ValidationError try: data = FeedbackAnalysis(**extracted_data) except ValidationError as e: # 记录原始输出、错误原因、模型名称,便于复盘 log_failure(extracted_data, e.errors()) raiseValidationError会告诉你哪个字段缺失、哪个字段类型不对,方便你精准定位是模型的错还是解析层的错。我在生产环境里会把全部分析失败的数据落库,每周看一眼分布,如果某类错误的比例持续偏高,就去调整 prompt 或字段描述。
6. 我在生产环境踩过的流式大坑:超时、断连与并发
前面讲的都是怎么"把功能做出来",这一节聊的才是"让它在生产环境活下来"。流式 SSE 在开发环境跑得再欢,也架不住线上出幺蛾子。我把这些坑按高频程度排个序,你以后的排查大概率用得上。
6.1 idle timeout:连接被代理层判定"空闲"而掐断
这是我在真实环境踩的第一个大坑,报错信息就是热搜词里那句"stream disconnected before completion: idle timeout waiting for sse"。发生场景是:模型思考很久才吐出第一个 token(比如调用了外部工具、或者用的是较慢的大模型),这期间连接上没有任何字节流动,代理层(Nginx 的proxy_read_timeout、负载均衡器、云厂商网关)就判定连接空闲,直接掐断。
解决方案有三个层次。第一,加心跳帧,也就是前面说的": heartbeat\n\n",15 到 30 秒一发,保证连接上有字节在流动,这是最根本的解法。第二,调大代理层的超时时间,比如 Nginx 里proxy_read_timeout 300s;。第三,前端自动重连,SSE 协议本身就支持断线重连,EventSource会自动处理,但用fetch流式解析就要自己实现重连逻辑,建议带上指数退避,别一断就连,把服务端打崩。
6.2 Nginx 缓冲导致"伪流式"
这个问题前面提到过:服务端明明在逐 tokenyield,但用户在前端等了几秒才一次性看到整段文字,打字机效果直接变成"打字机被砸了"的效果。原因就是 Nginx 在中间做了缓冲,攒够 4KB 或一定时间才统一转发。
解法就是响应头里加X-Accel-Buffering: no。注意这是 Nginx 特有的头,其他代理(比如云厂商的 SLB)可能有不同的关闭方式,排查时要确认你用的是哪一层代理。
6.3 多个并发请求串流:每一路都得独立 buffer
很多人以为 SSE 的解析是"每个请求一个响应",天然不会串线,但实际出过问题:用户在页面上一口气发起多个流式请求(比如对比两个模型的回答),前端的 buffer 状态如果设计成全局变量,两个流的 token 就会混在一起。教训是:每个流式请求必须有自己独立的 buffer、累积文本和渲染目标。我后来都封装成一个StreamSession类,一个实例对应一路流,内部状态互不干扰。
class StreamSession { constructor(onToken, onDone) { this.buffer = ""; this.accumulated = ""; this.onToken = onToken; this.onDone = onDone; } // 把读取循环里的逻辑都收进这个类 }6.4 客户端断开,服务端还在跑
用户在打字机效果进行到一半时关闭了页面,正常情况下浏览器会中断 TCP 连接,reader.read()会抛错,然后我们跳出循环。但如果用户是断网(而不是主动关闭),服务端可能一时半会儿感知不到连接断开,还在继续生成、继续yield。这会导致模型调用费的浪费、日志里出现一堆无意义的生成记录,严重的还会把服务端进程拖垮。
我现在的做法是给每个流式生成任务加一个"取消"机制。FastAPI 的StreamingResponse有个特点:如果客户端断开,生成器会在下一次yield时抛出GeneratorExit。利用这一点,在生成器里做个 try/except/finally,finally里清理状态、记录日志、终止上游 LLM 调用。LangChain 的astream如果不用了,记得要取消任务,避免 LLM 还在后台继续生成。
def generate_events(question: str): task = None try: task = asyncio.create_task(collect_tokens(question)) async for token in task: yield f"data: {json.dumps({'type': 'token', 'content': token})}\n\n" except asyncio.CancelledError: if task: task.cancel() raise finally: # 清理资源,记录中断日志 ...6.5 从 SSE 到 MCP 工具流式输出的拓展思考
最近在折腾 MCP(Model Context Protocol),有个体会是:工具调用的流式输出和 SSE 的流式输出,本质上是同一套哲学。MCP 里的工具返回数据,往往也是结构化 JSON,也需要经过"半截状态识别、最终校验、错误兜底"这条路。如果你把 SSE 这条链路彻底吃透了,再去看 MCP 的流式协议,会发现只是帧格式和事件类型变了,底层的"数据分块传输、缓冲切分、增量解析"完全是一回事。这也是为什么我建议在这类项目上多花时间把流式基建做好——它是一劳永逸的底层能力,不是一次性需求。
根据我这几个项目的实战经验,流式 + 结构化输出这套组合,最难的部分从来不是某个单一环节,而是整条链路的韧性:后端要扛住长时间连接和并发,前端要处理半截帧和渲染性能,模型输出要经过多层解析和校验,代理层要在中间"不捣乱"。你把这一整套都跑通了,再回头看"打字机效果"四个字,它不过是这条链路上最表层的一个 UI 效果而已。真正让人夜里睡得着的,是链路里每一个环节都有兜底方案。