☰
LangChain 流式输出与结构化输出实战:SSE 打字机效果与 JSON 解析
2026/10/6 5:03:42 网站建设 项目流程

1. 流式输出的本质:为什么我们需要 SSE

1.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(model="gpt-4o-mini", streaming=True) for chunk in llm.stream("给我讲讲 SSE 的原理"): print(chunk.content, end="", flush=True)

这段代码跑起来就能看到文字一个一个蹦出来。但这里有个坑,很多人第一次用的时候发现还是等全部生成完才输出,原因通常是忘了在初始化时设置streaming=True,或者用错了方法。invoke是同步阻塞的,stream才是流式的,astream是异步流式的。在 FastAPI 这类异步框架里,一定要用astream,否则会阻塞事件循环,导致整个服务卡住。

再往上一个层级,如果你用的是 Chain 或者 Agent,LangChain 也提供了统一的流式接口。Chain 有stream和astream,Agent 在 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 f"data: {json.dumps(payload, ensure_ascii=False)}\n\n"

这里有几个细节值得说。第一,ensure_ascii=False必须加,否则中文会被转义成\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_type="text/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 吐出能直接用的 JSON

3.1 为什么自由文本不够用

流式输出解决了体验问题,但还有一个更根本的问题:大模型默认吐出来的是自然语言,而程序需要的是结构化数据。比如你想让模型从一段用户评论里提取情感倾向、关键词、评分,如果它返回“这段评论看起来是正面的,用户提到了物流快和服务好,大概能打四星”,你没法直接拿这个结果去写数据库。

结构化输出要解决的就是这个问题:约束模型的输出格式,让它返回符合特定 schema 的 JSON。LangChain 在这方面提供了好几层工具,从最简单的PydanticOutputParser到更现代的with_structured_output方法,各有适用场景。

3.2 用 Pydantic 定义输出 schema

Pydantic 是 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的strict=True参数,让底层用更严格的约束。

5.3 中文乱码与编码问题

中文乱码通常出在两个地方。一是后端json.dumps没加ensure_ascii=False,导致中文被转义,虽然前端能解析但看着别扭。二是前端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, version="v2"): 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 f"data: {{\"type\":\"tool_start\",\"name\":\"{event['name']}\"}}\n\n" elif kind == "on_tool_end": yield f"data: {{\"type\":\"tool_end\",\"name\":\"{event['name']}\"}}\n\n"

astream_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。

这些坑的共同点是:文档里不会写,只有真正上手做才会遇到。所以我的建议是,流式应用一定要在接近生产的环境里充分测试,本地跑通不代表线上没问题。

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

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

立即咨询