做 AI 应用开发,尤其是涉及到 LangChain 和流式对话的产品,你会发现所有体验的根基都压在一个词上:SSE(Server-Sent Events)。无论是打字机效果、实时日志推送还是结构化输出的增量解析,背后都是这条看似简单的 HTTP 长连接在支撑。我最早接到一个需求——让大模型回答像 ChatGPT 那样逐字输出,同时把返回结果里的关键字段抽成 JSON 给前端做结构化渲染。项目做完后回头看,这条链路里能踩的坑几乎全踩了一遍:SSE 格式写错导致前端无法解析、LangChain 结构化输出和流式模式互相打架、流式 JSON 在途中断开导致整段结果丢失。这篇文章把我验证过的完整方案整理出来,覆盖服务端 FastAPI、LangChain 输出约束、前端打字机渲染和 JSON 解析,适合正在搭建 AI Agent 对话产品或准备对智能体平台做二次开发的同学参考。
1. 为什么 AI 对话应用绕不开 SSE:流式协议的核心机制
1.1 从一次请求-响应说起
传统 HTTP 接口的工作方式是一条"完成即交付"的流水线:浏览器发起请求后,服务端把整段响应体攒齐,一次性返回。这套模型做普通 API 没有任何问题,但放到大模型对话场景里就非常别扭——一个稍微复杂一点的回答可能要生成几秒钟甚至几十秒,如果让用户盯着一个 spinner 干等,体验很糟糕。
更关键的是,大模型本身就是按照 token(可以粗略理解为一个字或一个词)逐个生成的。既然生成是流式的,传输自然也应该跟着流式走。ChatGPT 之所以让人觉得"它在思考、在打字",并不是什么前端动画技巧,而是服务端真的把 token 一个一个推送了下来。
SSE 解决的就是这个"服务端主动推送"的问题。它基于普通的 HTTP 协议,只是约定了一种特殊的数据格式:服务端把响应的 Content-Type 设为text/event-stream,然后持续向连接中写入事件流,客户端一边收一边渲染。整个过程不需要额外的协议升级,也没有 WebSocket 那么重的握手开销。
我第一次实现时也以为这只是前端接一个接口的事,真正做下去才发现这玩意儿有几个容易忽略的约束:一是默认只支持服务端到客户端的单向推送,双向通信得靠客户端另外发请求;二是它建立在长连接上,中间任何一层代理(Nginx、网关、负载均衡)都可能因为超时策略断开这条连接,这就是后文会详聊的 idle timeout 问题。
1.2 SSE 的数据格式与事件语义
SSE 的格式比大多数人想象中简单。每条消息由若干"字段行"组成,字段之间用\n分隔,消息之间必须有一个空行(也就是连续的\n\n)作为结束标记。最基本的字段是data:
data: {"content": "你好"} data: {"content": ","}一个事件流里可以包含多种字段:
data::消息体内容,可以有多行,多行会被拼接成一个字符串event::事件类型,默认值为message,客户端可以用addEventListener监听不同事件id::事件 ID,客户端断线重连时会带上这个 ID 请求增量数据retry::断线后客户端重连的等待毫秒数: comment:以冒号开头的注释行,通常用作心跳包,客户端解析时会自动忽略
字段开头不允许有空格,冒号后面如果跟空格,会被当作数据内容的一部分处理。这是很多新手容易犯的格式错误——我在本地调试时也遇到过,自己手动拼的 SSE 字符串里data:后面多打了一个空格,结果前端读到的 JSON 带了一个空格前缀,JSON.parse直接报错。
还有一个关键点:事件流里传输的通常是一段文本,而不是浏览器自动解析好的数据对象。也就是说,前端拿到data字段后,永远要自己做一次解析。这也引出了后面 JSON 解析的一系列问题。
1.3 SSE 和 WebSocket 怎么选
在技术选型上,很多团队会纠结"SSE 还是 WebSocket"。我的建议很简单:如果只是服务端单向推送数据,比如大模型 token 流、日志流、通知推送,优先选 SSE;如果有双向交互需求,比如聊天室、实时协作文档、在线游戏,才需要考虑 WebSocket。
原因也很实在。SSE 基于普通 HTTP,可以直接复用现有的认证体系(Cookie、Token、网关鉴权),中间走 Nginx 也不需要额外做 Upgrade 配置。WebSocket 虽然也是 HTTP 升级而来,但它会占用独立的连接通道,Nginx 默认配置下 60 秒无数据也会断开,反而更容易踩坑。
SSE 还有一个 WebSocket 没有的优势:自带断线重连机制。浏览器端的EventSource对象在连接断开后,会根据服务端返回的retry字段自动重连,并自动携带Last-Event-ID头。WebSocket 要自己实现心跳和重连逻辑,复杂度高出一截。不过在 LangChain 对话场景里,我们往往需要 POST 请求携带上下文参数,EventSource 只支持 GET,所以更常见的做法是使用fetch配合ReadableStream手动解析。这个方案我会在第三章详细展开。
2. LangChain 结构化输出:把"自由文本"变成可用的 JSON
2.1 结构化输出到底解决什么问题
如果你只是做一个纯聊天机器人,模型的自然语言回复本身就是最终产物,不需要额外加工。但实际业务里,AI 的输出往往还要喂给下游系统:从一段客服回复中抽取"用户情绪标签"和"工单类型",从一段分析报告中提取"风险点"和"建议操作",或者让 Agent 输出一个完整的任务计划供前端渲染成步骤卡片。
这些都要求模型输出严格符合一个预定 JSON Schema。早期做法是在 prompt 里写"请以 JSON 格式返回,字段包括 xxx",模型大部分时候能配合,但只要字段一多、嵌套一深,总会出现漏字段、多逗号、引号未闭合这类错误。与其在 prompt 层面碰运气,不如用 LangChain 提供的能力把输出约束变成模型层面的强制行为。
2.2 with_structured_output 的内部机制
LangChain 的with_structured_output是对"结构化输出"能力的一层封装。用法很直白:先定义一个 Pydantic 模型,再把它传给链或模型:
from pydantic import BaseModel, Field from langchain_openai import ChatOpenAI class TaskPlan(BaseModel): title: str = Field(description="任务标题") steps: list[str] = Field(description="执行步骤列表") priority: int = Field(description="优先级,1-5,数值越高越优先") llm = ChatOpenAI(model="gpt-4o", temperature=0) structured_llm = llm.with_structured_output(TaskPlan) result = structured_llm.invoke("帮我规划一个周末学习计划") print(result.model_dump())这段代码背后发生了几件关键的事情:
第一,LangChain 会把TaskPlan转换成 JSON Schema,塞进模型接口的tools参数里。如果底层模型支持 function calling / tool calling(OpenAI 系、Claude 系、Qwen 系都支持),模型就会被强制按照这个 Schema 生成输出,LangChain 拿到响应后自动解析成 Pydantic 对象。
第二,如果底层模型不支持 tool calling,LangChain 会走一条降级路线:把 JSON Schema 写进系统提示词,再用输出解析器兜底。这条路线的稳定性完全取决于模型能力,所以我会在选择基础模型时优先确认它是否原生支持 tool calling。
第三,带with_structured_output的调用默认是非流式的。模型要等所有 token 生成完、解析成对象后一次返回。这跟对话场景需要的打字机效果天然冲突。
2.3 流式场景下结构化输出的取舍
回到实战场景:我需要打字机效果,又需要结构化字段。这时候有两条路可以走。
第一条路是双向奔赴——在同一套模型调用里既拿流式文本,又拿结构化结果。LangChain 在高版本里对部分模型支持流式结构化输出,但实际用下来兼容性参差不齐,API 还在演进中。我在一个 Agent 项目里试过用astream配合with_structured_output,有的模型能正常返回增量 JSON 片段,有的模型直接报错说 stream mode 与 structured output 不兼容。这种不确定性在生产环境里很难接受。
第二条路是拆分诉求——流式文本走一条链路,结构化字段走另一条链路。具体做法是:先用普通流式接口让模型逐 token 输出对话文本,同时让模型在文本中用特殊标记包裹一个 JSON 块;前端在流式渲染文本的同时,从增量数据中提取 JSON 块做解析。更稳妥的设计会在第四章讲。
这里有个值得记下的经验:结构化输出和流式输出本质上是两个需求。结构化的核心诉求是"严格可信",流式的核心诉求是"实时可见",两者不要强行糅在一个接口里。与其纠结 LangChain 的流式结构化 API,不如把数据协议设计成"增量事件流 + 最终结果事件",各取所需。这个思路在对接 LangGraph 或 DeerFlow 这类智能体框架时尤其管用。
3. 打字机效果全链路:FastAPI 流式接口 + 前端增量渲染
3.1 服务端:StreamingResponse 的正确姿势
服务端我用的是 FastAPI,配合 LangChain 的异步流式接口astream。一个最基础的 SSE 流式接口长这样:
import asyncio import json from fastapi import FastAPI from fastapi.responses import StreamingResponse from langchain_openai import ChatOpenAI app = FastAPI() llm = ChatOpenAI(model="gpt-4o", temperature=0.7) async def event_generator(prompt: str): # 先发一个事件,告知前端流开始 yield f"event: start\ndata: {json.dumps({'message': '开始生成'})}\n\n" async for chunk in llm.astream(prompt): content = chunk.content if content: # 这里把每个 token 包装成 SSE 事件 yield f"data: {json.dumps({'delta': content})}\n\n" # 流结束事件,前端收到后可以做收尾 yield f"event: done\ndata: {json.dumps({'message': '生成完成'})}\n\n" @app.post("/chat/stream") async def chat_stream(payload: dict): prompt = payload.get("prompt", "") return StreamingResponse( event_generator(prompt), media_type="text/event-stream", headers={ "Cache-Control": "no-cache", "X-Accel-Buffering": "no", "Connection": "keep-alive", }, )有几个细节必须交代清楚。
X-Accel-Buffering: no是给 Nginx 看的。Nginx 默认会缓冲上游响应,如果不关掉,SSE 的数据会被攒在代理层,前端收不到增量,打字机效果直接失效。这个头的作用就是告诉 Nginx"别缓冲,往客户端直推"。
为什么要用event: start和event: done区分事件类型?因为前端渲染逻辑需要知道边界——收到start时清空上一次的状态,收到done时停止 loading 并尝试解析最终 JSON。如果所有数据都混在一个message事件里,前端就不得不在数据内容里约定特殊标记,容易产生歧义。
还要注意别在流式生成器函数里做耗时初始化。StreamingResponse 一旦启动,生成器内的代码就开始执行,但如果你的生成器先花 10 秒去查数据库、调用外部认证服务,那么客户端 TTFB(首字节时间)会很难看,甚至触发前端的超时。正确做法是在生成器外先把需要的上下文、历史消息、链路 ID 都准备好,生成器只负责"等模型、透传 token"。
3.2 前端:用 fetch 读取 text/event-stream
浏览器端的EventSource只支持 GET 请求,而聊天接口往往需要 POST 携带 prompt、角色、上下文等参数,所以前端更通用的做法是直接用fetch读取流式响应:
async function streamChat(prompt) { const response = await fetch('/chat/stream', { method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify({ prompt }), }); if (!response.ok) { throw new Error(`HTTP ${response.status}`); } const reader = response.body.getReader(); const decoder = new TextDecoder('utf-8'); let buffer = ''; while (true) { const { value, done } = await reader.read(); if (done) break; buffer += decoder.decode(value, { stream: true }); // SSE 消息以空行分隔,按 \n\n 切割事件块 const chunks = buffer.split('\n\n'); buffer = chunks.pop(); // 最后一段可能是不完整的,留到下次拼接 for (const chunk of chunks) { handleSSEChunk(chunk); } } } function handleSSEChunk(chunk) { const lines = chunk.split('\n'); let eventType = 'message'; const dataLines = []; for (const line of lines) { if (line.startsWith('event:')) { eventType = line.slice(6).trim(); } else if (line.startsWith('data:')) { dataLines.push(line.slice(5).trimStart()); } } if (dataLines.length === 0) return; const dataStr = dataLines.join('\n'); let data; try { data = JSON.parse(dataStr); } catch (e) { console.error('SSE data 解析失败:', dataStr); return; } if (eventType === 'start') { resetUI(); } else if (eventType === 'done') { handleDone(data); } else if (data.delta) { appendDelta(data.delta); } }这段代码里最值得注意的是decoder.decode(value, { stream: true })。流式读取时,一个 UTF-8 字符可能被拆分在两个 chunk 里(尤其是中文、emoji 这类多字节字符),如果不加stream: true,边界上的字节被单独解码会变成乱码。另外,buffer的暂存逻辑很重要——TCP 传输不保证每帧数据正好以\n\n结尾,所以必须先把所有数据拼进缓冲区,再按完整分隔符切割,剩余未结束的部分留到下一轮。
3.3 渲染层:Markdown 增量渲染与性能优化
拿到delta之后怎么渲染,其实是个大坑。如果只是纯文本,直接往 DOM 里 append 就行。但大多数对话产品用的是 Markdown 格式,问题就来了:Markdown 是块级语法,一个代码块还没闭合、一个列表还没结束的时候,如果每一帧都重新解析整段 Markdown,会产生明显的闪烁和跳动;如果完全不重渲染,代码块的高亮又不会更新。
我实测下来最稳的做法是双轨制。对话内容区用一个只读容器承载已经确认完整的 Markdown 渲染结果,另用一个"增量缓冲层"积累最新的 1-2 个 token 原始文本。具体来说:
- 当缓冲区累计了足够长度(比如 50 个字符),或者遇到明显的语法边界(比如换行符、代码块结束标记
```),就把缓冲区内容拍平进完整文本,重新渲染整块 Markdown - 在缓冲区还没拍平之前,只在容器的末尾追加一小段纯文本占位,避免频繁全量重渲染
渲染时机上,建议用requestAnimationFrame做节流。模型生成 token 的速度远高于浏览器的渲染帧率,每收到一个 token 就同步更新 DOM 很容易造成页面卡顿。正确思路是把 token 先推入一个队列,在下一帧统一批量提交,这样 UI 的刷新节奏和浏览器渲染节奏保持一致。
还有一个经验:不要在流式过程中直接做代码高亮。高亮计算很消耗性能,尤其是长回复。我的做法是先渲染纯 Markdown 结构,流结束后再对代码块做一次高亮,用户几乎察觉不到这个延迟。
4. 流式 JSON 解析的实战方案:从暴力等待到增量解析
4.1 最简单的方案:等全部输出完再解析
如果聊天回复里末尾带一个 JSON 块,比如:
以下是你需要的任务计划: { "title": "学习计划", "steps": ["复习数学", "做物理题", "读英语"] }新手最常见的写法是:流式过程中只渲染文本,等done事件到达后,把完整文本里的 JSON 块用正则抠出来,再JSON.parse。这个方案实现简单,也不会出错,但有两个明显问题。
第一,"答案要等全部生成完才出现",用户能明显感觉结构化内容是"突然蹦出来"的,而不是跟随打字机效果逐步呈现。如果 JSON 块位于文本末尾,用户要看完一大段文字才能看到结构化结果。第二,对长输出来说,流结束后一次性解析大 JSON 会有轻微卡顿,体验不够顺滑。
但必须承认,对于一些边缘业务(比如只需要最终结果、中间状态无所谓的后端任务),"流式展示 + 最终解析"依然是最稳定、最值得推荐的方案。复杂的东西不一定就好,稳定可靠才是第一位的。
4.2 增量解析:括号配对与可恢复 JSON 解析
如果真的需要边接收边解析,主动权其实在服务端手里。如果服务端能保证流式数据是一个"正在生成中的合法 JSON",那么前端可以做增量解析——每次拿到新数据后,尝试解析当前缓冲区里的内容,成功则更新界面,失败则继续等待。
这里的关键是:判断当前缓冲区里的 JSON 是否"已经完整"。我写过一个简单的栈式判断器:
function isJsonLikelyComplete(str) { let depth = 0; let inString = false; let escaped = false; for (let i = 0; i < str.length; i++) { const ch = str[i]; if (escaped) { escaped = false; continue; } if (ch === '\\') { escaped = true; continue; } if (ch === '"') { inString = !inString; continue; } if (inString) continue; if (ch === '{' || ch === '[') depth++; else if (ch === '}' || ch === ']') depth--; } return depth <= 0 && !inString && str.trim().length > 0; }然后每次收到新 token,把增量追加到缓冲区,先判断"看起来完整",再尝试JSON.parse:
function tryParseIncremental(accumulated) { if (!isJsonLikelyComplete(accumulated)) { return { status: 'pending', data: null }; } try { return { status: 'done', data: JSON.parse(accumulated) }; } catch (e) { // 括号配平了但 parse 失败,说明中间有语法错误,等待更多字符 return { status: 'error', data: null }; } }这个方案在实际项目中能用,但要意识到模型输出的 JSON 可能存在各种脏数据:转义的反斜杠、字符串里的换行、未闭合的 Unicode 字符。一旦isJsonLikelyComplete返回 true 但JSON.parse一直失败,就不能死等,要设置一个最大等待帧数或超时时间,超时后放弃增量解析,等done事件后用完整内容做最终解析兜底。
4.3 工程推荐:事件分离的流式协议设计
做了几个项目之后,我越来越倾向一种更工程化的设计——不在数据层面硬挤,而是在协议层面把"结构化 JSON"和"文本增量"分离。
服务端先发一个meta事件,携带这次响应的整体结构说明(比如有哪些字段、每个字段的展示类型);然后流式发送delta事件,只携带文本增量;流结束后发一个done事件,携带最终完整 JSON。前端的行为就变得清晰了:
delta只负责打字机渲染meta驱动结构化组件的骨架done驱动结构化组件的最终数据填充
这种设计的最大好处是,结构化解析不再依赖"字符串里某个位置的 JSON 块"这种脆弱的约定,而是由事件类型明确保证。而且done事件里的 JSON 是服务端在拿到模型完整输出后自己解析好的,保证合法性。前端即使对增量内容一无所知,也能在流结束后正确展示结构化结果。
我自己在给 Agent 平台做二次开发时就是按这个协议设计的:工具调用信息、中间思考过程、最终回复三者的展示节奏完全不同,只有事件分离才能让前端逻辑保持清爽。
5. 踩坑实录:idle timeout 断连问题的完整排障链路
5.1 问题现象与根因分析
我在项目中第一次遇到这个报错时印象很深:前端控制台打出stream disconnected before completion: idle timeout waiting for sse,症状是流式回复生成到一半突然中断,后端看 LangChain 日志明明还在正常输出,前端却已经收不到任何数据。
根因其实分三层。
第一层是代理层超时。几乎所有代理/网关都有"空闲超时"配置:如果一段时间内连接上没有数据流动,代理就认为这条连接已经闲置,主动断开。Nginx 的proxy_read_timeout默认是 60 秒;AWS 的 ALB 默认 idle timeout 也是 60 秒;云厂商 API 网关的默认值普遍在 30-60 秒之间。
第二层是数据生成节奏太慢。大模型生成两个 token 之间可能有明显的停顿,尤其是开启 reasoning/thinking 能力时,模型可能先"思考"几十秒,再一次性输出。如果模型思考时间超过了代理的等待阈值,连接就会被断开。我曾经排查过一个案例:Agent 在调用工具时卡了 40 秒没输出,前端直接断流,用户体验就是"AI 转圈转到一半突然报错"。
第三层是服务端代码里有同步阻塞调用。Python 的 FastAPI 是异步框架,但如果你在流式生成器里用了requests.post、同步数据库查询这类阻塞操作,整个事件循环会被卡住,SSE 连接自然会"看起来很空闲"。
5.2 代理层配置修复
如果连接结构是"客户端 -> Nginx -> FastAPI -> LLM 服务",优先处理 Nginx 这一层:
location /chat/stream { proxy_pass http://backend; proxy_http_version 1.1; proxy_set_header Connection ""; proxy_buffering off; proxy_cache off; proxy_read_timeout 600s; proxy_send_timeout 600s; proxy_buffers 4 16k; proxy_buffer_size 8k; }proxy_buffering off是最关键的一项,它会强制 Nginx 拿到上游数据立刻转发给客户端,而不是攒在一起。proxy_read_timeout调大到 600 秒,给 LLM 的"思考停顿"留足空间。
如果你用的是云上负载均衡,在控制台里找到 idle timeout 配置,全部调大。数据库里这条经验我反复用:所有链路节点——LB、CDN、Nginx、网关——都要检查一遍,少一个都可能成为瓶颈。
5.3 服务端保活与生成节奏优化
代理层调完之后,服务端自己也要做好两件事。
第一件事是心跳保活。在生成器中周期性地发送注释行,让连接始终有数据在流动:
async def event_generator(prompt: str): async def heartbeat(): while True: yield ": ping\n\n" await asyncio.sleep(15) # 用 asyncio 把模型生成流和心跳流合并这里用的是 SSE 的注释行做保活——客户端解析时会忽略以冒号开头的行,但在网络层看来,这条连接一直有数据流动,就不会被判定为 idle。心跳间隔建议小于代理层超时时间的一半。
第二件事更本质:别让流式生成器出现"长时间无输出"。在接入 LangChain/Agent 时,如果中间要调用工具或进行多轮推理,应该在每个阶段切换时主动发一个status事件,比如event: status\ndata: {"stage": "calling_tool"}\n\n。这样前端能看到"AI 正在调用工具"的中间状态,用户感觉系统是活跃的,连接也因为有数据而保持存活。
我在接入 LangGraph 时会把astream的多个事件流(token、工具调用、状态更新)分发给不同的 SSE 事件类型,既能解决超时问题,又天然实现了第四章说的事件分离协议。
6. 进阶封装:把 SSE 调用、解析、重连封装成可复用模块
6.1 SSEStreamClient 的整体设计
踩完了所有坑之后,你会发现流式接入的代码其实高度重复,值得封装成一个内部通用模块。我在项目中维护了一个轻量级的SSEStreamClient,核心接口包括:
class SSEStreamClient: def __init__(self, base_url, on_event=None, ...): self.base_url = base_url self.on_event = on_event self.last_event_id = None self.retry_ms = 3000 async def connect(self, path, payload): # 发起 POST 请求,逐行解析 SSE ... async def _read_stream(self, response): # 解析缓冲区域、切割事件块、分发回调 ... async def _reconnect(self, path, payload): # 断线重连,根据 Last-Event-ID 续传 ...关键设计点有三个。
第一,_read_stream内部封装了多字节解码、缓冲区切分、事件类型分发,上层业务只注册on_event(event_type, data)回调,完全不关心底层 HTTP 细节。
第二,重连逻辑要有指数退避。第一次断线等 1 秒,第二次等 2 秒,最多等 30 秒,同时要规避服务端返回的retry字段覆盖这个节奏。重连时要带上Last-Event-ID头,这样后端可以实现增量续传,避免前端把整个长回复重新接收一遍。
第三,整个 client 要把"网络异常"和"业务异常"区分开。连接被断开是网络异常,走重连;服务端发了event: error且带错误码,走业务错误回调,此时重连没有意义。
6.2 对接 LangGraph/DeerFlow Agent 的二次开发思路
现在很多 Agent 平台(包括 LangGraph 的 Agent 服务、DeerFlow 这类智能体框架)已经自带了流式事件能力,二次开发时最常做的事就是把这些内部事件转换成我们自定义的 SSE 协议。
以 LangGraph 为例,调用图时可以指定stream_mode:
events = graph.astream( {"messages": [{"role": "user", "content": prompt}]}, config={"recursion_limit": 50}, stream_mode="messages" # 取 token 增量 )在这个流式迭代里,每个节点执行时都会产生不同的事件。我通常的做法是:在 FastAPI 层写一个适配器,把 LangGraph 的事件流转换为标准 SSE 事件:
- 图开始运行时发送
event: graph_start - 节点切换时发送
event: node_status,附带节点名和状态 - 模型产出 token 时发送
event: delta,附带文本增量 - 图运行结束发送
event: done,附带最终完整结果
这样做的好处是,前端只需要对接自己的 SSEStreamClient,后续无论底层换成 LangGraph、DeerFlow 还是别的智能体框架,前端代码几乎不用改。把"框架输出"和"产品协议"解耦,是所有二次开发项目里最值得投入的一步。
我对这套封装最满意的地方在于:它把服务端的流式能力变成了一个统一出口,所有业务模块共用同一条流式管道,新接入一个 Agent 只需要写一段适配器,不需要再碰前端。实际上这条经验已经帮我在多个项目里省掉了大量重复排障的时间——连接超时、缓冲区截断、重连失败这些常见问题,在封装层就被兜住了,业务侧基本不会感知到。