☰
SSE与LangChain流式输出实战:AI对话打字机效果全链路指南
2026/10/6 15:25:26 网站建设 项目流程

做 AI 应用做了快两年,我最大的感受是:能不能把一个流式对话做顺,决定了用户打开你产品的第一印象。你辛辛苦苦调好了 Prompt,结果前端一直在转圈,用户等了三秒没有反应,直接关页面了——这类问题十有八九不是模型不行,而是你压根没把数据流的链路打通。最近在做项目的过程中,我专门把 SSE 流式传输、LangChain 结构化输出、前端打字机效果这一整条链路重新梳理了一遍,踩了不少坑,也沉淀了一套可以直接抄走的方案。这篇文章就把这套东西完整拆开,从协议原理到前后端代码,再到那些"文档里绝对不会告诉你"的排查经验,一次讲透。

这篇实战内容适合谁?后端同学想给 AI 接口加流式输出但还没想清楚消息格式,前端同学被「怎么把流式数据渲染成打字机效果」卡住,或者你已经在用 LangChain 但发现chain.stream()出来的东西和想象的不一样——那这篇文章就是给你准备的。我会尽量说人话,复杂的地方用类比讲明白,保证你看完能直接动手改代码。

1. AI 对话为什么绕不开 SSE?流式协议的本质与选型

1.1 从轮询到长连接:SSE 到底解决了什么问题

传统 Web 应用里,前端想知道后端有没有新数据,最笨的办法是轮询——每隔一两秒发一次请求问"好了吗?好了吗?"。这个模式在老系统里很常见,但放到大模型对话场景里就是灾难。LLM 生成一段 300 字的回答,哪怕速度已经很快,也要好几秒才能吐完。如果前端傻等全部生成完再一次性展示,用户在这几秒里面对的就是一片空白,既不知道请求到底有没有成功,也不知道模型是不是卡住了。

SSE(Server-Sent Events,服务器发送事件)就是专门解决这个问题的。它让服务器把响应内容拆成一个个小片段,通过一条 HTTP 长连接持续推送。前端建好连接之后,不需要反复发请求,只要坐在那里接收数据就行。用大白话说,轮询是你每隔一会儿跑去问店家"饭好了没",SSE 是店家做好一道菜就端一道菜上来,你坐那儿等着吃。

SSE 的技术基础是 HTTP 协议本身,不需要额外安装 WebSocket 之类的依赖,而且它的消息格式极其简单:每条消息用data:开头,结尾是一个空行。比如:

data: 你好 data: 世界

这样一个文本流,前端收到后按空行切分,就能拿到两条消息。SSE 还支持event:字段自定义事件类型、id:字段做断点续传,机制虽然简单,但足够支撑大部分实时推送场景。

1.2 SSE 协议格式与 LangChain 流式输出的天然契合

LangChain 的流式接口astream()返回的是一个异步迭代器,模型每生成一个 Token(或者一小段 Token),这个迭代器就往外吐一次数据。这种"一边生成一边吐"的模式,恰好和 SSE 的"一端推一端收"是天生一对。

打个比方:LLM 是个很能说的人,SSE 是一条电话线。LangChain 负责让这个人不停嘴地往外说话,SSE 负责把每一句话实时传到电话那头,前端的"打字机效果"就是听众看着记录员一个字一个字地记下来。

实际编码里,你只需要在 FastAPI 的StreamingResponse里放一个异步生成器,把 LangChain 吐出来的每一个 chunk 按 SSE 格式包装,网络层就会把这些消息逐条推给前端。这里有一个关键点:SSE 的心跳机制。协议规定,如果服务器持续 15 秒以上没有发送任何数据,连接就可能被中间的网络设备关闭。所以规范的 SSE 服务器即使没有业务数据,也应该定期发送一个注释行(:开头的一行)作为心跳,相当于对着电话线"喂喂喂"地确认连接还活着。

1.3 SSE 与 WebSocket 的选型对比

很多同学第一反应是"流式传输那必然 WebSocket"。实际做下来你会发现,在 AI 对话这个场景里,SSE 的优势比 WebSocket 更明显。我整理了一张对比表,方便你根据场景选型:

