流式输出和结构化输出,这两个词放在一起本身就有点矛盾。SSE 追求的是"边生成边推送",一个字一个字往外蹦,用户体感好;而结构化输出要求的是"整段 JSON 完整可解析",字段一个不能少、类型一个不能错。我在实际项目里踩过最典型的坑就是:前端用 EventSource 接流,后端用 LangChain 的 PydanticOutputParser 做解析,结果流式返回的 JSON 是分片的,解析器拿到半截字符串直接抛异常,整个链路崩掉。后来才搞明白,这两件事需要分开处理——流式负责传输体验,结构化负责数据落地,中间得靠 OutputParser 和 ToolCall 做桥接。
这篇内容适合正在用 LangChain 做 AI 应用、被流式输出和结构化解析同时折磨的开发者。不管你是刚接触 LangChain 的新手,还是已经写过几版 Agent 但总在解析环节翻车的老手,下面这些从实战里抠出来的细节应该都能对上你的痛点。我会把 SSE 流式的底层机制、LangChain 三大 OutputParser 的适用边界、ToolCall 的结构化方案,以及它们怎么串起来用,全部拆开讲清楚。
1. 先搞清楚 SSE 流式到底在传什么
很多人一上来就写EventSource接流,但根本没想过 SSE 协议层到底在传什么格式的数据。这就像你开车不看仪表盘,油没了才知道慌。
1.1 SSE 的数据帧格式与 LangChain 的流式适配
SSE 的全称是 Server-Sent Events,它本质上就是一个长连接,服务端往客户端持续推送文本。每一条消息的格式是固定的:
event: message data: {"content": "你"} data: {"content": "好"}注意那个空行,它是消息分隔符。没有空行,浏览器端的 EventSource 就不会触发onmessage。我在早期项目里自己手写 SSE 服务端时,就是因为忘了加空行,前端一直收不到消息,排查了两个小时才发现是格式问题。
LangChain 的流式输出最终也是要转成这个格式。它的astream方法返回的是一个异步生成器,每次 yield 出来的是一个AIMessageChunk对象。你需要自己把它序列化成 SSE 帧。常见的做法是用 FastAPI 的StreamingResponse:
from fastapi import FastAPI from fastapi.responses import StreamingResponse from langchain_openai import ChatOpenAI app = FastAPI() llm = ChatOpenAI(model="gpt-4o-mini", streaming=True) async def event_generator(): async for chunk in llm.astream("用一句话介绍 LangChain"): if chunk.content: yield f"data: {chunk.content}\n\n" @app.get("/stream") async def stream(): return StreamingResponse(event_generator(), media_type="text/event-stream")这里有个细节:media_type必须设成text/event-stream,否则浏览器不会把它当成 SSE 流来处理。另外chunk.content要做空判断,因为有些 chunk 只包含元数据没有实际文本。
1.2 流式场景下结构化输出的核心矛盾
为什么流式和结构化输出会打架?因为 JSON 的语法结构决定了它不能像自然语言那样随意分片。你想想,{"name": "张这半截 JSON 是没法解析的,解析器会直接报JSONDecodeError。
我见过有人试图在流式过程中不断尝试解析累积的字符串,一旦解析成功就返回。这个思路理论上可行,但实际用起来问题很多:第一,JSON 的嵌套结构可能导致部分解析成功但字段不完整;第二,每次 chunk 到达都做一次完整解析,性能开销随文本长度线性增长;第三,如果模型输出的 JSON 格式有轻微偏差,累积解析会一直失败直到超时。
提示:流式场景下不要试图对每个 chunk 做结构化解析,正确的做法是流式传输原始文本,等流结束后再做一次完整的结构化解析。
那有没有办法既流式又结构化?有,但需要换思路。LangChain 提供了几种方案,下面会逐一拆解。
1.3 前端消费 SSE 流的常见坑
前端用 EventSource 接流时,有几个坑几乎每个人都会踩:
第一个坑是自动重连。EventSource 默认在连接断开后会自动重连,如果你的接口不支持断点续传,重连后模型会重新生成一遍,用户看到重复内容。解决办法是在服务端发送完最后一条消息后,主动发送一个event: done的事件,前端收到后调用eventSource.close()。
第二个坑是中文乱码。SSE 默认使用 UTF-8 编码,但如果你的服务端没有正确设置charset=utf-8,中文会变成乱码。FastAPI 的StreamingResponse默认就是 UTF-8,但如果你用的是 Flask 或其他框架,需要手动指定。
第三个坑是代理缓冲。某些反向代理会缓冲 SSE 流,导致前端收到消息的延迟很大。解决办法是在响应头里加上X-Accel-Buffering: no,告诉代理不要缓冲。
const eventSource = new EventSource('/stream'); let buffer = ''; eventSource.onmessage = (event) => { buffer += event.data; document.getElementById('output').textContent = buffer; }; eventSource.addEventListener('done', () => { eventSource.close(); // 流结束后,把 buffer 发给后端做结构化解析 parseStructured(buffer); });这个模式是我目前用得最顺手的:流式阶段只做展示,流结束后再调一次结构化解析接口。用户体验和数据准确性都能兼顾。
2. LangChain 三大 OutputParser 的适用边界
LangChain 的 OutputParser 是结构化输出的核心工具,但很多人只知道 PydanticOutputParser,不知道另外两个的存在,更不知道它们各自的适用场景。选错了 Parser,轻则解析失败率飙升,重则整个链路不可用。
2.1 PydanticOutputParser:强类型场景的首选
PydanticOutputParser 是三个里面最严格的,它要求模型输出的 JSON 必须完全符合你定义的 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}") ]) chain = prompt | llm | parser result = chain.invoke({ "input": "张三今年28岁,会Python和Java", "format_instructions": parser.get_format_instructions() })get_format_instructions()这个方法很关键,它会把 Pydantic 模型的 JSON Schema 转成一段自然语言描述,塞进 prompt 里告诉模型该怎么输出。很多人忘了加这个,结果模型输出的字段名跟模型定义对不上,解析直接失败。
PydanticOutputParser 适合什么场景?适合字段固定、类型明确、不允许有额外字段的场景。比如从简历里提取姓名、年龄、技能,从合同里提取甲方、乙方、金额。这些场景下数据质量要求高,宁可解析失败也不能出错。
但它也有明显的短板:对模型能力要求高。如果你用的是小参数量的本地模型,它经常输出不符合 Schema 的 JSON,解析失败率可能超过 30%。这时候就需要考虑下面两种方案。
2.2 StructuredOutputParser:灵活性与结构化的折中
StructuredOutputParser 不像 PydanticOutputParser 那样要求严格的类型定义,它只需要你告诉它有哪些字段、每个字段是什么类型,然后它会生成一个 ResponseSchema 列表。
from langchain.output_parsers import StructuredOutputParser, ResponseSchema schemas = [ ResponseSchema(name="product_name", description="产品名称"), ResponseSchema(name="price", description="产品价格,数字类型"), ResponseSchema(name="tags", description="产品标签,逗号分隔的字符串") ] parser = StructuredOutputParser.from_response_schemas(schemas) format_instructions = parser.get_format_instructions() prompt = ChatPromptTemplate.from_template( "从以下描述中提取产品信息:\n{input}\n\n{format_instructions}" ) chain = prompt | llm | parser result = chain.invoke({ "input": "这款无线蓝牙耳机售价299元,主打降噪和长续航", "format_instructions": format_instructions })StructuredOutputParser 的输出是一个字典,字段值都是字符串类型(即使你描述里写了"数字类型",它返回的也是字符串)。这个特性在某些场景下反而是优势——你不需要提前定义严格的类型,拿到结果后自己做类型转换就行。
它的适用场景是:字段不固定、需要动态调整、或者模型能力有限无法稳定输出严格 JSON 的情况。比如做一个通用的信息抽取工具,用户自己定义要抽取哪些字段,这时候用 StructuredOutputParser 就比 PydanticOutputParser 灵活得多。
但要注意,它的解析容错性虽然比 PydanticOutputParser 好,但也不是万能的。如果模型输出的 JSON 缺少某个字段,它还是会报错。我一般的做法是在 prompt 里加一句"如果某个字段无法从原文中提取,请填写'未知'",这样能大幅降低解析失败率。
2.3 OutputFunctionsParser:ToolCall 场景的底层支撑
OutputFunctionsParser 是三个里面最特殊的一个,它不解析 JSON 文本,而是解析模型返回的 function call 参数。在 OpenAI 的 Function Calling 机制下,模型不会直接输出 JSON 文本,而是返回一个function_call对象,里面包含函数名和参数。
from langchain.output_parsers.openai_functions import OutputFunctionsParser from langchain_core.prompts import ChatPromptTemplate function_schema = { "name": "extract_info", "description": "提取人物信息", "parameters": { "type": "object", "properties": { "name": {"type": "string", "description": "姓名"}, "age": {"type": "integer", "description": "年龄"} }, "required": ["name", "age"] } } parser = OutputFunctionsParser() prompt = ChatPromptTemplate.from_template("提取以下文本中的人物信息:{input}") chain = prompt | llm.bind(functions=[function_schema]) | parser result = chain.invoke({"input": "李四,35岁,工程师"})OutputFunctionsParser 的优势在于:它绕过了 JSON 文本解析这一步。模型返回的 function call 参数已经是结构化的对象了,不需要再做字符串解析。这就从根本上避免了 JSON 格式错误的问题。
但它的限制也很明显:依赖模型支持 Function Calling。OpenAI 的 GPT 系列、Claude、以及部分开源模型支持这个能力,但很多小模型不支持。另外,Function Calling 的输出是"全有或全无"的,模型要么返回完整的参数,要么不返回,没法做流式输出。
2.4 三种 Parser 的选型对照
| 维度 | PydanticOutputParser | StructuredOutputParser | OutputFunctionsParser |
|---|---|---|---|
| 类型严格度 | 高,强类型校验 | 中,字段值为字符串 | 高,由 Schema 定义 |
| 模型要求 | 高,需要较强 JSON 能力 | 中,容错性较好 | 需要支持 Function Calling |
| 流式支持 | 不支持 | 不支持 | 不支持 |
| 适用场景 | 字段固定的信息提取 | 动态字段的通用抽取 | 工具调用、Agent 场景 |
| 解析失败率 | 较高(小模型上) | 中等 | 低(依赖模型能力) |
选型逻辑其实很简单:如果你的模型支持 Function Calling 且场景适合,优先用 OutputFunctionsParser;如果需要严格的类型校验且模型能力够强,用 PydanticOutputParser;如果字段不固定或者模型能力有限,用 StructuredOutputParser。
3. ToolCall 方案:让模型主动输出结构化数据
ToolCall(工具调用)是 LangChain 里做结构化输出的另一条路。它的思路跟 OutputParser 不同:不是让模型输出 JSON 文本再解析,而是让模型直接调用一个预定义的工具,工具的参数就是结构化的数据。
3.1 ToolCall 与 Function Calling 的关系
先理清一个概念:Function Calling 是模型层面的能力,ToolCall 是 LangChain 层面的封装。OpenAI 的 API 里叫functions,LangChain 把它包装成了Tool对象。
from langchain_core.tools import tool from langchain_openai import ChatOpenAI from langchain_core.prompts import ChatPromptTemplate @tool def extract_person(name: str, age: int, skills: list[str]) -> dict: """提取人物信息""" return {"name": name, "age": age, "skills": skills} llm_with_tools = ChatOpenAI(model="gpt-4o-mini").bind_tools([extract_person]) prompt = ChatPromptTemplate.from_template("从以下文本提取人物信息:{input}") chain = prompt | llm_with_tools result = chain.invoke({"input": "王五,42岁,擅长数据分析和机器学习"}) tool_calls = result.tool_calls # tool_calls[0]["args"] 就是结构化数据bind_tools会把工具的定义转成模型能理解的 Schema,模型在生成时会判断是否需要调用工具。如果需要,它会返回一个tool_calls列表,每个元素包含工具名和参数。
这个方案最大的好处是:模型输出的参数已经是解析好的 Python 对象,不需要你再做 JSON 解析。而且 LangChain 会自动校验参数类型,如果模型输出的参数类型不对,会直接报错而不是静默失败。
3.2 用 ToolCall 做结构化输出的完整链路
实际项目里,ToolCall 做结构化输出通常需要配合一个"执行器"来跑通完整链路。因为模型只是"决定调用工具",真正执行工具的是你的代码。
from langchain_core.runnables import RunnablePassthrough from langchain_core.output_parsers import JsonOutputParser def execute_tool_calls(message): results = [] for tool_call in message.tool_calls: tool_name = tool_call["name"] tool_args = tool_call["args"] # 这里可以根据 tool_name 分发到不同的处理逻辑 results.append({"tool": tool_name, "args": tool_args}) return results chain = prompt | llm_with_tools | execute_tool_calls result = chain.invoke({"input": "赵六,30岁,会前端开发和UI设计"})这个链路的关键在于execute_tool_calls这个函数。它接收模型返回的AIMessage,提取tool_calls字段,然后做后续处理。你可以把结果存数据库、发给前端、或者触发其他业务流程。
注意:ToolCall 返回的参数虽然已经是结构化对象,但不代表它一定符合你的业务校验规则。比如模型可能返回一个负数年龄,或者一个不存在的技能名称。业务层面的校验还是得自己做。
3.3 ToolCall 在 Agent 场景下的结构化优势
在 Agent 场景下,ToolCall 的优势更加明显。Agent 需要根据用户输入决定调用哪个工具、传什么参数,这本身就是结构化输出的过程。
from langchain.agents import AgentExecutor, create_openai_tools_agent tools = [extract_person, search_database, send_email] agent = create_openai_tools_agent(llm, tools, prompt) agent_executor = AgentExecutor(agent=agent, tools=tools) result = agent_executor.invoke({"input": "帮我查一下张三的信息,然后发邮件给他"})Agent 会自动规划:先调用search_database查张三,拿到结果后调用send_email发邮件。每一步的工具调用参数都是结构化的,不需要你手动解析。
我实测下来,ToolCall 方案在 Agent 场景下的结构化输出成功率明显高于 OutputParser 方案。原因很简单:模型在 Function Calling 模式下经过了专门的训练,输出格式的稳定性远高于自由文本生成 JSON。
3.4 ToolCall 方案的局限与应对
ToolCall 不是银弹,它有几个硬性限制:
第一,模型必须支持 Function Calling。如果你用的是本地部署的小模型,比如 Qwen 7B 或 Llama 3 8B,它们可能不支持或者支持得不好。这种情况下只能退回 OutputParser 方案。
第二,ToolCall 不支持流式输出。模型要么返回完整的工具调用参数,要么不返回。你没法像文本生成那样一个字一个字往外推。如果业务要求流式展示,需要做特殊处理——比如流式展示"正在分析..."的提示,等工具调用完成后再一次性展示结果。
第三,工具定义会增加 token 消耗。每个工具的定义都会作为 system prompt 的一部分发给模型,工具越多,token 消耗越大。我一般建议单次绑定的工具不超过 10 个,超过的话考虑做工具分组或动态绑定。
4. 流式与结构化的桥接方案
前面分别讲了 SSE 流式、OutputParser 和 ToolCall,现在到了最关键的部分:怎么把它们串起来,既保证流式体验,又拿到结构化数据。
4.1 两阶段方案:流式展示 + 流后解析
这是我最推荐的方案,也是实际项目里用得最多的。核心思路是:流式阶段只负责把模型的原始输出推给前端展示,流结束后再把完整文本发给解析器做结构化处理。
async def stream_and_parse(user_input: str): full_text = "" # 第一阶段:流式输出 async for chunk in llm.astream(user_input): if chunk.content: full_text += chunk.content yield f"data: {chunk.content}\n\n" # 第二阶段:结构化解析 yield f"event: done\ndata: {json.dumps({'status': 'stream_end'})}\n\n" # 解析完整文本 try: structured = parser.parse(full_text) yield f"event: structured\ndata: {json.dumps(structured.dict())}\n\n" except Exception as e: yield f"event: error\ndata: {json.dumps({'error': str(e)})}\n\n"这个方案的好处是职责清晰:流式负责体验,解析负责数据。前端收到done事件后知道流结束了,收到structured事件后拿到结构化数据。
但有个细节要注意:解析失败的处理。如果模型输出的 JSON 格式有问题,解析会失败。这时候不能直接报错给用户,而是要有降级方案。我的做法是:解析失败时,把原始文本返回给前端,同时记录日志,后续人工排查。
4.2 流式 JSON 的增量解析思路
如果你确实需要在流式过程中就拿到部分结构化数据,可以考虑增量解析。这个方案的核心是:维护一个累积的 JSON 字符串,每次收到新 chunk 后尝试解析,如果解析成功就返回部分结果。
import json class IncrementalJSONParser: def __init__(self): self.buffer = "" self.parsed_keys = set() def feed(self, chunk: str): self.buffer += chunk # 尝试补全 JSON 并解析 try: # 简单场景下,尝试在 buffer 末尾补全括号 candidate = self.buffer open_braces = candidate.count('{') - candidate.count('}') candidate += '}' * open_braces data = json.loads(candidate) # 提取新增的字段 new_fields = {k: v for k, v in data.items() if k not in self.parsed_keys} self.parsed_keys.update(new_fields.keys()) return new_fields except json.JSONDecodeError: return None这个方案在简单场景下能用,但有几个明显的限制:第一,只适合扁平 JSON,嵌套结构补全逻辑很复杂;第二,如果模型输出的 JSON 有语法错误,增量解析会一直失败;第三,频繁的 JSON 解析有性能开销。
我个人的建议是:除非业务强需求,否则不要用增量解析。两阶段方案已经能满足 90% 的场景,而且稳定性和可维护性都好得多。
4.3 用 ToolCall 做流式场景下的结构化兜底
还有一个折中方案:流式阶段用普通文本生成,流结束后用 ToolCall 做结构化提取。这样流式体验有了,结构化数据的准确性也有保障。
async def stream_with_toolcall_fallback(user_input: str): # 第一阶段:流式生成 full_text = "" async for chunk in llm.astream(user_input): if chunk.content: full_text += chunk.content yield f"data: {chunk.content}\n\n" # 第二阶段:用 ToolCall 做结构化提取 extraction_prompt = f"从以下文本中提取结构化信息:\n{full_text}" result = await llm_with_tools.ainvoke(extraction_prompt) if result.tool_calls: structured_data = result.tool_calls[0]["args"] yield f"event: structured\ndata: {json.dumps(structured_data)}\n\n"这个方案相当于做了两次模型调用:第一次生成内容,第二次提取结构。成本翻倍,但准确性最高。适合对数据质量要求极高的场景,比如合同信息提取、财务数据录入。
4.4 三种桥接方案的对比与选型
| 方案 | 流式体验 | 结构化准确性 | 成本 | 适用场景 |
|---|---|---|---|---|
| 两阶段方案 | 好 | 中 | 低 | 通用场景 |
| 增量解析 | 好 | 低 | 低 | 简单扁平 JSON |
| ToolCall 兜底 | 好 | 高 | 高 | 高准确性要求 |
选型建议:先用两阶段方案跑通,如果解析失败率超过 10%,再考虑 ToolCall 兜底。增量解析除非有特殊需求,否则不建议用。
5. 实战中踩过的坑与排查链路
这一节不讲理论,只讲我在实际项目里踩过的坑和排查过程。这些经验在官方文档里找不到,但每一个都让我多花了好几个小时。
5.1 SSE 流中断:idle timeout 的完整排查过程
有一次线上环境,用户反馈流式输出经常中断,前端报错stream disconnected before completion: idle timeout waiting for sse。这个报错信息很明确:SSE 连接因为空闲超时被断开了。
排查第一步:确认超时时间。我查了 Nginx 的配置,proxy_read_timeout默认是 60 秒。如果模型生成一个长回复超过 60 秒没有新数据推送,Nginx 就会断开连接。
排查第二步:确认模型是否真的有空闲期。我加了日志,发现模型在生成过程中确实有超过 60 秒没有输出任何 chunk 的情况。原因是模型在处理复杂推理时,会在内部"思考"很久,期间不产生任何输出。
排查第三步:解决方案。有两个方向:一是调大 Nginx 的超时时间,二是让服务端定期发送心跳。我两个都做了:
location /stream { proxy_pass http://backend; proxy_read_timeout 300s; proxy_set_header Connection ''; proxy_http_version 1.1; chunked_transfer_encoding off; proxy_buffering off; }async def event_generator(): last_heartbeat = time.time() async for chunk in llm.astream(user_input): if chunk.content: yield f"data: {chunk.content}\n\n" last_heartbeat = time.time() elif time.time() - last_heartbeat > 15: # 超过15秒没有内容,发送心跳 yield f": heartbeat\n\n" last_heartbeat = time.time()心跳的格式是: heartbeat\n\n,以冒号开头的行在 SSE 协议里是注释,前端不会触发onmessage,但能保持连接活跃。
提示:心跳间隔建议设为超时时间的 1/3 到 1/2。比如超时 60 秒,心跳设 20-30 秒一次。
5.2 PydanticOutputParser 解析失败的三种典型情况
PydanticOutputParser 的解析失败是我遇到最多的坑,总结下来有三种典型情况:
情况一:模型输出了 markdown 代码块包裹的 JSON。模型经常输出```json\n{...}\n```这种格式,解析器拿到带反引号的字符串直接报错。解决办法是在 prompt 里明确说"直接输出 JSON,不要用代码块包裹",或者在解析前做一次字符串清洗:
def clean_json_string(text: str) -> str: text = text.strip() if text.startswith("```json"): text = text[7:] if text.startswith("```"): text = text[3:] if text.endswith("```"): text = text[:-3] return text.strip()情况二:模型输出的字段名跟定义不一致。比如你定义的是user_name,模型输出的是userName。这种问题在 prompt 里加一句"字段名必须严格使用以下名称"能缓解,但不能完全避免。更稳妥的做法是在 Pydantic 模型里加 alias:
from pydantic import BaseModel, Field class UserInfo(BaseModel): user_name: str = Field(alias="userName") class Config: populate_by_name = True情况三:模型输出的类型不对。比如你定义age: int,模型输出了"28"(字符串)。Pydantic 默认会尝试做类型转换,但有些情况下会失败。解决办法是在 Field 里加strict=False,或者在 prompt 里强调类型。
5.3 ToolCall 参数校验的边界情况
ToolCall 虽然比 OutputParser 稳定,但也不是没有坑。我遇到过几种边界情况:
情况一:模型返回了工具定义里不存在的参数。比如工具定义只有name和age,模型却返回了name、age、email。LangChain 默认会忽略多余参数,但如果你开了严格模式,会直接报错。
情况二:必填参数缺失。模型有时候会漏掉某个必填参数。这种情况下 LangChain 会报ValidationError。我的做法是在工具函数里给所有参数设默认值,然后在业务逻辑里判断是否为空。
情况三:参数类型不匹配。比如定义age: int,模型返回了"二十八"。这种在 Function Calling 模式下比较少见,但一旦出现就是硬错误。解决办法是在工具函数里做类型转换和校验。
5.4 流式与结构化并发的资源竞争问题
最后一个坑比较隐蔽:当多个用户同时请求流式接口时,如果每个请求都创建一个新的 LLM 连接,资源消耗会很大。我一开始没注意这个问题,上线后发现并发稍微高一点就出现连接超时。
解决办法是用连接池 + 信号量控制并发:
import asyncio semaphore = asyncio.Semaphore(10) # 最多10个并发 async def stream_handler(user_input: str): async with semaphore: async for chunk in llm.astream(user_input): yield f"data: {chunk.content}\n\n"信号量的值根据你的模型服务承载能力来定。如果是调 OpenAI 的 API,一般 10-20 个并发没问题;如果是本地部署的模型,要根据 GPU 显存来调整。
另外,LangChain 的astream方法在连接断开时不会自动清理资源,需要在finally块里手动关闭:
async def stream_handler(user_input: str): try: async for chunk in llm.astream(user_input): yield f"data: {chunk.content}\n\n" finally: await llm.async_client.close()这些坑每一个都让我在深夜排查过,希望你看完能少走点弯路。流式和结构化输出本身不复杂,复杂的是它们之间的衔接和边界情况的处理。把上面这些方案和排查思路吃透,大部分场景都能覆盖了。