ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

从SSE原理到LangChain结构化输出:流式AI应用实战

从SSE原理到LangChain结构化输出:流式AI应用实战 做流式 AI 应用的人应该都有同感需求文档写着实现打字机效果一开始以为只是前端response.text边到边渲染真正动手才发现这句话背后牵扯的是 HTTP 协议、SSE 报文格式、后端事件流设计、LangChain 结构化输出甚至还有网关超时、半截 JSON 这些生产环境才会冒出来的麻烦。这篇文章我就以从 SSE 原理到 LangChain 结构化输出这条技术线为主线把我在 FastAPI LangChain LangGraph 项目里反复折腾打字机效果 JSON 解析的经验完整拆开讲。适合正在做 AI Agent 应用、需要让模型既流式输出文本、又能稳定返回结构化数据给下游系统的开发者阅读。1. 为什么打字机效果必须重新理解 HTTP 响应1.1 普通 HTTP 请求为什么做不到逐字返回绝大多数人第一次接触后端接口时认知模型是请求一次响应一次。服务端把完整数据计算完一次性通过 HTTP Response Body 返回给前端前端拿到JSON.parse(response)之后渲染。这个模式对普通业务接口完全够用但到了大模型对话场景就非常别扭。假设模型生成一篇 500 字的回答如果等服务端全部生成完再返回普通场景下可能需要 5 到 10 秒用户盯着空白页面干等体验极差。更重要的是大模型接口本身的流式返回是模型推理时的天然能力——每一个 token 生成后就可以发出去没有必要等完整句子拼完。我的理解是所谓打字机效果本质上不是前端动画技巧而是把响应时间从总耗时拆成首 token 时间的问题。只要第一条数据能尽快到达浏览器用户就会觉得系统开始工作了后续逐字渲染只是锦上添花。所以核心转变在这里不能把 AI 响应当成一次性 JSON 返回而是当成一个持续打开的管道服务端不断往管道里写入数据前端持续读取。这正是 Server-Sent Events 的用武之地。1.2 浏览器 EventSource 和 fetch 流读取有什么本质区别很多教程一提 SSE 就默认用EventSource但实际做 AI 项目我很快就发现EventSource有先天限制它只支持 GET 请求无法自定义 Header。而很多 AI 接口需要带 Authorization 认证你把 token 放在 query 参数里既不安全也会被日志系统记得到处都是。所以更合理的方案是用浏览器的fetchAPI 配合ReadableStream手动读取 SSE 流。fetch 支持 POST 和自定义 Header也能拿到response.body上的流对象。区别在于EventSource帮你做了协议解析和自动重连而 fetch 流需要自己处理缓冲区、拆分事件、识别断流。后者灵活度更高也更适合封装成一个统一的流式调用函数。一句话总结小 Demo 用 EventSource 可以真实项目请用 fetch ReadableStream。1.3 为什么大模型场景普遍选 SSE 而不是 WebSocket我刚接触这个领域时有一个疑惑WebSocket 不是也能双向通信吗做打字机效果不是更合适后来被几个现实问题教育了。WebSocket 是双向全双工协议连接建立后客户端服务端可以互相推消息能力确实更强。但它带来的复杂度也是成倍的需要处理握手升级、心跳保活、消息帧格式、断线重连状态同步。而 AI 对话这个场景绝大多数时间都是客户端发一个问题、服务端持续推一堆 token请求方向是单向的客户端并不需要频繁向服务端推消息。这种情况下用 WebSocket 属于高射炮打蚊子。SSE 是建立在 HTTP 之上的轻量协议服务端可以随时通过同一个连接把数据推给客户端。它有几个非常务实的特点自动重连通过 Last-Event-ID 断点续传、文本格式易懂、可以穿过大多数 HTTP 代理、不需要额外协议升级。OpenAI、Anthropic 这些大模型 API 的输出流也基本都走 SSE 或类似的分块传输流所以学习成本很低。2. SSE 协议拆解data、event、id、retry 到底在干什么2.1 SSE 的数据格式远比想象中简单SSE 不是一个复杂协议它本质上就是 HTTP 响应里Content-Type被设为text/event-stream之后服务端持续输出特定格式的文本块。每个事件由若干字段组成常见的有event: message data: {type: token, content: 你好}中间的空行代表一个事件的结束。如果一段数据里有多行data:它们会被拼接成一个字段中间以换行符连接。event:用来给事件命名比如我常用event: token和event: result区分中间输出和最终结果。id:则是给事件编号用于客户端断线重连时通过Last-Event-ID告诉服务端我收到哪条了从它后面继续发。实际写后端时最简单的 SSE 输出就是循环yield字符串async def event_stream(): for token in [你, 好, 世, 界]: yield fevent: token\ndata: {token}\n\n注意一个细节SSE 的字段类型是按行的文本协议不是 JSON。所以把data里的内容当成 JSON 字符串传递是完全可行的但data:本身不是 JSON。很多初次接触的人会把data: {type: token}整体当成一个 JSON 对象去解析容易在data:前缀上栽跟头。2.2 心跳注释、retry 重连和 Last-Event-IDSSE 规范里还有一个容易被忽略的元素以冒号:开头的行是注释行它不会被当作数据触发事件主要作用是让连接保持活跃。为什么需要它因为很多代理服务器、负载均衡器如果看到一段时间内没有任何数据传输会主动掐断连接。这个问题在 AI Agent 场景特别突出后面我会单独说。所以在长耗时连接里服务端最好每隔十几秒发一行注释作为心跳: ping另外retry:字段可以告诉客户端如果断开了等待多少毫秒再重新连接。比如retry: 3000客户端收到这个字段后会以 3 秒为间隔尝试重连。在真实聊天应用里这个值不宜太短否则服务端一旦重启客户端会疯狂打请求。我自己的实践是事件流中既有消息类事件也有心跳事件前端解析时必须忽略没有任何data字段的注释行否则会把空事件当成无效 JSON 抛异常。3. 用 FastAPI 后端真正把打字机效果跑起来3.1 StreamingResponse 输出 SSE 流的最小实现FastAPI 做 SSE 非常方便核心是StreamingResponse。它接受一个异步生成器或普通生成器设置media_typetext/event-stream框架就会按流式响应输出。下面这个例子是从 FastAPI 接口返回一个模拟 token 流from fastapi import FastAPI from fastapi.responses import StreamingResponse import asyncio import json app FastAPI() def format_sse(data: dict, event: str message) - str: return fevent: {event}\ndata: {json.dumps(data, ensure_asciiFalse)}\n\n async def generate_tokens(): for token in [我, 是, 一, 个, 智, 能, 助, 手]: await asyncio.sleep(0.1) yield format_sse({type: token, content: token}, eventtoken) app.post(/chat) async def chat(): headers { Cache-Control: no-cache, Connection: keep-alive, X-Accel-Buffering: no, } return StreamingResponse(generate_tokens(), media_typetext/event-stream, headersheaders)这里有三点必须解释第一Cache-Control: no-cache防止浏览器或中间代理缓存整个响应。如果被缓存客户端要等流结束才能拿到数据打字机效果直接报废。第二X-Accel-Buffering: no是给 Nginx 这类反向代理看的。Nginx 默认会缓冲后端响应到一个阈值才转发这会让 SSE 失去实时性。显式关掉缓冲数据才能逐条到达前端。第三Connection: keep-alive告诉网络链路不要关闭 TCP 连接避免流被中途切断。3.2 前端 fetch 流读取处理跨 chunk 的拆包问题前端用 fetch 读取流时最容易踩的坑是网络包并不是按\n\n整齐切分的。一次read()可能返回半行、上一行的尾巴加下一行的开头也可能一次返回多个完整事件。如果直接对每次读取的字符串做 JSON.parse大概率会频繁报错。我习惯用一个SSEParser类来维护缓冲区class SSEParser { constructor() { this.buffer ; } push(text) { this.buffer text; const events []; let index; while ((index this.buffer.indexOf(\n\n)) ! -1) { const raw this.buffer.slice(0, index); this.buffer this.buffer.slice(index 2); events.push(this.parseEvent(raw)); } return events.filter(Boolean); } parseEvent(raw) { const data []; let eventName message; for (const line of raw.split(\n)) { if (line.startsWith(:)) continue; if (line.startsWith(event:)) { eventName line.slice(6).trim(); } else if (line.startsWith(data:)) { data.push(line.slice(5).trim()); } } if (data.length 0) return null; try { return { event: eventName, data: JSON.parse(data.join(\n)) }; } catch { return { event: eventName, raw: data.join(\n) }; } } }然后用它连接 fetch 的 ReadableStreamconst response await fetch(/chat, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ query: 你好 }), }); const reader response.body.getReader(); const decoder new TextDecoder(utf-8); const parser new SSEParser(); while (true) { const { done, value } await reader.read(); if (done) break; const text decoder.decode(value, { stream: true }); for (const event of parser.push(text)) { if (event.event token) { appendToUI(event.data.content); } else if (event.event result) { handleFinalResult(event.data); } } }注意TextDecoder的stream: true参数。中文 UTF-8 编码在三字节字符可能被拆到两个网络 chunk 里如果不流式解码尾部会出现乱码。这是一个非常细节但常见的问题。3.3 为什么还要封装一层SSE 调用逻辑做多个会话功能时我不会在每个页面都写一遍 fetch 解析逻辑。更好的做法是把请求 解析 回调 重连封装成一个统一的流式客户端函数对外暴露onToken、onResult、onError回调。这样产品页面只关心业务回调不感知 SSE 协议细节。这也能解决热词里提到的封装 SSE 流式接口调用逻辑完成流式消息解析这一诉求。一个最小封装示意async function streamChat(url: string, payload: object, handlers: { onToken: (text: string) void; onResult: (data: unknown) void; onError: (err: Error) void; }) { // fetch parser 分发回调 }4. LangChain 结构化输出让模型不只聊天还能交数据4.1 为什么需要结构化输出而不是让模型自己吐 JSON单纯让模型生成一段看起来像 JSON 的文本非常不可靠。大模型会在 JSON 外面加 Markdown 代码块、加解释性文字、甚至把字段名写错。如果下游系统要拿这个 JSON 入库或者触发工具调用解析失败会导致整个流程崩溃。LangChain 提供的结构化输出能力本质上是利用当前大模型的 function calling / tool calling 能力让模型在生成层就输出符合某个 Schema 的 JSON。我不需要自己去猜模型输出而是定义好 Pydantic 类让它按照类的字段生成。我常用的写法是from pydantic import BaseModel, Field from langchain_openai import ChatOpenAI class WeatherResponse(BaseModel): city: str Field(description城市名) temperature: float Field(description当前温度) condition: str Field(description天气状况如晴、雨、多云) llm ChatOpenAI(modelgpt-4o, temperature0) structured_llm llm.with_structured_output(WeatherResponse) result: WeatherResponse structured_llm.invoke(北京今天怎么样) print(result.city, result.temperature)with_structured_output会从 Pydantic 类生成 JSON Schema再通过 tool calling 机制让模型返回满足 Schema 的参数。拿到的是已经校验过的 Pydantic 对象不是需要自己json.loads的字符串。这样无疑稳得多。4.2 输出解析器路线传统 OutputParser 与 JSON 修复LangChain 里还有一套更底层的输出解析器比如JsonOutputFunctionsParser、PydanticOutputParser。如果你的模型不支持 tool calling可以走提示词 输出解析器路线。但那条路可靠性低需要额外做 JSON 修复。实际项目里我遇到最多的两类问题一是模型在 JSON 前后加了多余内容二是 JSON 本身不完整。对前者常做的是在 System Prompt 里强调只输出 JSON不要任何解释性文字再用解析器去除代码块标记。对后者可以用json_repair这类库修复部分错误。from json_repair import repair_json raw json\n{city: 北京, temperature: 25, fixed repair_json(raw) print(fixed) # 可能修复为 {city: 北京, temperature: 25, ...}但我要提醒一点json_repair能修复一些常见错误但不是万能的。如果需要 100% 可靠性最好在生成侧用 tool calling 机制提前约束。4.3 流式输出和结构化输出天然存在矛盾这是整篇最关键的问题。打字机效果要求流式输出结构化输出要求完整稳定的 JSON。两者方向不一样。with_structured_output在多数模型上并不支持直接在流式中间过程给出结构化对象。你可能拿到的是一个一个 token需要自己在后端拼出完整 JSON。如果模型选择 tool calling 模式LangChain 提供stream时会给出包含函数调用参数片段的增量块这些片段往往不是合法的 JSON不能直接JSON.parse。所以在做 Agent 应用时我通常会做一个取舍业务层把中间 token 流和最终结构化结果分成两件事件输出。中间过程走打字机最终结果走event: result。这样既保证用户体验又保证下游数据结构完整。5. 增量 JSON 解析让结构化数据和 token 一起到达5.1 为什么不能等完整 JSON 到达后再渲染结构化内容有些场景下用户希望前端能实时看到结构化卡片在逐步成形。比如一个 Agent 要返回航班、酒店、景点三个对象如果能边生成边渲染卡片体验会更好。这时就需要增量 JSON 解析。问题是增量 JSON 可能出现严重的不完整例如模型只生成了{flight: {airline: 中国国航, departure: 北京, arrival: 上这个字符串不是合法 JSON普通JSON.parse直接抛异常。于是需要增量解析器它能在数据不完整时尽力提取出已经定义完整的字段。Python 端我知道的库有partial-json-parser前端 JS 也有类似思路。5.2 稳定性优先后端累积 JSON前端渲染已确认字段我对增量解析的态度是可以在前端做但一定要降级。最稳妥的方案不是直接解析半截 JSON而是在后端维护一个结构化结果缓存。每个 chunk 到达时尝试把累计字符串解析成 JSON如果成功就把这次解析得到的对象推给前端如果失败就只发送普通 token 事件。这样前端永远只处理解析成功的结构化片段不会把错误状态展示给用户。举个例子async def generate_structured_stream(): full_text async for token in langchain_stream: full_text token yield format_sse({type: token, content: token}, eventtoken) try: partial json.loads(full_text) yield format_sse({type: partial, data: partial}, eventpartial) except json.JSONDecodeError: pass这样做有个额外好处JSON 一旦完整最后一次 partial 事件就是最终结构化结果。前端不需要在完全独立的事件里重复接收。5.3 设计事件类型时一定要分清楚语义我见过不少项目把 token、message、final result 全部塞进同一个message事件前端靠JSON.parse后猜字段名。这在字段结构简单时还能跑一旦业务复杂就非常痛苦。我比较推荐用 SSE 的event字段区分语义event: token普通文本片段用于打字机渲染。event: partial当前已解析成功的结构化对象片段用于增量卡片。event: result最终完整结构化结果用于提交数据或落库。event: error业务错误信息。event: done整个流结束。这样的好处是前端代码逻辑清晰每个事件类型对应一个处理函数不需要再做一堆 if 判断。也是封装 SSE 流式接口调用逻辑时最值得统一的部分。6. 集成 LangGraph Agent 时的真实踩坑空闲超时、断流、半截 JSON6.1 stream disconnected before completion: idle timeout waiting for SSE这个报错如果做生产部署大概率会碰到。它的含义是某个中间代理或网关在长时间没有读到数据时主动关闭了流。LLM 直出文本时token 会一直有输出不太容易触发空闲超时。但换成 LangGraph Agent 后情况完全不同。Agent 在内部进行工具调用、规划下一步、等待子任务返回时可能好几秒没有任何输出。比如我需要 Agent 调用一个外部搜索接口搜索耗时 8 秒这 8 秒里 SSE 连接上没有任何数据。Nginx 默认proxy_read_timeout往往只有 60 秒AWS 的 ALB 超时甚至更短。一旦超过时限代理就会掐断连接前端收到stream disconnected before completion。解决办法主要有三个第一在后端生成器里必须实现应用层心跳。每隔固定时间比如 15 秒发一个 SSE 注释行或一个event: heartbeat的空事件。关键是让网络链路感知到连接还活着而不是真的没有数据。async def generate_with_heartbeat(): while True: if new_token_available(): yield format_sse({type: token, content: token}, eventtoken) else: yield : heartbeat\n\n await asyncio.sleep(0.1)注意不要 sleep 太久心跳间隔太长等于没有。第二调整反向代理的超时配置。如果前面有 Nginx至少要像下面这样设置proxy_buffering off; proxy_cache off; proxy_read_timeout 300s; proxy_send_timeout 300s; proxy_http_version 1.1;proxy_buffering off是为了不缓冲 SSE 分片proxy_http_version 1.1是为了支持 keep-alive 长连接。第三前端要有自动重连机制。可以用Last-Event-ID让重连后的连接从断开位置继续接收。但如果 Agent 已经执行完了重连只会拿到一个结束标记前端要做好处理不能把重复事件叠加渲染。6.2 半截 JSON工具调用流式参数片段LangGraph 里如果 Agent 返回结构化数据另一个坑是stream输出里函数调用的参数是分片出现的。比如某个返回经过with_structured_output的节点stream()产生的事件可能是{function_call: {name: search, arguments: {\query\: }} {function_call: {name: search, arguments: 上海天气}}这些片段不是完整的 JSON如果直接在回调里尝试解析会收获一堆异常。我的处理方式是在后端节点里就把最终工具调用的完整结果收集起来等AIMessageChunk积累完整之后再序列化输出。不要把流式分片直接暴露给前端。6.3 通知型 SSE 和对话型 SSE 不要混在一起热词里出现了通知sse指的是那种服务器主动推送通知、客户端只需要接收的场景。它和对话型 SSE 最大的区别是通知型没有请求-响应的对应关系往往只有一个共享连接。对话型则每个请求都独立开流。我在一个项目里曾经试图用同一个连接处理所有 Agent 任务的消息结果消息串流、前端无法分辨哪条属于哪个会话。后来老老实实改成每个对话请求一个独立 EventSource/fetch 流并发连接数确实会高但逻辑清晰得多。如果你要做实时通知可以单独建一条通知 SSE如果要做对话 AI不要混用。两者混在一起最后维护成本高到你想重构。收尾一点个人的工程体会做完整套方案后我的体会是别把打字机效果当作前端问题也别把JSON 解析当作后端问题。它们本质上是同一个问题——如何在一个持续打开的流式通道里稳定地传递和还原不同粒度的数据。如果让我重新设计一个最小可用的 AI Agent 流式接口我会这样组合后端用 FastAPI 的StreamingResponse输出 SSE事件类型严格区分为 token / partial / result / error / doneLangChain 层用with_structured_output定义好 Pydantic Schema最终结果只在 result 事件里完整发送前端用 fetch 流读取封装一个全局 SSE Parser处理编码、拆包和心跳网关层配置关闭缓冲、加长读取超时。这个组合能跑过绝大多数真实业务场景。最后一个小建议做流式接口时尽量从第一天就把日志打清楚。事件类型、时间戳、chunk 长度、解析状态都记录下来。否则断流问题一出现你会在前端没收到和后端没发出去之间反复排查那种状态太痛苦了。
返回列表