维度SSEWebSocket
传输方向服务器单向推送给客户端双向实时通信
底层协议HTTP,天然穿过网关和代理TCP,需要协议升级
消息格式纯文本,data:前缀,可读性极强二进制或文本,需要自定义协议
自动重连协议内置,连接断开会自动重试需要自己实现重连逻辑
携带参数依赖 URL/Query,头部受限可以自定义请求头
浏览器支持现代浏览器全支持全支持,但需要额外的事件处理

单看最后两行你可能觉得 WebSocket 更强,但 AI 对话这个场景里,数据传输方向就是单向的——服务器把模型生成的 Token 推给前端,前端基本不需要反向发送数据(请求参数已经在 POST 里带过去了)。用 WebSocket 等于杀鸡用牛刀,还要额外操心二进制分帧、粘包、重连这些破事。SSE 自带断线重传机制,浏览器 EventSource 对象甚至能配置重试间隔,这对我来说已经足够香了。唯一的硬伤是 EventSource 只能发 GET 请求、不能自定义 Header——所以我后面实战部分直接绕开了它,用 fetch 流式读取来拿 SSE,这样既保留 POST 传参灵活性,又能按需处理各个事件。

2. LangChain 流式接口逐层拆解:stream、astream 与消息结构

2.1 Runnable 协议的四种调用方式

LangChain 从 0.1.x 开始主推 Runnable 协议,所有 Chain 组件都实现了统一的调用接口。和咱这篇文章相关的,就是下面四种:

  • invoke():同步调用,输入一个字典,输出最终结果。适合调试和脚本场景。
  • ainvoke():异步调用,还用 await 接。适合在 FastAPI 请求处理函数里用。
  • stream():同步流式,返回一个迭代器,每次吐一块中间结果。
  • astream():异步流式,返回异步迭代器,配合async for使用。后台接口推荐用这个。

关键是理解:invoke和stream的区别不是"快和慢",而是是否暴露中间过程。invoke把整条链路当成一个黑盒,你只能拿到最终输出;stream把链路里的每一步都摊开给你看,模型每生成一个 Token,迭代器就吐一次。

