ARTICLE DETAIL

资讯详情

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

SSE流式输出与LangChain结构化输出实战:FastAPI+Vue智能体开发

SSE流式输出与LangChain结构化输出实战:FastAPI+Vue智能体开发 1. 从“打字机”到“结构化”为什么这套组合拳值得死磕大模型应用做久了你会发现一个很割裂的现象前端想要那种一个字一个字往外蹦的“打字机效果”后端却想要一份规规矩矩、能直接入库的 JSON。这两件事单独做都不难难的是把它们揉在一起——流式输出到一半你根本不知道后面会蹦出什么等全部吐完了再去解析 JSON打字机效果就没了边流边解析又容易在 JSON 还没闭合的时候把解析器搞崩。我最近在做一个基于 FastAPI LangChain 的智能体项目前端是 Vue后端要同时满足两个诉求对话过程必须流式让用户看到实时生成而工具调用、意图识别这些环节又必须拿到结构化的数据不能是一坨自由文本。踩了几轮坑之后我把 SSE 流式原理、LangChain 的结构化输出、以及两者的衔接方案完整跑通了一遍。这篇文章就把这套东西从头到尾拆开讲包括 SSE 的帧格式、LangChain 里with_structured_output的几种实现路径、流式场景下 JSON 增量解析的坑以及前端怎么接。适合谁看如果你正在做 AI 对话类产品或者用 LangChain 搭 Agent 但被“流式 结构化”卡住过这篇应该能帮你省下不少试错时间。哪怕你只是好奇 SSE 到底是怎么工作的前面那部分原理也值得一读——理解了帧格式很多“stream disconnected before completion”之类的报错你一眼就能定位。2. SSE 流式原理先把底层那层窗户纸捅破2.1 SSE 到底是什么和 WebSocket 差在哪SSE全称 Server-Sent Events翻译过来就是“服务器发送事件”。它本质上是 HTTP 协议的一个“长连接变种”客户端发起一个普通的 GET 请求服务器返回的响应头里带上Content-Type: text/event-stream然后这条连接就不关了服务器可以持续往客户端推送数据。很多人第一反应是“那它和 WebSocket 有啥区别”。我一开始也纠结过后来想明白了WebSocket 是双向的全双工客户端和服务器随时都能发消息SSE 是单向的只有服务器往客户端推客户端要发消息得另开一个 HTTP 请求。这个差异决定了它们的适用场景——聊天这种“用户发一句、AI 回一段”的模式其实天然就是“请求-响应”结构用户发消息用普通 POSTAI 回复用 SSE 推流完全够用而且比 WebSocket 简单太多。SSE 最大的好处是它跑在 HTTP 上不需要协议升级不需要额外的握手Nginx、网关、负载均衡这些中间件基本都能直接过。WebSocket 遇到某些代理配置还得专门调SSE 就省心很多。当然代价就是单向但这个代价在大多数 AI 对话场景里根本不算代价。2.2 帧格式data、event、id、retry 四个字段SSE 的协议格式其实特别简单简单到很多人第一次看会觉得“就这”。服务器推送的每一条消息由若干行组成行与行之间用换行符分隔一条消息以两个连续换行符结束。字段就四个data:消息内容这是最核心的可以有多行多行会被拼接event:事件类型客户端可以据此区分不同种类的消息id:消息 ID用于断线重连时告诉服务器“我从哪条之后开始收”retry:重连等待时间单位毫秒一个典型的帧长这样event: message id: 1 data: {content: 你} data: {content: 好}注意最后那个空行它是消息结束的标志。没有这个空行客户端会一直等以为消息还没发完。我见过有人调试 SSE 调了半天收不到数据最后发现就是少了个空行——这种坑真的只有踩过才知道。data:后面如果内容本身包含换行需要拆成多个data:行客户端解析时会用换行符把它们重新拼起来。这个细节在做多行文本流式的时候特别重要比如 AI 输出的代码块里面全是换行如果不按规范拆行客户端收到的内容就会错乱。2.3 浏览器端 EventSource 的自动重连机制浏览器原生提供了EventSource对象来消费 SSE用起来极其简单const es new EventSource(/api/chat/stream); es.onmessage (e) { console.log(e.data); }; es.onerror (err) { console.error(连接出错, err); };EventSource有个很贴心的特性连接断了它会自动重连而且会带上Last-Event-ID请求头把最后收到的id告诉服务器服务器理论上可以从断点续传。这个机制在移动端网络不稳定的场景下非常有用。但这里有个大坑自动重连是“无条件”的。如果服务器返回的是 4xx 错误比如鉴权失败EventSource也会傻乎乎地一直重连导致请求风暴。所以生产环境里我一般会在onerror里判断readyState如果是CLOSED就手动关掉别让它无限重试。另外EventSource不支持自定义请求头这意味着你没法在 header 里塞 token只能把鉴权信息放 URL 参数或者 cookie 里——这也是为什么很多项目干脆不用EventSource改用fetchReadableStream自己解析。2.4 为什么 AI 对话场景偏爱 SSE回到 AI 场景。大模型的生成是逐 token 的一个稍微长点的回答可能要好几秒甚至十几秒。如果等全部生成完再一次性返回用户盯着转圈圈会以为卡死了。SSE 让每个 token 生成出来就立刻推给前端用户看到文字一个个蹦出来心理上就觉得“它在思考、它在干活”体验完全不一样。而且 SSE 的文本特性天然适合传文本。相比 WebSocket 要处理二进制帧、要维护心跳SSE 就是纯文本流调试的时候curl一下就能看到原始数据排查问题不要太方便。我调试流式接口的时候经常直接curl -N -H Accept: text/event-stream http://localhost:8000/api/chat/stream-N参数关掉 curl 的缓冲数据一来就打印配合data:前缀一眼就能看出帧结构对不对。3. LangChain 结构化输出让模型“说人话”也“说格式话”3.1 结构化输出解决的是什么问题自由文本好用但不好“用”。你让模型分析一段用户评论的情感它回你“这条评论整体是正面的用户对物流速度表示满意但对包装有点意见”——人看着挺好程序要提取“情感极性正面、满意度高、问题点包装”就费劲了。结构化输出就是让模型直接返回一个符合预定义 schema 的对象程序拿到就能用。LangChain 在这方面提供了好几层抽象。最直接的是with_structured_output你给它一个 Pydantic 模型或者 JSON Schema它返回的就是一个解析好的对象不用你自己去json.loads。底层它其实做了两件事一是把 schema 塞进 prompt 或者用模型的 function calling / tool use 能力约束输出格式二是把模型返回的文本解析成对象。3.2 Pydantic 模型定义与字段约束用 Pydantic 定义 schema 是最舒服的方式因为类型校验、默认值、字段描述都能写在一起from pydantic import BaseModel, Field from typing import Literal, Optional class SentimentResult(BaseModel): 用户评论的情感分析结果 polarity: Literal[positive, negative, neutral] Field( description情感极性 ) confidence: float Field( ge0.0, le1.0, description置信度0到1之间 ) aspects: list[str] Field( default_factorylist, description评论涉及的具体方面如物流、包装、客服 ) summary: Optional[str] Field( defaultNone, description一句话总结 )这里有几个经验点。第一Field的description不是写给人看的是写给模型看的模型会读这些描述来理解每个字段该填什么所以描述要写清楚别偷懒。第二能用Literal就用Literal把取值范围卡死模型就不容易乱填。第三Optional和default要慎用字段一多模型可能会“偷懒”不填可选字段如果你需要它一定输出就别给默认值。3.3 with_structured_output 的三种底层实现路径with_structured_output看起来是一个方法底层其实有好几种实现取决于你用的模型支持什么第一种是function calling / tool use。像 OpenAI 的模型、Claude、以及很多国产模型都支持工具调用LangChain 会把你的 schema 转换成一个“工具定义”让模型以调用工具的形式输出结构化参数。这条路最稳因为格式约束是模型层面保证的解析成功率最高。第二种是JSON mode。有些模型支持强制输出 JSON你只要在 prompt 里说明要什么格式模型就会返回合法 JSON。这条路比 function calling 弱一点因为 schema 约束靠 prompt模型偶尔会漏字段或者类型不对。第三种是prompt 解析器。模型啥都不支持的时候就只能把 schema 写进 prompt让它按格式输出然后自己解析。这条路最不稳但兼容性最好。LangChain 会根据模型能力自动选你也可以通过method参数强制指定。我的建议是能用 function calling 就用实在不行退到 JSON mode最后才考虑纯 prompt。3.4 结构化输出和流式的天然矛盾问题来了结构化输出要求“完整、合法”流式输出要求“边生成边给”。这俩天生打架。你不可能在 JSON 还没闭合的时候就说它是合法的但用户又不想等到全部生成完才看到东西。这个矛盾在 Agent 场景里更明显。Agent 经常要“先思考、再调工具、再总结”思考过程可以是自由文本流式输出但工具调用的参数必须是结构化的。如果整个流程都等结构化结果那用户就只能干等如果全流式工具调用又没法解析。我试过几种方案最后发现比较靠谱的是“分层”把需要结构化的部分工具调用、意图识别和需要流式的部分最终回答分开处理。中间用 SSE 的不同event类型区分前端根据事件类型决定是渲染到对话气泡里还是拿去触发别的逻辑。这个思路后面会详细展开。4. 流式 结构化把两者缝起来的实战方案4.1 整体架构FastAPI 后端 Vue 前端 SSE 通道先说我这套项目的整体结构。后端 FastAPI 提供两个接口一个是普通的 POST 接口接收用户消息返回一个task_id另一个是 SSE 接口前端拿着task_id去订阅后端把生成过程推过来。为什么拆成两个因为 SSE 是 GET 请求把用户消息塞进 URL 不太优雅而且消息可能很长。拆开之后POST 负责“提交任务”SSE 负责“订阅结果”职责清晰。后端内部用 LangChain 的 Agent 或者 Chain 处理任务处理过程中通过一个队列把事件推给 SSE 接口。队列可以用asyncio.Queue也可以用 Redis 的 pub/sub看你要不要跨进程。单机的话asyncio.Queue就够了。前端 Vue 这边用fetchReadableStream手动解析 SSE而不是用EventSource原因前面说了——要带 token要更灵活地控制重连。4.2 后端用 asyncio.Queue 桥接生成器和 SSE核心代码大概长这样。先定义一个事件类型from dataclasses import dataclass from typing import Literal dataclass class StreamEvent: event: Literal[token, structured, done, error] data: str然后 SSE 接口从队列里取事件按 SSE 格式写出去from fastapi import FastAPI from fastapi.responses import StreamingResponse import asyncio, json app FastAPI() queues: dict[str, asyncio.Queue] {} app.get(/api/chat/stream/{task_id}) async def stream(task_id: str): queue queues.get(task_id) if not queue: return {error: task not found} async def event_generator(): try: while True: evt await asyncio.wait_for(queue.get(), timeout30) if evt.event done: yield fevent: done\ndata: {{}}\n\n break payload json.dumps({data: evt.data}, ensure_asciiFalse) yield fevent: {evt.event}\ndata: {payload}\n\n except asyncio.TimeoutError: yield event: error\ndata: {\msg\:\idle timeout\}\n\n return StreamingResponse( event_generator(), media_typetext/event-stream, headers{ Cache-Control: no-cache, X-Accel-Buffering: no, }, )这里有几个关键点。X-Accel-Buffering: no是给 Nginx 看的告诉它别缓冲这个响应否则 SSE 会被 Nginx 攒着一起发打字机效果就没了。asyncio.wait_for加超时是为了防止连接一直挂着不释放超时了就发个 error 事件让前端知道。ensure_asciiFalse是为了中文不被转义成\uXXXX省点带宽也方便调试。4.3 生成侧LangChain 回调把 token 塞进队列LangChain 的流式靠回调callback实现。你可以自定义一个AsyncCallbackHandler在on_llm_new_token里把 token 塞进队列from langchain_core.callbacks import AsyncCallbackHandler class QueueCallback(AsyncCallbackHandler): def __init__(self, queue: asyncio.Queue): self.queue queue async def on_llm_new_token(self, token: str, **kwargs): await self.queue.put(StreamEvent(eventtoken, datatoken))然后在调用 LLM 的时候把这个 callback 传进去。这样模型每生成一个 token回调就被触发一次token 立刻进队列SSE 接口那边就能马上推给前端。整条链路是异步的不会阻塞。4.4 结构化结果怎么在流里“插队”工具调用或者结构化结果怎么处理我的做法是当 Agent 决定调用工具、并且工具参数已经解析完成时往队列里塞一个structured事件data是序列化后的 JSON。前端收到这个事件就知道“哦模型要调工具了”可以显示一个“正在查询…”的提示而不是把它当成对话内容渲染。async def on_tool_start(self, serialized, input_str, **kwargs): await self.queue.put(StreamEvent( eventstructured, datajson.dumps({tool: serialized.get(name), args: input_str}, ensure_asciiFalse) ))这样整个流里就混着两种事件token是给用户看的文本structured是给程序用的数据。前端按event字段分流互不干扰。这个设计我觉得是整套方案里最舒服的地方——它没有强行把结构化和流式揉成一个东西而是让它们各走各的道在传输层汇合。5. 增量 JSON 解析流式场景下最容易被低估的坑5.1 为什么不能等 JSON 完整了再解析有人会想结构化结果反正要完整才能用那我等流结束再解析不就行了理论上可以但实际场景里往往不行。比如 Agent 调工具工具参数一解析出来就该立刻执行等整个流结束再执行用户等待时间就长了。再比如有些场景需要“边生成边校验”发现模型跑偏了要及时中断等结束就晚了。所以增量解析是有价值的。但增量解析 JSON 有个根本困难JSON 是上下文相关的一个{后面可能跟键、可能跟值一个字符串里的}不算结构结束。你不能简单地按字符扫得维护状态。5.2 手写一个容错的增量 JSON 解析器最土但最可控的办法是自己写一个状态机。核心思路是维护一个括号栈和字符串状态遇到{[入栈遇到}]出栈遇到引号切换“是否在字符串内”的状态只有在字符串外遇到的括号才算结构括号。当栈为空且不在字符串内时说明一个完整的 JSON 对象结束了。class IncrementalJSONParser: def __init__(self): self.buffer self.depth 0 self.in_string False self.escape False def feed(self, chunk: str): results [] for ch in chunk: self.buffer ch if self.escape: self.escape False continue if ch \\ and self.in_string: self.escape True continue if ch : self.in_string not self.in_string continue if self.in_string: continue if ch in {[: self.depth 1 elif ch in }]: self.depth - 1 if self.depth 0: results.append(self.buffer) self.buffer return results这个解析器不校验 JSON 合法性只负责“切出完整的 JSON 片段”。切出来之后再交给json.loads去校验如果非法就丢弃或者报错。这样职责分离逻辑清晰。5.3 用 json-repair 之类的库兜底手写解析器能处理大部分情况但模型有时候会输出一些“接近合法但不完全合法”的 JSON比如尾逗号、单引号、缺引号的键。这时候可以用json-repair这类库兜底它会尝试修复常见的 JSON 错误再解析。我的策略是先用标准json.loads失败了再用 repair 库还失败就记录日志、跳过这条。别指望 100% 成功留好降级路径。5.4 流式解析的边界情况清单实际跑下来我遇到过这些边界情况列出来给大家避坑情况表现处理方式字符串内含转义引号状态机误判字符串结束处理\转义字符串内含{}括号栈误计数字符串内不计数多个 JSON 对象连在一起一次 feed 切出多个返回列表JSON 跨 chunk 断开单个 chunk 不完整用 buffer 累积模型输出 markdown 代码块包裹前后有 json先剥离代码块标记尾逗号json.loads 报错repair 库兜底提示模型输出 JSON 时经常喜欢用 json 包起来解析前一定要先做一次“去代码块”处理否则你的解析器永远在等一个不存在的结束符。6. 前端 Vue 侧fetch 流式读取与打字机渲染6.1 为什么不用 EventSource前面提过EventSource不支持自定义 headertoken 只能塞 URL 或 cookie。URL 塞 token 不安全cookie 在跨域场景下又麻烦。而且EventSource的自动重连在出错时会变成负担。所以我用fetchReadableStream自己读流、自己解析 SSE 帧。6.2 手动解析 SSE 帧的完整代码async function streamChat(taskId, token, handlers) { const resp await fetch(/api/chat/stream/${taskId}, { headers: { Authorization: Bearer ${token} }, }); const reader resp.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 }); let idx; while ((idx buffer.indexOf(\n\n)) ! -1) { const raw buffer.slice(0, idx); buffer buffer.slice(idx 2); const evt parseSSEFrame(raw); if (evt) handlers[evt.event]?.(evt.data); } } } function parseSSEFrame(raw) { let event message; const dataLines []; for (const line of raw.split(\n)) { if (line.startsWith(event:)) event line.slice(6).trim(); else if (line.startsWith(data:)) dataLines.push(line.slice(5).trim()); } if (!dataLines.length) return null; return { event, data: dataLines.join(\n) }; }这里的关键是decoder.decode(value, { stream: true })stream: true保证多字节字符比如中文跨 chunk 时不会被截断成乱码。这个参数不加中文流式输出会出现“半个字”的乱码非常隐蔽。6.3 打字机效果的实现细节打字机效果本身不难难的是“平滑”。如果每个 token 到了就直接text token遇到模型一次吐好几个 token 的时候文字会“跳”一下不够顺。我的做法是维护一个待显示队列用一个定时器每隔 20-30ms 从队列里取一个字符渲染这样无论后端来得多快多慢前端都是匀速出字。const displayQueue []; let timer null; function pushChar(ch) { displayQueue.push(ch); if (!timer) timer setInterval(flush, 25); } function flush() { if (!displayQueue.length) { clearInterval(timer); timer null; return; } displayText.value displayQueue.shift(); }这个方案还有个好处如果用户滚动页面或者切换标签渲染不会因为后端推得太快而卡顿。代价是显示会稍微滞后于实际生成但用户感知不到。6.4 结构化事件在前端怎么用前端收到structured事件时不要往对话气泡里塞而是单独处理。比如工具调用事件可以显示一个“正在查询天气…”的卡片意图识别事件可以用来切换 UI 状态。这样用户看到的是“AI 在思考、在调工具、在回答”的完整过程而不是一堆看不懂的 JSON。7. 常见问题与排查技巧实录7.1 “stream disconnected before completion: idle timeout” 怎么破这个报错我遇到好几次本质是连接空闲太久被中间层掐了。可能的原因有几个一是后端生成太慢超过网关的空闲超时二是 Nginx 的proxy_read_timeout默认 60 秒超了就断三是后端自己的asyncio.wait_for超时设太短。排查顺序先看后端日志有没有正常推数据再看 Nginx 配置最后看客户端有没有及时消费。解决上Nginx 加proxy_read_timeout 300s;和proxy_buffering off;后端定期发心跳比如每 15 秒发一个: heartbeat\n\n注释帧客户端确保读取循环不阻塞。7.2 中文乱码与 chunk 截断中文乱码几乎都是TextDecoder没加stream: true或者后端json.dumps没加ensure_asciiFalse导致前端拿到一堆\uXXXX。前者是解码问题后者是编码问题症状不同但都好排查。后端统一ensure_asciiFalse前端统一stream: true基本就不会有乱码。7.3 结构化输出解析失败的降级策略模型不是每次都能输出合法 JSON尤其是用 prompt 方式约束的时候。我的降级策略分三层第一层标准解析第二层 repair 库修复第三层如果还失败就把原始文本作为普通消息推给前端同时记录日志。千万别因为解析失败就让整个请求挂掉用户体验会崩。7.4 常见问题速查表现象可能原因排查方向前端收不到任何数据响应头不对 / 被缓冲检查 Content-Type 和 X-Accel-Buffering数据攒着一起到Nginx 缓冲关 proxy_buffering中文乱码解码未用 stream 模式TextDecoder stream: trueJSON 解析报错模型输出不合法加 repair 兜底连接频繁断开空闲超时加心跳 调超时打字机卡顿渲染未节流加显示队列匀速渲染提示调试 SSE 时先在浏览器 Network 面板看 EventStream 标签页能看到每一条原始帧比 console.log 直观得多。8. 一些踩坑之后的个人体会这套方案跑下来我最大的感受是流式和结构化不是非此即彼而是应该分层处理。把“给用户看的”和“给程序用的”用不同的事件类型分开传输层用 SSE 统一承载前端按类型分流整个架构就清爽了。硬要把两者揉成一个数据流只会让解析逻辑越来越复杂。另外SSE 这东西看着简单但生产环境的坑基本都在中间件和超时上。本地跑得好好的一上 Nginx 就出问题十有八九是缓冲和超时。我现在的习惯是任何 SSE 接口上线前先在 Nginx 后面压一遍确认流是“实时”出来的而不是攒一批发一批。最后分享一个小技巧如果你不确定后端推的帧格式对不对直接curl -N打出来看比任何调试工具都快。看到data:一行行往外冒心里就踏实了。
返回列表