1. 为什么流式输出和结构化输出总是打架
做过大模型应用的人多半遇到过这个场景:前端用 SSE 接收流式内容,一个字一个字往外蹦,体验很顺滑;可一旦业务需要拿到结构化的 JSON,比如要提取字段、要触发下游工具调用,流式就变成了麻烦。你没法在流还没结束时解析 JSON,因为 JSON 是残缺的;可如果等流结束再解析,那流式的意义又没了一半。
这个矛盾在 LangChain 里体现得特别明显。LangChain 的 OutputParser 体系天生是为"拿到完整输出再解析"设计的,而 SSE 流式是"边生成边消费"。两者要配合,就得在中间做文章。再加上 ToolCall 这一层——模型不只是输出文本,还要决定调用哪个工具、传什么参数——结构化输出的要求就更高了,因为工具调用的参数必须是合法 JSON,错一个括号整个调用就废了。
我最近在一个基于 FastAPI + LangChain 的项目里把这三块彻底捋了一遍:SSE 流式接口怎么封装、OutputParser 怎么选、ToolCall 怎么和流式共存。踩的坑不少,有些是 LangChain 版本差异导致的,有些是 SSE 协议本身的特性决定的。这篇就把整个实战过程拆开讲,从最基础的流式封装讲到 ToolCall 的流式解析,尽量把每个"为什么这么设计"说清楚。
适合的读者是已经用过 LangChain 基础功能、想把它接到真实 Web 服务里的人。如果你还在纠结 LangChain 怎么入门,建议先把 Chain 和 Prompt 跑通再来看这篇,不然有些细节会显得突兀。
2. SSE 流式接口的封装逻辑与常见断流问题
2.1 SSE 到底是个什么东西,为什么大模型场景偏爱它
SSE 全称 Server-Sent Events,是 HTTP 协议下的一种单向推送机制。服务端保持一个长连接,持续往客户端写data: xxx\n\n格式的文本块,客户端用 EventSource 或者 fetch 的 ReadableStream 来读。它和 WebSocket 的区别在于:SSE 是单向的(服务端到客户端),基于纯 HTTP,不需要额外的协议升级,浏览器原生支持自动重连。
大模型场景偏爱 SSE 的原因很实际。第一,大模型的输出本身就是单向的,用户发一次请求,模型吐一串 token,不需要双向通信。第二,SSE 基于 HTTP,能直接穿过大多数网关和负载均衡,部署成本低。第三,SSE 的文本格式天然适合传 token,每个 token 包一个data:就行。
但 SSE 有个容易被忽略的点:它是纯文本协议,没有二进制帧的概念。这意味着你不能在里面塞复杂的二进制结构,所有东西都得序列化成字符串。这一点在后面讲结构化输出时会变成一个关键约束。
2.2 用 FastAPI 封装 SSE 接口的最小可用骨架
先给一个能跑的最小骨架。FastAPI 里返回 SSE 用StreamingResponse,媒体类型设成text/event-stream。
from fastapi import FastAPI from fastapi.responses import StreamingResponse import asyncio app = FastAPI() async def event_generator(): for i in range(5): yield f"data: chunk-{i}\n\n" await asyncio.sleep(0.5) yield "data: [DONE]\n\n" @app.get("/stream") async def stream(): return StreamingResponse( event_generator(), media_type="text/event-stream", headers={ "Cache-Control": "no-cache", "Connection": "keep-alive", "X-Accel-Buffering": "no", }, )这里有几个细节值得说。X-Accel-Buffering: no这个头是给 Nginx 看的,告诉它不要缓冲这个响应。如果你部署在 Nginx 后面又没加这个头,会发现流式变成了"攒一批再发",前端看起来就是卡顿的。Cache-Control: no-cache防止中间层缓存。Connection: keep-alive保持长连接。
data:后面跟内容,然后必须有两个换行\n\n,这是 SSE 协议的帧分隔符。少一个换行,客户端就认为这一帧还没结束,会一直等。这个坑我见过不止一次,尤其是拼接字符串时手滑。
2.3 "stream disconnected before completion" 的根因排查
热词里有个很典型的报错:stream disconnected before completion: idle timeout waiting for sse。这个错误的字面意思是:SSE 连接在完成前断开了,因为等待 SSE 时空闲超时了。
根因通常有三个方向。第一是服务端或中间层设了空闲超时。Nginx 默认的proxy_read_timeout是 60 秒,如果模型生成慢,60 秒内没有任何数据写出,连接就被掐了。解决办法是在 Nginx 配置里调大这个值,或者更优雅的做法是定期发送心跳包。
async def event_generator_with_heartbeat(chain, inputs): queue = asyncio.Queue() async def produce(): async for chunk in chain.astream(inputs): await queue.put(chunk) await queue.put(None) producer = asyncio.create_task(produce()) while True: try: chunk = await asyncio.wait_for(queue.get(), timeout=15) except asyncio.TimeoutError: yield ": heartbeat\n\n" continue if chunk is None: break yield f"data: {chunk}\n\n" await producer心跳包用:开头,这是 SSE 的注释语法,客户端会忽略它,但它能让连接保持活跃,避免空闲超时。15 秒发一次是个比较稳妥的值,比大多数默认超时都短。
第二个方向是客户端主动断开。用户切了页面、关了标签,EventSource 会断开,服务端如果还在往一个已断开的连接写数据,就会抛异常。这个要在生成器里捕获asyncio.CancelledError或者ConnectionResetError,做好清理。
第三个方向是代理层。有些云厂商的负载均衡对长连接有硬性时长限制,比如 5 分钟强制断开。这种只能靠客户端重连 + 服务端记录断点来续传,或者干脆改成轮询。选型时要先确认部署环境的限制。
2.4 前端消费 SSE 的两种姿势与踩坑
前端消费 SSE 有两种主流方式。一种是用原生EventSource,简单但有局限:它只支持 GET 请求,不能自定义请求头,没法传 Authorization。另一种是用fetch+ReadableStream,灵活但需要自己解析帧。
async function consumeSSE(url, body) { const resp = await fetch(url, { method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify(body), }); const reader = resp.body.getReader(); const decoder = new TextDecoder(); let buffer = ""; while (true) { const { done, value } = await reader.read(); if (done) break; buffer += decoder.decode(value, { stream: true }); const parts = buffer.split("\n\n"); buffer = parts.pop(); for (const part of parts) { if (part.startsWith("data: ")) { const data = part.slice(6); if (data === "[DONE]") return; handleChunk(data); } } } }这里的关键是buffer的处理。网络传输是分片的,一个 SSE 帧可能被切成两半到达,所以不能假设每次read()拿到的就是完整帧。正确做法是累积到 buffer,按\n\n切分,最后一段可能不完整,留在 buffer 里等下一片。这个逻辑写错了,就会出现 JSON 解析到一半报错的情况,而且报错是间歇性的,很难复现。
decoder.decode(value, { stream: true })里的stream: true也很重要。UTF-8 一个汉字占 3 字节,如果分片正好切在汉字中间,不加这个参数会解码出乱码。加了之后 TextDecoder 会保留不完整的字节序列,等下一片拼上再解码。
3. LangChain 三大 OutputParser 的选型与实战差异
3.1 PydanticOutputParser:强类型场景的首选
LangChain 的 OutputParser 家族里,PydanticOutputParser 是最"重"的一个,也是结构化要求最高时的首选。它的工作方式是:你定义一个 Pydantic 模型,它自动生成一段格式说明塞进 prompt,模型按这个格式输出,解析器再把文本转回 Pydantic 对象。
from langchain_core.output_parsers import PydanticOutputParser from langchain_core.prompts import ChatPromptTemplate from pydantic import BaseModel, Field class PersonInfo(BaseModel): name: str = Field(description="人物姓名") age: int = Field(description="年龄,整数") skills: list[str] = Field(description="技能列表") parser = PydanticOutputParser(pydantic_object=PersonInfo) prompt = ChatPromptTemplate.from_messages([ ("system", "从文本中提取人物信息。\n{format_instructions}"), ("human", "{input}"), ]).partial(format_instructions=parser.get_format_instructions()) chain = prompt | llm | parser result = chain.invoke({"input": "张三今年 28 岁,会 Python 和 Go。"})get_format_instructions()会生成一段类似"请输出符合以下 JSON Schema 的内容"的说明。这段说明的质量直接决定解析成功率。Pydantic 的Field(description=...)会变成 JSON Schema 里的 description,模型靠这个理解每个字段的含义。所以字段描述写得越清楚,解析越稳。
PydanticOutputParser 的优势是类型安全。解析出来的对象有完整的类型校验,age是 int 不是 str,skills是 list 不是字符串。下游代码可以直接用属性访问,不用手动转换。缺点是它对格式要求严格,模型稍微跑偏就解析失败,而且失败时抛的是OutputParserException,需要自己捕获处理。
3.2 JsonOutputParser:轻量灵活但需要自己兜底
JsonOutputParser 比 Pydantic 轻,它不要求你定义模型,直接输出 dict。适合结构不固定、或者你懒得定义模型的场景。
from langchain_core.output_parsers import JsonOutputParser parser = JsonOutputParser() prompt = ChatPromptTemplate.from_messages([ ("system", "输出 JSON 格式。\n{format_instructions}"), ("human", "{input}"), ]).partial(format_instructions=parser.get_format_instructions()) chain = prompt | llm | parser result = chain.invoke({"input": "..."}) # result 是 dict它的解析逻辑比 Pydantic 宽松一些,内部会尝试从文本里提取 JSON 块,即使模型在 JSON 前后加了废话也能捞出来。但宽松也意味着不可控:模型可能输出一个字段名拼错的 JSON,JsonOutputParser 照样解析成功,错误要到下游用的时候才暴露。
我的经验是,如果结构固定,优先用 PydanticOutputParser,把校验前移。如果结构确实灵活,用 JsonOutputParser,但一定要在拿到 dict 后自己再做一层校验,别直接信任。
3.3 StructuredOutputParser:多字段场景的折中方案
StructuredOutputParser 是介于两者之间的方案。你用ResponseSchema定义字段,它生成格式说明,解析出来是 dict。
from langchain.output_parsers import StructuredOutputParser, ResponseSchema schemas = [ ResponseSchema(name="summary", description="内容摘要"), ResponseSchema(name="sentiment", description="情感倾向,positive/negative/neutral"), ResponseSchema(name="keywords", description="关键词,逗号分隔"), ] parser = StructuredOutputParser.from_response_schemas(schemas)它和 JsonOutputParser 的区别在于:StructuredOutputParser 会强制要求模型按指定的字段名输出,格式说明更明确。JsonOutputParser 更自由,模型可能自己发明字段名。
三个解析器的选型可以总结成一张表:
| 解析器 | 类型安全 | 格式约束 | 适用场景 | 解析失败率 |
|---|---|---|---|---|
| PydanticOutputParser | 强 | 严格 | 结构固定、需要类型校验 | 中 |
| JsonOutputParser | 弱 | 宽松 | 结构灵活、快速原型 | 低但错误隐蔽 |
| StructuredOutputParser | 中 | 中等 | 多字段、字段名固定 | 中 |
3.4 解析失败时的重试与修复策略
不管用哪个解析器,解析失败都是常态。模型不是编译器,它输出的 JSON 偶尔会多一个逗号、少一个引号、或者把true写成True。处理策略有三层。
第一层是 prompt 层面预防。在格式说明里明确写"只输出 JSON,不要有任何其他文字",并且给一两个示例。示例比说明管用,模型模仿示例的能力很强。
第二层是解析器层面兜底。LangChain 提供了OutputFixingParser,它包一层 LLM,解析失败时把错误信息和原始输出一起丢给模型,让它修复。
from langchain.output_parsers import OutputFixingParser base_parser = PydanticOutputParser(pydantic_object=PersonInfo) fixing_parser = OutputFixingParser.from_llm(parser=base_parser, llm=llm)代价是多一次 LLM 调用,延迟和成本都上去了。适合对成功率要求高、能接受额外开销的场景。
第三层是业务层面降级。解析失败时返回一个默认结构,或者把原始文本透传给前端让用户自己看。这层最土但最稳,生产环境一定要有。
4. ToolCall 与流式输出的共存方案
4.1 ToolCall 的本质:模型输出的是结构化参数
ToolCall 看起来是个新概念,本质上还是结构化输出。模型决定调用某个工具时,它输出的是一段 JSON,包含工具名和参数。LangChain 把这部分封装成tool_calls字段,挂在 AIMessage 上。
from langchain_core.tools import tool @tool def get_weather(city: str) -> str: """查询指定城市的天气。""" return f"{city} 今天晴,25 度。" llm_with_tools = llm.bind_tools([get_weather]) response = llm_with_tools.invoke("北京天气怎么样?") print(response.tool_calls) # [{'name': 'get_weather', 'args': {'city': '北京'}, 'id': 'call_xxx'}]bind_tools会把工具的定义(名字、描述、参数 schema)转成模型能理解的格式塞进请求。模型返回的tool_calls里,args就是结构化参数。所以 ToolCall 的可靠性,本质上还是结构化输出的可靠性。
4.2 流式场景下 ToolCall 参数是分片到达的
这是最容易踩坑的地方。流式模式下,tool_calls的参数不是一次性给你的,而是一片一片拼出来的。第一片可能是{"name": "get_weather", "args": ""},第二片是{"args": "{\"ci"},第三片是{"args": "ty\": \"北"},以此类推。
如果你在流式过程中直接读chunk.tool_calls[0]["args"],拿到的是残缺的 JSON 字符串,json.loads必然报错。正确做法是累积所有分片,等流结束后再解析。
async def stream_with_tools(chain, inputs): tool_call_buffer = {} async for chunk in chain.astream(inputs): if chunk.tool_call_chunks: for tc in chunk.tool_call_chunks: idx = tc.get("index", 0) if idx not in tool_call_buffer: tool_call_buffer[idx] = {"name": "", "args": "", "id": ""} if tc.get("name"): tool_call_buffer[idx]["name"] += tc["name"] if tc.get("args"): tool_call_buffer[idx]["args"] += tc["args"] if tc.get("id"): tool_call_buffer[idx]["id"] = tc["id"] if chunk.content: yield f"data: {chunk.content}\n\n" # 流结束后统一解析 for idx, tc in tool_call_buffer.items(): args = json.loads(tc["args"]) # 执行工具...index字段很关键。模型可能一次返回多个工具调用,每个有自己的 index,必须按 index 分组累积,不能混在一起。这个细节在文档里不显眼,但不处理就会串参数。
4.3 流式过程中如何给前端反馈工具调用状态
用户等工具调用时是焦虑的,因为界面上什么都没发生。好的做法是把工具调用的状态也通过 SSE 推给前端,让用户知道"模型正在查天气"。
可以自定义事件类型,SSE 支持event:字段:
yield f"event: tool_start\ndata: {json.dumps({'name': tc['name']})}\n\n" # 工具执行完 yield f"event: tool_end\ndata: {json.dumps({'name': tc['name'], 'result': result})}\n\n" # 正常文本 yield f"event: message\ndata: {chunk.content}\n\n"前端按event字段分流处理,tool_start时显示"正在调用 XX 工具",tool_end时显示结果,message时追加文本。这样整个交互过程是透明的,用户不会觉得卡住了。
4.4 工具执行结果回灌模型时的格式陷阱
工具执行完,结果要作为 ToolMessage 回灌给模型,让它继续生成。这里有个格式陷阱:ToolMessage 必须带上对应的tool_call_id,否则模型不知道这个结果对应哪个调用。
from langchain_core.messages import ToolMessage tool_messages = [] for tc in response.tool_calls: result = execute_tool(tc["name"], tc["args"]) tool_messages.append( ToolMessage(content=str(result), tool_call_id=tc["id"]) )tool_call_id必须和模型返回的id完全一致。如果模型一次调了多个工具,每个结果都要配对正确的 id,配错了模型会混乱。这个在单工具场景下不容易出错,多工具并发时就要小心。
另外,工具返回的内容如果是复杂结构,建议序列化成 JSON 字符串再塞进content。ToolMessage 的 content 是字符串类型,直接塞 dict 会被转成 Python 的 repr 格式,模型读起来别扭。
5. 把三者串起来:一个完整的流式结构化输出链路
5.1 整体架构:SSE 层、Chain 层、Parser 层怎么分工
把前面几块拼起来,一个完整的链路是这样的:FastAPI 提供 SSE 端点,接收请求后启动一个 LangChain Chain,Chain 内部可能触发 ToolCall,最终输出经过 OutputParser 结构化。SSE 层负责把过程中的每个事件(文本 token、工具调用、最终结构)推给前端。
分工上,SSE 层只管传输,不关心内容语义;Chain 层负责编排 LLM 和工具;Parser 层负责把最终文本转成结构化对象。三层解耦的好处是,换模型不影响 SSE 层,换传输方式不影响 Chain 层。
5.2 流式过程中做增量结构化解析的可行性
有人会想:能不能在流式过程中就做增量解析,边收边出结构化字段?技术上可行,但很麻烦。JSON 是上下文相关的,一个字段的值没结束时,你无法确定它是什么类型。比如{"age": 2后面可能是8}也可能是8.5},前者是 int 后者是 float。
有个取巧的办法是用支持流式解析的库,比如ijson,它能处理不完整的 JSON。但收益有限,因为大部分业务场景下,用户要的是最终结果,中间过程用文本展示就够了。我的建议是:流式阶段只做文本展示和工具状态反馈,结构化解析放到流结束后统一做。这样逻辑简单,出错也好排查。
5.3 完整代码:从请求到结构化结果的端到端实现
import json import asyncio from fastapi import FastAPI from fastapi.responses import StreamingResponse from langchain_core.prompts import ChatPromptTemplate from langchain_core.output_parsers import PydanticOutputParser from langchain_core.messages import ToolMessage from pydantic import BaseModel, Field app = FastAPI() class ExtractResult(BaseModel): title: str = Field(description="标题") tags: list[str] = Field(description="标签列表") parser = PydanticOutputParser(pydantic_object=ExtractResult) prompt = ChatPromptTemplate.from_messages([ ("system", "提取信息。\n{format_instructions}"), ("human", "{input}"), ]).partial(format_instructions=parser.get_format_instructions()) async def run_chain(input_text: str): chain = prompt | llm_with_tools tool_buffer = {} full_text = "" async for chunk in chain.astream({"input": input_text}): if chunk.content: full_text += chunk.content yield f"event: message\ndata: {json.dumps({'text': chunk.content})}\n\n" if chunk.tool_call_chunks: for tc in chunk.tool_call_chunks: idx = tc.get("index", 0) buf = tool_buffer.setdefault(idx, {"name": "", "args": "", "id": ""}) buf["name"] += tc.get("name") or "" buf["args"] += tc.get("args") or "" buf["id"] = tc.get("id") or buf["id"] yield f"event: tool_delta\ndata: {json.dumps({'name': buf['name']})}\n\n" # 处理工具调用 if tool_buffer: tool_messages = [] for idx, tc in tool_buffer.items(): args = json.loads(tc["args"]) result = execute_tool(tc["name"], args) tool_messages.append(ToolMessage(content=str(result), tool_call_id=tc["id"])) yield f"event: tool_end\ndata: {json.dumps({'name': tc['name'], 'result': str(result)})}\n\n" # 回灌模型继续生成 follow_up = await llm_with_tools.ainvoke( prompt.format_messages(input=input_text) + tool_messages ) full_text = follow_up.content yield f"event: message\ndata: {json.dumps({'text': full_text})}\n\n" # 最终结构化解析 try: parsed = parser.parse(full_text) yield f"event: result\ndata: {parsed.model_dump_json()}\n\n" except Exception as e: yield f"event: error\ndata: {json.dumps({'msg': str(e)})}\n\n" yield "event: done\ndata: [DONE]\n\n" @app.post("/extract") async def extract(payload: dict): return StreamingResponse( run_chain(payload["input"]), media_type="text/event-stream", headers={"X-Accel-Buffering": "no", "Cache-Control": "no-cache"}, )这段代码把前面讲的所有点都串起来了:心跳没加(可以按需补)、工具参数按 index 累积、工具结果带 tool_call_id 回灌、最终统一解析。实际用时还要加上异常处理和日志。
5.4 实测中的性能与稳定性观察
实测下来有几个观察值得分享。第一,SSE 的首字节延迟主要取决于模型的首 token 时间,和 SSE 封装本身关系不大。如果首字节慢,先查模型而不是查 SSE 代码。
第二,工具调用会显著拉长总时长,因为多了一次 LLM 往返。如果工具执行本身也慢,用户等待时间会很难看。建议给工具执行加超时,超时就返回一个"查询超时"的结果让模型继续,而不是一直挂着。
第三,结构化解析的失败率在流式场景下比非流式略高。原因是流式拼接的文本偶尔会有细微的截断问题,尤其是网络抖动时。所以流式场景下更要做好解析失败的兜底。
6. 几个容易翻车的细节和我的处理习惯
6.1 中文编码在 SSE 分片时的乱码问题
前面提过 TextDecoder 的stream: true,这里再强调一次服务端侧。Python 里如果手动拼接字节再 decode,同样要注意。用yield f"data: {text}\n\n"这种字符串 yield,FastAPI 会自己处理编码,一般没问题。但如果你在中间做了text.encode().decode()之类的操作,就可能把多字节字符切断。
我的习惯是全程用 str,不碰 bytes,让框架处理编码。只有在明确需要控制字节流时才手动处理,且一定用增量解码器。
6.2 工具参数里嵌套对象时的解析顺序
工具参数如果是嵌套结构,比如{"filter": {"city": "北京", "date": "2024-01-01"}},流式分片可能在任何位置切断。累积的时候是按字符串拼,所以顺序天然是对的,只要 index 分组正确就没问题。但解析时要注意,json.loads要求完整 JSON,所以必须等所有分片到齐。判断"到齐"的标志是流结束,不是某个特定字段出现。
6.3 多工具并发调用时的 index 管理
模型一次返回多个工具调用时,tool_call_chunks里的每个 chunk 都带 index。不同模型的 index 行为不完全一致,有的从 0 开始连续,有的可能跳号。稳妥的做法是用 dict 按 index 存,不要用 list 按顺序 append,否则跳号时会错位。
6.4 生产环境必须加的几道保险
第一道是超时。整个请求要有总超时,工具执行要有单独超时,LLM 调用也要有超时。任何一层没超时,都可能把连接挂死。
第二道是限流。SSE 长连接很占资源,没有限流的话并发一高就崩。按用户或 IP 做并发数限制。
第三道是日志。流式请求的日志要记录:请求开始、首字节时间、工具调用、解析结果、异常。出问题时这些日志是唯一的线索,因为流式请求很难复现。
第四道是降级。解析失败、工具超时、模型报错,都要有对应的降级响应,不能让前端一直转圈。
这几道保险加完,整个链路的稳定性会有质的提升。我自己的项目里,加之前线上偶发卡死,加之后基本没再出现过连接层面的问题。