from langchain_openai import ChatOpenAI from langchain_core.prompts import ChatPromptTemplate from langchain_core.output_parsers import StrOutputParser llm = ChatOpenAI( model="gpt-4o-mini", temperature=0.7, streaming=True, # 强烈建议在生产环境按实际模型服务商配置 base_url、api_key ) prompt = ChatPromptTemplate.from_messages([ ("system", "你是一个简洁的中文写作助手,回答控制在300字以内。"), ("human", "{question}"), ]) chain = prompt | llm | StrOutputParser() # 同步逐 Token 输出 for chunk in chain.stream({"question": "用一句话解释什么是SSE"}): print(chunk, end="", flush=True)

跑一下这段代码,你会看到控制台像打字机一样逐字输出。那个streaming=True很关键——它告诉 OpenAI SDK 开启流式模式,服务端用 chunked 编码逐个发 Token,LangChain 才能逐块往外吐。如果漏掉这个参数,chain.stream()依然会"一次性把结果给你",所谓流式就名存实亡了。

2.2 流式链里真正会吐出来的数据类型

刚开始做流式接入的同学,往往会犯一个错误:以为astream()每次吐出来的都是纯字符串。实际上不一定。这取决于你的 Chain 最后挂的是什么组件。

如果 Chain 结构是prompt | llm | StrOutputParser(),那么最后一步StrOutputParser会把 AIMessageChunk 的内容转成字符串,所以astream()吐出来的是字符串,直接拼起来就是最终文本。

但如果你的 Chain 结构是prompt | llm(没有 StrOutputParser),那么astream()吐出来的是AIMessageChunk对象——它里面有content属性,还可能有tool_calls(如果模型调用了工具)。这时候你要么自己加一个(msg) -> msg.content的魔法方法把那层壳剥掉,要么在业务代码里用hasattr(chunk, "content")判断一下,再决定怎么取文本。

顺着这个思路往下走,你还会发现,结构化的输出解析器(比如PydanticOutputParser)和流式组合时有个特性:解析器本身不是流式的。它必须等模型把整段 JSON 吐完,才能一次性解析出结构化对象。这意味着整个 Chain 的最后一段数据是在"憋大招",它会先攒着,然后突然蹦出一个大对象。我用一个简单 Chain 验证过:prompt | llm | PydanticOutputParser走astream(),前面模型吐了十几个 token 数据,解析器那边一个都不会冒出来,最后一次性输出最终结果。做前端的同学如果发现"流式到一半就停了,然后又一下子出现完整内容",问题十有八九出在这里。

2.3 非 Token 数据的处理:元数据与工具调用片段

另外一个容易被忽略的环节是:现代 LLM 的流式响应里不一定只有文本 Token。如果你开启了函数调用(Function Calling / Tool Calling),模型可能先吐出一段tool_calls的增量参数,再吐出文本;如果你在 LangChain 里挂了回调处理器(callbacks),流式过程中还会穿插各种事件的回调。

我在实际项目里处理过一种脏数据:模型同时返回content和tool_calls的 chunk,前端只关心content,但如果不加过滤,这些tool_calls的字符串会被直接拼进正文里,用户看到的就是满屏乱码结构。

解决方案是做一个统一的流式消息格式化层,把不同类型的输出拆成不同的 SSE 事件:

  • token事件:纯文本增量,直接拼到正文;
  • tool事件:工具调用信息,前端可以展示"正在检索资料"之类的状态;
  • meta事件:最终的结构化数据,比如 JSON 提取结果、意图分类、引用列表。

这么设计的好处是,前端只负责展示,后端告诉它"这块是正文""这块是状态""这块是 JSON",两边职责清楚,调试也方便。后面第 4 节我会给出完整代码。

3. FastAPI + LangChain 流水线实战:后端 SSE 封装与心跳机制

3.1 用 StreamingResponse 正确构建 SSE 响应

FastAPI 可以说是和 LangChain 流式配合得最顺的 Python 框架。它的StreamingResponse接收一个可迭代对象,然后通过 HTTP 分块传输把数据一块一块发出去。我们要做的,就是把 LangChain 的异步生成器包装成 SSE 格式。

核心代码如下,我尽量把注释写清楚:

import asyncio import json import time from fastapi import FastAPI from fastapi.responses import StreamingResponse from pydantic import BaseModel from langchain_openai import ChatOpenAI from langchain_core.prompts import ChatPromptTemplate from langchain_core.output_parsers import StrOutputParser app = FastAPI() llm = ChatOpenAI( model="gpt-4o-mini", temperature=0.7, streaming=True, ) prompt = ChatPromptTemplate.from_messages([ ("system", "你是专业的中文助手。回答要结构清晰、语言自然。"), ("human", "{question}"), ]) # 这个 chain 只负责流式文本 chat_chain = prompt | llm | StrOutputParser() def sse_format(event: str, payload: dict) -> str: """把事件和数据包装成 SSE 协议的一帧""" return f"event: {event}\ndata: {json.dumps(payload, ensure_ascii=False)}\n\n" async def sse_generator(question: str): """生成器:把 LangChain 的流式输出转成 SSE 帧""" # 1. 先发一个多行注释作为心跳,顺便告诉前端“我开始了” yield ": connected\n\n" yield sse_format("meta", {"type": "start", "timestamp": time.time()}) try: # 2. 核心循环:逐 token 包装 async for chunk in chat_chain.astream({"question": question}): # 兼容 AIMessageChunk 的情况 if isinstance(chunk, str): content = chunk elif hasattr(chunk, "content"): content = chunk.content or "" else: content = str(chunk) # 有真实内容才发 token 事件,空白跳过 if content.strip(): yield sse_format("token", {"content": content}) # 3. 结束标记 yield sse_format("end", {"type": "done"}) except Exception as exc: # 生成器内捕获异常,否则客户端只会看到连接中断,什么都拿不到 yield sse_format("error", {"message": str(exc)}) class ChatRequest(BaseModel): question: str @app.post("/chat") async def chat_endpoint(req: ChatRequest): return StreamingResponse( sse_generator(req.question), media_type="text/event-stream", headers={ "Cache-Control": "no-cache", "Connection": "keep-alive", "X-Accel-Buffering": "no", # 告诉 Nginx 别缓冲 }, )

有几个细节我要特别强调,这些是我实际运行时踩出来的坑:

第一,media_type必须是text/event-stream。这是浏览器识别 SSE 响应的硬条件。如果你漏掉了,前端即使收到了数据,EventSource 也不会正确解析。用 fetch 流式读取时,这个字段依然重要,因为有些代理服务器会根据 Content-Type 决定要不要缓冲。

第二,X-Accel-Buffering: no这个响应头不能省。当你把服务部署到 Nginx 后面时,Nginx 默认会缓存 FastAPI 发出来的所有数据,攒够一定量再一次性发给客户端。那打字机效果就废了——用户看到的还是"转圈半天,猛一下全出来"。加了这个头,就是告诉 Nginx"这条响应你不要碰,来一个帧发一个帧"。

第三,生成器里必须 try-except。模型生成过程中抛异常是常有的事(上下文超限、内容过滤、网络抖动)。不 catch 的话,FastAPI 会把连接直接断开,前端拿到的是半截消息,什么错误信息都没有。把异常包装成 SSE 的error事件发给前端,至少用户能知道"出错了",不至于以为是网络挂了。

3.2 同步生成器和异步生成器:前端卡顿的隐患

如果你的业务代码里拿的是同步的chain.stream(),在 FastAPI 里有两种用法:

  • 把同步生成器传给StreamingResponse,FastAPI 会在线程池里跑它,不会阻塞主事件循环;
  • 但要注意,同步生成器里如果调用了阻塞的 I/O 或 CPU 密集操作,还是会占住线程池资源。

我踩过的一个坑是这样的:一开始我图省事,在async def端点里直接用chain.stream()(同步方法)循环,结果事件循环被阻塞,前端同时发起的第二个请求要等第一个跑完才开始响应。这个问题特别隐蔽,因为单个请求测试时根本感觉不出来,一旦并发量上来,接口响应曲线直接爆炸。

正确姿势是用astream()配合async for,或者在生成器函数前面加一个线程池装饰器/asyncio.to_thread包装。而且我建议把"生成器本身"也设计成异步的,让 LangChain 的异步流式能力真正发挥出来。对于大多数生产环境来说,这就是你要的代码。

3.3 心跳机制与空闲超时问题

先给你复现一个真实场景。我用 LangChain 接了个长思考链路的模型(内部套了好几个子 Agent),有一次用户问了一个复杂计算题,模型在内部工具推理阶段沉默了 40 秒没吐任何 Token。连接中断,前端报了一个让我们排查了很久的错误——消息流在完成前就断开了,像极了超时等待 SSE 数据。

这个问题的根因是空闲超时(idle timeout)。网络设备、Nginx、Gunicorn 都以"这条连接最近有没有数据流动"作为判断依据。模型 40 秒没吐数据,Nginx 觉得这连接废了,直接掐断。

解决思路有两个层次:

  • 治本:在生成器里加心跳。每隔一段时间(我一般设 15 秒)主动发送一个 SSE 注释帧: heartbeat\n\n,让连接一直有数据流动,网络设备就不会误判。
  • 治标:调大 Nginx 和 Gunicorn 的超时时间,比如proxy_read_timeout 300s、Gunicorn--timeout 300。但这只是给服务器"更多的耐心",如果模型真的长时间无响应,该断还是会断。

心跳帧用:开头是 SSE 规范规定的注释格式,接收方看到会无视它,不会当成业务数据。所以它可以安全地混在真实数据流中间。

async def sse_generator_with_heartbeat(question: str): last_heartbeat = time.time() async for chunk in chat_chain.astream({"question": question}): content = chunk if isinstance(chunk, str) else getattr(chunk, "content", "") if content.strip(): yield sse_format("token", {"content": content}) last_heartbeat = time.time() elif time.time() - last_heartbeat > 15: yield ": heartbeat\n\n" last_heartbeat = time.time()

这样即使模型长时间思考,连接也不会因为空闲而中断。

4. 结构化输出与 JSON 解析:让模型不只“说话”还能“填表单”

4.1 LangChain 结构化输出的三种主流方式

流式文本解决的是"用户看得到响应"的问题,但 AI 应用落地时还面临另一个问题:模型输出的是自然语言,系统怎么才能拿到结构化的数据去做后续处理?比如从简历里提取姓名、工作年限、技能列表,或者让 Agent 返回一个可执行的 JSON 指令。

LangChain 提供三种主流方式,我一一说清楚它们各自的适用场景:

方式一:Pydantic 输出解析器(PydanticOutputParser)这是最经典的方式。你先定义一个 Pydantic 模型(就是数据结构),LangChain 会把它的 schema 描述自动塞进 Prompt 的format_instructions部分,要求模型"按这个格式输出 JSON"。模型输出后,解析器把 JSON 转成 Pydantic 对象。整个过程不依赖特定模型的函数调用能力,适用面广。

方式二:with_structured_output()方法这是 ChatOpenAI 等模型封装类上的推荐方法。底层原理是让模型走 Function Calling / 原生 JSON 模式,把 Pydantic schema 转成工具的 JSON Schema,模型直接以"工具参数"的形式返回结构化数据。这种方式在 OpenAI 系模型上准确率极高,因为模型天生被训练过"调用函数时严格按 schema 输出"。

方式三:自定义response_format参数如果你用的是 OpenAI 兼容接口,可以直接在ChatOpenAI里传model_kwargs={"response_format": {"type": "json_object"}},强制模型输出 JSON。这个方式的约束力不如前两种,适合场景极简单的需求。

我用一张表对比它们的差异:

方式约束力依赖模型能力适用场景
PydanticOutputParser中等,靠 Prompt 引导所有模型通用场景,兼容各类模型
with_structured_output强,走工具调用协议OpenAI 及支持工具的模型需要高准确率解析
response_format弱,只强制 JSON 格式OpenAI 兼容接口简单、无嵌套结构的输出

4.2 用 Pydantic 定义输出 Schema 并接入链路

我实际项目里最常用的组合是"方式二 + Pydantic"。下面是一个从面试评价文本中提取结构化信息的例子:

from langchain_core.pydantic_v1 import BaseModel, Field class InterviewEvaluation(BaseModel): """面试评价结构化提取""" candidate_name: str = Field(description="候选人姓名") overall_score: int = Field(description="综合评分,0到100分") strengths: list[str] = Field(description="候选人的优点,列出2到4条") risks: list[str] = Field(description="候选人的风险点,列出1到3条") hiring_verdict: str = Field(description="录用建议,只能是hire/no_hire/on_hold之一") # 直接使用 with_structured_output,把 Pydantic 模型转成 OpenAI 工具协议 structured_llm = llm.with_structured_output(InterviewEvaluation) result = structured_llm.invoke( "张三的编程能力很强,算法题全过,但沟通表达存在明显欠缺" ) print(result.candidate_name) # 张三 print(result.overall_score) # 某个整数 print(result.hiring_verdict) # 某个枚举值

这里要解释两个细节。第一,Pydantic 模型里的Field(description=...)不是写给人看的注释,LangChain 会把 description 拼进给模型的工具 schema,模型靠它理解字段含义,所以描述写得越具体,输出正确率越高。第二,枚举值(hire/no_hire/on_hold)直接用字段值和字符串描述,比单纯让模型"自由发挥"可靠得多。

4.3 流式文本和结构化输出如何共存

项目做到一半你会发现一个很现实的需求:既要让用户看到逐字生成的回复,又要让后端拿到一个结构化的 JSON 做业务逻辑。这两个需求能不能同时满足?能,但需要设计。

我的做法是:在一条 SSE 流里,同时发两种事件。token事件走流式文本,让前端做打字机效果;meta事件在文本流结束后发送结构化结果。示例代码如下:

@router.post("/chat_interview") async def chat_interview(req: InterviewRequest): async def gen(): # 先流式输出评价文本 async for token in interview_chain.astream({"question": req.question}): yield sse_format("token", {"content": token}) # 再调用结构化模型,提取评价要点 eval_result = await structured_llm.ainvoke(req.question) yield sse_format("meta", { "type": "structured", "data": eval_result.dict(), }) yield sse_format("end", {"type": "done"}) return StreamingResponse(gen(), media_type="text/event-stream")

前端拿到meta事件后,可以把它渲染成侧边栏的"智能摘要"卡片,或者直接塞进表单。这样用户看到的是自然流畅的对话,系统拿到的是可入库的结构化数据,一举两得。

不过要提醒你,这种"流式文本 + 尾部结构化"的模式有一个小问题:结构化结果要等模型完整生成才能出,所以如果是长文本,前端可能要等好几秒才看到侧边栏更新。如果对实时性要求更高,你可以考虑让结构化模型和流式模型并行跑(两个 LLM 调用同时发出去),结构化结果先到就先发,文本流自然也不会受影响。代价是成本翻倍。大多数场景下"先文本后结构化"已经够用,没必要提前优化。

4.4 JSON 解析的兜底容错方案

理论上with_structured_output已经保证了模型按 schema 返回 JSON,但现实世界永远有意外。比如你接入了一个第三方模型,它照样走 OpenAI 兼容接口,但偶尔会把 JSON 包在 markdown 代码块里,或者输出到一半被截断。这种时候如果你直接json.loads(),就是一个JSONDecodeError等着你。

我的兜底解析流程一共四层,按顺序尝试,大量减少了线上报错:

import json import re def robust_json_parse(text: str): """四层兜底解析模型输出的 JSON""" # 第一层:去掉 markdown 代码块包裹 if text.startswith("```"): text = re.sub(r"^```(?:json)?\s*|\s*```$", "", text, flags=re.MULTILINE).strip() # 第二层:标准解析 try: return json.loads(text) except json.JSONDecodeError: pass # 第三层:截取最外层大括号,忽略前后杂音 match = re.search(r"\{.*\}", text, re.DOTALL) if match: try: return json.loads(match.group(0)) except json.JSONDecodeError: pass # 第四层:用 json_repair 做宽松修复(处理缺失引号、末尾逗号等) try: from json_repair import repair_json repaired = repair_json(text) return json.loads(repaired) except Exception: return None

这里有个容易踩的坑:第三层的正则\{.*\}用的是贪婪模式,如果 JSON 里有嵌套的对象,它会一直匹配到最后一个},正好是完整 JSON 的外层,所以能正确提取。但如果模型输出里同时有两个 JSON 对象,这个正则也会把它们拼在一起导致解析失败。所以这层只是兜底,不能作为主路径。

