ARTICLE DETAIL

资讯详情

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

LangChain 流式输出与结构化输出实战:SSE 打字机效果与 JSON 解析

LangChain 流式输出与结构化输出实战:SSE 打字机效果与 JSON 解析 1. 流式输出的本质为什么我们需要 SSE1.1 从“等一锅饭”到“边炒边上桌”的思维转变做过大模型应用的人都有一个共同体会用户等一个完整回答的耐心远比我们想象的要短。早期做对话产品时我试过让前端一直转圈等后端把整段回答生成完再一次性返回结果就是超过三秒用户就开始怀疑是不是卡死了超过五秒直接关页面走人。这个体验问题不是靠优化模型推理速度能解决的因为大模型逐 token 生成的物理特性摆在那里你不可能让一个需要生成五百字的回答在一瞬间全部蹦出来。流式输出解决的正是这个“等待焦虑”问题。它的核心思路很简单模型每生成一小段内容就立刻推给前端渲染而不是攒齐了再发。用户看到文字一个一个蹦出来哪怕总时长没变主观感受上也会觉得“它在思考、它在回应”这就是所谓的打字机效果。而实现这种效果最成熟、最通用的底层协议就是 SSE全称 Server-Sent Events。SSE 本质上是一个基于 HTTP 长连接的单项推送协议。客户端发起一个普通 HTTP 请求服务端在响应头里声明Content-Type: text/event-stream然后保持这个连接不关闭持续往客户端写数据。每一条数据以data:开头以两个换行符结束格式非常朴素。浏览器端有原生的EventSourceAPI 可以直接消费但实际项目里我们更多用fetch配合ReadableStream来手动解析因为EventSource只支持 GET 请求没法携带复杂的请求体这在需要传对话历史的场景下是硬伤。1.2 SSE 与 WebSocket 的选型逻辑很多人一提到实时推送就想到 WebSocket觉得双向通信肯定比单向强。但在大模型对话这个场景里这个想法是错的。WebSocket 建立的是全双工连接协议更重需要额外的握手升级过程服务端维护连接的成本也更高。而大模型对话的数据流向是典型的“客户端发一次请求服务端持续推多次响应”本质上是单向的。用 WebSocket 就像为了送一趟快递专门修了一条双向高速公路杀鸡用牛刀。SSE 的优势在于它复用了 HTTP 协议栈不需要额外的协议升级穿透代理和网关的能力更强断线重连机制也是浏览器原生支持的。当然它也有短板比如默认不支持二进制传输、连接数在 HTTP/1.1 下有限制但这些在大模型文本对话场景里都不是问题。我个人的经验是纯文本流式推送用 SSE需要双向实时交互比如协同编辑、游戏才上 WebSocket不要为了技术时髦而过度设计。1.3 一次完整的 SSE 数据流长什么样在动手写代码之前先把 SSE 的数据格式彻底搞清楚后面解析才不会踩坑。服务端推给客户端的数据在网络上实际传输的样子是这样的data: {type:token,content:你} data: {type:token,content:好} data: {type:done,finish_reason:stop}注意几个关键细节。第一每条消息以data:开头冒号后面有一个空格这个空格是规范的一部分解析时要去掉。第二每条消息以两个换行符\n\n结尾这是消息之间的分隔符。第三如果一条消息内容很长可以分成多个data:行客户端会把它们用换行符拼接起来。第四服务端可以发送event:字段来指定事件类型发送id:字段来标记消息序号发送retry:字段来指定重连间隔。实际项目中OpenAI 兼容的接口返回格式通常是每个 chunk 一个 JSON里面包含choices[0].delta.content这样的结构。而 LangChain 的流式输出会把这些 chunk 统一封装成AIMessageChunk对象。理解这个底层格式是后面所有解析工作的基础。2. LangChain 流式输出的接入与封装2.1 LangChain 的流式接口到底怎么用LangChain 从 0.1 版本开始对流式输出的支持已经相当完善了。最基础的用法是调用模型的stream方法它会返回一个生成器每次 yield 一个AIMessageChunk。我拿 OpenAI 兼容的模型举例代码大概长这样from langchain_openai import ChatOpenAI llm ChatOpenAI(modelgpt-4o-mini, streamingTrue) for chunk in llm.stream(给我讲讲 SSE 的原理): print(chunk.content, end, flushTrue)这段代码跑起来就能看到文字一个一个蹦出来。但这里有个坑很多人第一次用的时候发现还是等全部生成完才输出原因通常是忘了在初始化时设置streamingTrue或者用错了方法。invoke是同步阻塞的stream才是流式的astream是异步流式的。在 FastAPI 这类异步框架里一定要用astream否则会阻塞事件循环导致整个服务卡住。再往上一个层级如果你用的是 Chain 或者 AgentLangChain 也提供了统一的流式接口。Chain 有stream和astreamAgent 在 LangGraph 体系下也有对应的流式方法。但 Agent 的流式输出比单纯 LLM 复杂得多因为它中间可能涉及工具调用、多轮推理流出来的不只是最终回答的 token还有中间步骤的事件。这个后面单独讲。2.2 把 LangChain 的 chunk 转成 SSE 格式LangChain 的AIMessageChunk对象不能直接扔给前端必须转成 SSE 格式的字符串。我封装过一个通用的转换函数核心逻辑就是把 chunk 的内容包装成 JSON再套上data:前缀和双换行后缀import json def chunk_to_sse(chunk): payload { type: token, content: chunk.content, finish_reason: chunk.response_metadata.get(finish_reason) } return fdata: {json.dumps(payload, ensure_asciiFalse)}\n\n这里有几个细节值得说。第一ensure_asciiFalse必须加否则中文会被转义成\uXXXX的形式虽然前端也能解析但传输体积会变大调试时看着也难受。第二finish_reason要透传出去前端需要知道什么时候流结束了才能关闭连接、停止 loading 动画。第三如果 chunk 的 content 是空字符串比如第一个 chunk 通常只有 role 信息可以选择跳过不发送减少无效传输。在 FastAPI 里返回 SSE 响应用的是StreamingResponse配合一个异步生成器from fastapi import FastAPI from fastapi.responses import StreamingResponse app FastAPI() async def event_generator(prompt: str): async for chunk in llm.astream(prompt): if chunk.content: yield chunk_to_sse(chunk) yield data: {\type\:\done\}\n\n app.get(/chat) async def chat(prompt: str): return StreamingResponse( event_generator(prompt), media_typetext/event-stream, headers{ Cache-Control: no-cache, Connection: keep-alive, X-Accel-Buffering: no } )X-Accel-Buffering: no这个头非常关键如果你前面挂了 Nginx不加这个头 Nginx 会默认缓冲响应导致流式效果失效用户还是等全部生成完才看到内容。这个坑我踩过不止一次排查了半天才发现是网关层在缓冲。2.3 封装一个可复用的 SSE 流式接口调用逻辑后端封装好了前端消费也不能马虎。浏览器原生EventSource只支持 GET传不了复杂的请求体所以实际项目里我推荐用fetch加ReadableStream手动解析。下面是我常用的一个封装async function streamChat(prompt, onToken, onDone) { const response await fetch(/chat, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ prompt }) }); const reader response.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 }); const lines buffer.split(\n\n); buffer lines.pop(); for (const line of lines) { if (!line.startsWith(data: )) continue; const data JSON.parse(line.slice(6)); if (data.type token) onToken(data.content); if (data.type done) onDone(); } } }这段代码的核心在于buffer的处理。网络传输是分片的一个 SSE 消息可能被拆到两个 TCP 包里所以不能假设每次read()拿到的都是完整消息。正确做法是把已接收的内容拼到 buffer 里按\n\n切分最后一段可能不完整留在 buffer 里等下次拼接。这个细节如果处理不好会出现 JSON 解析报错而且报错是偶发的特别难排查。3. 结构化输出让 AI 吐出能直接用的 JSON3.1 为什么自由文本不够用流式输出解决了体验问题但还有一个更根本的问题大模型默认吐出来的是自然语言而程序需要的是结构化数据。比如你想让模型从一段用户评论里提取情感倾向、关键词、评分如果它返回“这段评论看起来是正面的用户提到了物流快和服务好大概能打四星”你没法直接拿这个结果去写数据库。结构化输出要解决的就是这个问题约束模型的输出格式让它返回符合特定 schema 的 JSON。LangChain 在这方面提供了好几层工具从最简单的PydanticOutputParser到更现代的with_structured_output方法各有适用场景。3.2 用 Pydantic 定义输出 schemaPydantic 是 Python 生态里做数据校验的事实标准LangChain 的结构化输出深度集成了它。定义一个 schema 非常直观from pydantic import BaseModel, Field from typing import List class ReviewAnalysis(BaseModel): sentiment: str Field(description情感倾向只能是 positive/negative/neutral) score: int Field(description评分1 到 5 的整数) keywords: List[str] Field(description评论中提到的关键词列表) summary: str Field(description一句话总结)每个字段的description非常重要它不是给人看的注释而是会作为提示词的一部分发给模型告诉模型这个字段该填什么。description 写得越清楚模型填错格式的概率越低。我见过很多人 schema 定义得很随意description 空着不写然后抱怨模型输出不稳定其实问题出在自己这边。3.3 with_structured_output 的实战用法LangChain 现在主推的是with_structured_output方法它比老的 Parser 方案更简洁而且底层会根据模型能力自动选择最佳实现方式。对于支持 function calling 的模型它会用工具调用的方式约束输出对于不支持的模型它会退化成提示词约束加解析。structured_llm llm.with_structured_output(ReviewAnalysis) result structured_llm.invoke(这个产品太棒了物流超快客服也很耐心五星好评) print(result.sentiment) # positive print(result.score) # 5返回的result直接就是ReviewAnalysis类型的对象字段访问用点号IDE 有自动补全类型检查也能过。这比手动json.loads再取字段舒服太多了。但这里有个关键限制with_structured_output默认是非流式的。因为结构化输出需要等模型把整个 JSON 生成完才能解析中途的片段是不完整的 JSON没法解析。这就产生了一个矛盾既要结构化又要流式打字机效果怎么办3.4 结构化输出与流式的矛盾及折中方案这个矛盾的本质是JSON 的语法要求完整性而流式输出的特点是渐进性。一个 JSON 对象在生成到一半的时候{sentiment: pos这样的片段是没法解析的。我实践下来有三种折中方案。第一种是“先流式后结构化”让模型先用自然语言流式回答回答完再单独调一次结构化接口提取数据。缺点是调了两次模型成本和延迟都翻倍。第二种是“流式 JSON 增量解析”用一个能容忍不完整 JSON 的解析器边流边尝试解析能解析出多少算多少。这种方案技术含量高但体验最好。第三种是“字段级流式”把结构化输出拆成多个字段每个字段单独流式生成前端按字段逐个渲染。我目前项目里用得最多的是第二种配合一个叫partial-json-parser的库它能解析不完整的 JSON 片段返回已经完整的部分。比如{sentiment: positive, score:这样的片段它能解析出{sentiment: positive}。前端拿到部分数据就能先渲染等完整了再补全。4. 打字机效果的前端实现细节4.1 逐字渲染还是逐块渲染后端推过来的 chunk 粒度是不固定的有时候一个 chunk 是一个字有时候是一整句。如果直接按 chunk 渲染会出现“有时候一个字一个字蹦有时候一整句突然出现”的不均匀感。要做出丝滑的打字机效果前端需要做一层缓冲和匀速输出。我的做法是维护一个待渲染队列后端每来一个 chunk 就入队然后用requestAnimationFrame或者setInterval以固定速度从队列里取字符渲染。这样无论后端推得快还是慢视觉上都是匀速的。速度一般控制在每帧 1 到 3 个字符太快了没有打字感太慢了用户着急。let queue ; let rendering false; function enqueue(text) { queue text; if (!rendering) renderLoop(); } function renderLoop() { rendering true; if (queue.length 0) { rendering false; return; } const char queue[0]; queue queue.slice(1); outputElement.textContent char; setTimeout(renderLoop, 30); }这个 30 毫秒的间隔是调出来的经验值对应大约每秒 33 个字符接近正常人阅读速度看起来比较自然。4.2 自动滚动与用户打断的处理打字机效果还有一个容易被忽略的细节自动滚动。内容越来越多容器要自动滚到底部否则用户得手动往下拉。但这里有个坑如果用户主动往上滚动去看之前的内容你还强制滚到底部用户会很烦躁。正确做法是判断当前滚动位置只有当用户已经在底部附近时才自动滚动。function autoScroll() { const el document.getElementById(chat-container); const isAtBottom el.scrollHeight - el.scrollTop - el.clientHeight 50; if (isAtBottom) { el.scrollTop el.scrollHeight; } }这个 50 像素的阈值也是经验值太小了稍微滚一点就触发太大了用户滚上去了还会被拉下来。4.3 流中断与异常状态的 UI 反馈流式输出最怕的就是中途断了。网络抖动、服务端超时、模型报错都可能导致流中断。这时候前端不能一直转圈等必须给用户明确的反馈。我在实际项目里遇到过stream disconnected before completion: idle timeout waiting for sse这个报错原因是服务端超过一定时间没有推送任何数据网关判定连接空闲就掐断了。解决办法有两个一是服务端定期发送心跳注释以:开头的行客户端会忽略保持连接活跃二是前端设置超时检测超过一定时间没收到数据就主动断开并提示用户重试。async def event_generator(prompt: str): last_heartbeat time.time() async for chunk in llm.astream(prompt): if chunk.content: yield chunk_to_sse(chunk) if time.time() - last_heartbeat 15: yield : heartbeat\n\n last_heartbeat time.time() yield data: {\type\:\done\}\n\n心跳间隔设 15 秒比较稳妥大部分网关的空闲超时都在 30 秒以上留一半余量。5. 常见问题排查与避坑实录5.1 流式失效的排查思路流式失效是最常见的问题表现就是用户等半天然后所有内容一次性出现。排查要按链路逐段确认。先确认模型层是不是真的在流式可以在后端加日志看astream是不是逐个 yield 的。如果模型层没问题再确认 FastAPI 的StreamingResponse有没有被中间件缓冲。最后确认网关层Nginx 需要关proxy_buffering加X-Accel-Buffering: no头。下面这张表是我整理的排查清单按顺序过一遍基本能定位问题排查环节检查项常见问题模型层是否用 stream/astream误用 invoke 导致阻塞框架层StreamingResponse 配置media_type 写错中间件是否有缓冲中间件GZip 中间件会缓冲网关层Nginx 缓冲配置proxy_buffering 默认开前端层是否正确解析流按 chunk 而非按消息解析5.2 JSON 解析失败的典型场景结构化输出解析失败十有八九是模型输出的 JSON 不合法。常见的有多了 markdown 代码块标记json 包裹、字段类型不对该是整数给了字符串、缺少必填字段、JSON 后面跟了多余的解释文字。LangChain 的解析器对 markdown 代码块标记有一定容错但类型错误和缺字段是没法自动修复的。我的经验是在 schema 的 description 里把约束写死比如“只返回 JSON不要有任何其他文字”、“score 必须是 1 到 5 的整数不要加引号”。另外可以用with_structured_output的strictTrue参数让底层用更严格的约束。5.3 中文乱码与编码问题中文乱码通常出在两个地方。一是后端json.dumps没加ensure_asciiFalse导致中文被转义虽然前端能解析但看着别扭。二是前端TextDecoder没指定utf-8或者解码时没加{ stream: true }参数导致多字节字符被截断。{ stream: true }这个参数特别重要。UTF-8 编码的中文一个字占三个字节如果网络分片正好切在一个字的中间不加这个参数就会解码出乱码。加了之后TextDecoder会把不完整的字节序列缓存起来等下一个分片到了再一起解码。5.4 并发场景下的连接管理多个用户同时对话时每个用户一个 SSE 连接服务端要维护大量长连接。这里要注意几个点。一是连接要有超时机制用户关了页面但连接没断的情况很常见需要服务端定期清理。二是要限制单用户的最大并发连接数防止恶意占用。三是如果用异步框架确保生成器里没有阻塞操作否则会拖垮整个事件循环。我在一个项目里遇到过连接泄漏原因是用户关闭页面后后端的生成器还在跑因为模型还在生成。解决办法是在生成器里检测客户端断开FastAPI 里可以通过request.is_disconnected()来判断断开就停止生成释放资源。6. 从单轮到多轮Agent 场景下的流式挑战6.1 Agent 流式输出的特殊性前面讲的都是单轮对话的流式Agent 场景要复杂得多。一个 Agent 处理用户请求时可能先思考、再调用工具、拿到结果再思考、最后才给出回答。这个过程中用户希望看到的不只是最终回答还有中间的推理步骤和工具调用状态这样才有“AI 在干活”的感知。LangGraph 体系下Agent 的流式输出有几种模式。values模式每次输出完整状态updates模式只输出变化的部分messages模式专门输出消息 token。实际项目里我通常用messages模式拿 token 流同时用updates模式拿工具调用事件两者结合给用户完整的反馈。6.2 工具调用事件的透传工具调用是 Agent 的特色也是流式处理的难点。当 Agent 决定调用某个工具时流里会出现一个带有tool_calls的 chunk这时候前端应该显示“正在调用 XX 工具”的提示而不是继续渲染文字。async for event in agent.astream_events(input, versionv2): kind event[event] if kind on_chat_model_stream: chunk event[data][chunk] if chunk.content: yield chunk_to_sse(chunk) elif kind on_tool_start: yield fdata: {{\type\:\tool_start\,\name\:\{event[name]}\}}\n\n elif kind on_tool_end: yield fdata: {{\type\:\tool_end\,\name\:\{event[name]}\}}\n\nastream_events是 LangChain 提供的统一事件流接口能拿到模型流、工具开始、工具结束等各种事件。用这个接口就不用自己去猜 chunk 的类型了事件类型是明确的。6.3 多轮对话历史的流式处理多轮对话时每次请求都要把历史消息带上。历史消息可能很长如果每次都全量传输请求体会很大。我的做法是后端维护会话状态前端只传一个 session_id后端根据 id 取出历史。这样请求体小也避免了历史被篡改的风险。但会话状态存哪里是个问题。存内存最简单但服务重启就丢了多实例部署也不共享。存 Redis 是更稳妥的方案设置合理的过期时间比如 30 分钟无活动就清理。如果对话很重要不能丢那就得落库但落库会增加延迟需要权衡。7. 性能优化与生产环境注意事项7.1 减少首字延迟首字延迟是流式体验的关键指标用户从点击发送到看到第一个字的时间超过一秒就会觉得慢。影响首字延迟的因素有几个模型本身的推理启动时间、网络往返、后端处理逻辑。优化手段上模型层可以选更快的模型或者用推理加速服务。网络层可以把服务部署在离用户近的区域。后端层要确保在调用模型之前没有耗时操作比如查数据库、做复杂计算这些都应该提前做好或者异步做。我见过有人在生成器里先查一次用户信息再调模型白白增加了几百毫秒延迟。7.2 背压与流量控制流式输出是服务端推、客户端收如果客户端消费慢服务端推得快数据就会在缓冲区堆积。Python 的异步生成器天然有背压机制yield会等待消费者取走才继续所以一般不用担心。但如果中间加了队列做缓冲就要注意队列长度限制防止内存暴涨。7.3 日志与可观测性生产环境一定要有完善的日志。每次请求记录请求 id、用户 id、prompt 长度、首字延迟、总时长、token 数、是否异常中断。这些数据是排查问题和优化性能的基础。我习惯在 SSE 流里也带上请求 id前端报错时可以把 id 给到后端直接定位到具体那次请求的日志。另外要监控异常中断率如果这个指标突然升高说明可能有网络问题或者服务端问题。中断率超过 5% 就值得警惕了。8. 我踩过的几个印象深刻的坑第一个坑是 Nginx 缓冲。本地开发一切正常部署到测试环境流式就失效了排查了一下午才发现是 Nginx 默认开启了proxy_buffering。这个坑的教训是流式应用部署时网关层的配置一定要单独确认不能假设默认配置就是对的。第二个坑是TextDecoder的stream参数。前端偶尔出现乱码特别是中文概率大概百分之几。查了很久才定位到是解码时没加{ stream: true }导致多字节字符被网络分片截断。这个 bug 的隐蔽性在于它是概率性的本地测试很难复现。第三个坑是结构化输出的流式矛盾。一开始我想当然地以为with_structured_output也能流式结果发现它内部是等完整 JSON 才返回的。后来改用增量 JSON 解析才解决。这个坑让我明白不是所有 LangChain 的方法都支持流式用之前要确认清楚。第四个坑是连接泄漏。用户关闭页面后后端生成器还在跑因为模型还在生成生成器不知道客户端已经走了。时间一长大量僵尸连接占满资源。解决办法是在生成器循环里定期检查request.is_disconnected()断开就break。这些坑的共同点是文档里不会写只有真正上手做才会遇到。所以我的建议是流式应用一定要在接近生产的环境里充分测试本地跑通不代表线上没问题。
返回列表