还有一个思路很多人会忽略:利用错误反馈让模型自我修正。当解析失败时,把异常信息和提取到的残缺 JSON 一起塞回 Prompt,让模型在不重新生成全部内容的前提下,只补全出错的部分。我在带错误重试机制的场景里测试过,修复成功率能到 80% 以上。

5. 前端打字机效果实现:Vue + fetch 流式读取全流程

5.1 为什么最终选了 fetch 而不是 EventSource

浏览器自带 SSE 客户端就是EventSource,它开箱即用、自动重连,很多教程都推荐它。但我在项目里实际对接 LangChain 的后端接口时,踩了三个让它不那么好用的坑:

  • EventSource 只能发 GET 请求,可我需要在请求体里带对话上下文、用户配置、历史消息等一堆 JSON,硬塞到 URL 里既丑陋又容易超长。
  • EventSource 不能自定义请求头。我想在 Authorization 头里带用户的登录凭证,EventSource 就是加不上去,只能退而求其次用 Cookie 或 URL 参数,既不安全又别扭。
  • EventSource 对事件处理的粒度不够。它是面向"接收推送"设计的,对于token、meta、end这种需要区分处理的自定义事件,你得写一堆 event listener 去注册,不如手写 fetch 解析来得直观。

所以我最终采用的方案是:用 fetch 的 POST 请求去打 SSE 接口,然后通过response.body.getReader()读取流式数据,自己解析 SSE 帧。看起来多写了几行代码,但换来的是完全可控的请求参数和响应处理。

5.2 ReadableStream 解析 SSE 帧:核心代码逐段说明

前端解析的关键在于,SSE 数据不是一次到位的,而是一块一块地来。网络字节流可能把一条消息拆成两半,也可能把多条消息粘在一起。所以解析器必须维护一个buffer,每收到一块数据就追加进 buffer,然后按\n\n(空行)切分完整帧,剩余的半截留到下一轮继续拼。

一段可用的 Vue 3 组合式函数代码如下:

export function useSSE() { const fullText = ref('') const isStreaming = ref(false) const structuredData = ref(null) async function sendMessage(question) { fullText.value = '' structuredData.value = null isStreaming.value = true const response = await fetch('/api/chat', { method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify({ question }) }) if (!response.ok) { throw new Error(`请求失败: ${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 // 把 Uint8Array 解码成字符串,注意 stream: true 处理多字节字符 buffer += decoder.decode(value, { stream: true }) // 按空行切分 SSE 帧,最后一个可能是半截帧,先放回 buffer const frames = buffer.split('\n\n') buffer = frames.pop() for (const frame of frames) { const parsed = parseSSEFrame(frame) if (parsed.event === 'token') { fullText.value += parsed.data.content } else if (parsed.event === 'meta') { structuredData.value = parsed.data } else if (parsed.event === 'error') { console.error('服务端错误:', parsed.data.message) } } } isStreaming.value = false } return { fullText, structuredData, isStreaming, sendMessage } } function parseSSEFrame(frame) { const lines = frame.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:')) { // SSE 允许多行 data,拼一个换行再 trim 保持 JSON 可解析 data += line.slice(5) + '\n' } } // 去掉末尾多余的换行位再解析 JSON const trimmed = data.replace(/\n$/, '') return { event, data: trimmed ? JSON.parse(trimmed) : {} } }

这段代码里我特别想提醒一个细节:TextDecoder的第二个参数{ stream: true }千万别漏。UTF-8 的中文字符在二进制流里可能横跨两个 chunk,比如第一个 chunk 来了 2 个字节,第二个 chunk 才来剩余的 1 个字节。如果你在每一轮都单独decode(value),中文就会在中间断开,变成乱码。{ stream: true }会告诉解码器"数据没完,我帮你缓存到完整字符再输出",这是中文流式场景下的必备参数。

5.3 打字机效果与 Markdown 实时渲染的协调

拿到流式文本后,最朴素的打字机效果就是把fullText直接绑定到某个 DOM 节点。但如果你在文本里混了 Markdown 语法,直接v-html绑定原始文本会暴露一堆##、**符号,体验极差。正确做法是每一次 token 增量后重新跑一次 Markdown 渲染,用渲染结果替换 DOM 内容。

这时候会有新的问题:重新渲染会让光标跳动或者滚动位置不稳定。原因很简单,Markdown 解析器输出的是完整的 HTML 字符串,替换 innerHTML 会重置 DOM 结构,导致用户阅读位置跳到页面顶部。解决思路是:

  • 用nextTick等 DOM 更新完成后,手动把滚动容器 scrollTop 设置到 scrollHeight;
  • 如果渲染量大,考虑做"节流渲染",比如每秒最多重新渲染两次,让打字效果不至于因为频繁解析 markdown 而卡顿;
  • 如果文本较长且包含表格、代码块这类复杂 Markdown,建议前端用marked.js加highlight.js的组合,解析性能比直接套 Vue 的v-html强很多。

在这个环节你要注意:异步流式状态下,不要直接用computed对fullText做marked.parse()然后绑到 v-html 上。因为每次fullText更新,computed 都会同步重新计算,一旦文本量大了,UI 线程会被 Markdown 解析阻塞,打字机效果直接卡成 PPT。我的处理方案是把marked.parse放到watch的回调里,或者干脆在 Message 对象上用requestIdleCallback做延迟解析,把渲染压力分散到浏览器空闲时段。

6. 线上常见问题排查实录:SSE 断连、超时与解析失败

6.1 "stream disconnected before completion" 的根因排查

这是我在评论区里被问得最多的一个错误。完整报错是类似stream disconnected before completion: idle timeout waiting for sse。我排查这类问题有一套固定流程,按顺序操作能快速定位:

  1. 先看是不是空闲超时。打开浏览器开发者工具,看 Network 面板里这条 SSE 请求的响应时间线。如果请求在某个平静期突然结束,曲线变平很久后断开,基本可以锁定是中间设备在超时。直接在后端生成器里加心跳帧(参考 3.3 节)。
  2. 再看代理层。把服务直连和走 Nginx 各测一次。如果直连正常、走 Nginx 断连,那九成是 Nginx 对长连接的 read timeout 设置太短。用我下面的配置:
    location /api/ { proxy_pass http://backend; proxy_http_version 1.1; proxy_set_header Connection ''; proxy_buffering off; proxy_cache off; proxy_read_timeout 300s; }
    注意proxy_http_version 1.1很关键,HTTP/1.0 不支持分块传输编码,SSE 会直接失效。
  3. 最后看应用层超时配置。如果你用 Gunicorn 跑 FastAPI,默认 worker 超时是 30 秒,模型生成超过 30 秒 Gunicorn 会强制杀 worker。启动命令里加--timeout 300。

6.2 常见问题速查表

下面这张表是我在实际项目中整理出来的高频问题,按"现象—原因—解决方案"列出来,排查时直接对照:

现象根因解决方案
前端转圈很久,最后一次性显示全部内容Nginx 缓冲导致数据攒批发送加X-Accel-Buffering: no响应头,proxy_buffering off
对话进行中突然断连,前端无错误提示空闲超时(idle timeout)生成器内加: heartbeat心跳帧
流式接口同时并发请求时变慢同步chain.stream()阻塞事件循环改用chain.astream()配合异步生成器
前端收到中文乱码TextDecoder没有设置{ stream: true }按 5.2 节代码加上
JSON.parse报错模型输出被 markdown 包裹或截断用 4.4 节的四层兜底解析
网页开着多个会话时,某些流式连接不定时断开浏览器对同一域名连接数限制收紧页面数量,或用 HTTP/2
LangChain 流式中间有长停顿模型内部在等待工具调用心跳帧 + 前端展示"正在处理"状态提示

6.3 一个让我印象深刻的坑:EventSource 的隐式缓冲陷阱

最后分享一个非常隐蔽的坑。某个页面我一开始用的 EventSource 实现打字机效果,单条消息没有任何问题。但在一个后台管理页里同时开了 5 个 EventSource(监听不同 schema 的推送),结果其中两个连接经常悄悄断掉,怎么查都查不出原因。

后来才反应过来:HTTP/1.1 协议中,浏览器对同一个域名的并发连接数限制是 6。页面里其他资源请求占了一部分,再开多个 EventSource,连超了之后浏览器就会挂起连接。加上我用的是 Docker 部署的 FastAPI,服务端自身没有维护长连接池,压力一大连接就有可能被系统层面回收。

解决方案有两个方向:要么改用 HTTP/2(多路复用让连接数限制不再成为瓶颈),要么像我后来做的,把多个 SSE 订阅合并成一个接口,用多 event 类型区分业务。对大部分中小项目来说,后者实现成本更低,效果也很直接。

写在最后,也是我踩过坑之后最想告诉你的

把这套方案完整落地之后,我回头看的最大的体会是:流式传输这件事,难的不是协议本身,而是链路中每一层的配合。LangChain 输出 Token、FastAPI 包装 SSE、Nginx 不捣乱、前端增量渲染——任何一环想当然,最后的体验就会打折扣。我自己最开始就是被"流式接口能跑通"这个表象骗了,前端拿到 token 就直接渲染,结果 JSON 全被拼在正文里、Nginx 缓冲把打字机效果吞成了整块输出、结构化解析失败直接让 Agent 中间状态崩掉。所以建议各位动手之前,先画一张自己项目的 SSE 链路图,标清楚每一环负责什么、可能在哪里出问题,再动手写代码,会少走很多弯路。

另外一个后来觉得非常值的工作,是把 SSE 的消息格式从"服务端想发啥就发啥"收敛成了"事件类型 + JSON 数据"的规范。只要消息流动的地方都遵循这个规范,前端就能精准地处理 token、meta、end、error 四类事件,后端加新功能(比如加个中间状态推送)也不需要前端跟着改逻辑。如果你打算长期迭代一个 AI 产品,我强烈建议你一开始就把这套消息约定定下来。

最后再分享一个日常调试的实用小技巧:先用curl测后端,再连前端。curl -N http://localhost:8000/chat -d '{"question":"你好"}',-N参数会禁用缓冲,让你直接在终端看到每一个 SSE 帧的到达时间。如果终端里能逐帧刷出来,后端基本没问题,再去排查前端;如果终端也是一次性输出,问题就在后端或网络层。这个习惯帮我省了至少一半的联调时间。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询