ARTICLE DETAIL

资讯详情

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

LangChain+FastAPI实战:SSE流式与结构化输出完整链路

LangChain+FastAPI实战:SSE流式与结构化输出完整链路 你最近有没有被这个场景支配模型话还没说完UI 就开始等字段名倒是说完了可返回的是一段带着自然语言的散文前端愣是没法直接把那段 JSON 取出来。前阵子我帮一个团队把 LangChain 出的 AI 运营助手接到他们自己的 Vue 前端后端用的是 FastAPI。需求拆到最后其实就是三件事SSE 流式打底、OutputParser 负责把模型输出变成结构化数据、ToolCall 让模型真的去查数据库、调接口。这篇文章把完整链路给你跑通一遍也把我在真实项目里踩过的坑一并写了。不管你是刚入门 LangChain还是已经在 LangGraph 里搭过 Agent只要你想把“文本输出”变成“系统能直接消费的数据”这篇都能给你省几天测试时间。文章不会绕弯子所有代码都是我实测过的版本你直接抄就能跑。1. 为什么非要把“流式”和“结构化”放在一起聊1.1 两个词天生有冲突不是所有流都能停流式输出的天然节奏是“一个字一个字往外蹦”而结构化输出要求你交付一个“完整、自洽的对象”。这两件事看起来是冲突的解析器要拿到的是一整段 JSON但用户在前端看到的却是逐渐出现的内容。我在第一次做这个需求时也纠结过这个问题要么等模型完整输出后再一次性返回损失用户体验要么按 token 流式返回但前端拿到的全是没有边界的部分文本无法渲染结构化内容。后来想明白了关键不是“要不要流式”而是“流式里每一帧有没有固定 schema”。聊天对话框里可以流式显示文本但系统要看的是事件流不是文字流。也就是说后端需要沿着 SSE 推出一系列有类型标记的 JSON 事件例如工具开始、工具结束、内容增量、结构化结果完成。前端按类型消费这些事件该渲染的渲染该保存的保存。打个比方这就好比你点了一份外卖。你不需要等厨师把菜完完整整端到你面前才知道进度平台会不断推“商家接单”“骑手取餐”“配送中”这些结构化事件。这些事件有统一格式你的 App 才能把它们翻译成进度条。如果平台只给你一段“餐马上就好”的语音你是没法做进度条显示的。SSE 做的就是这个配送通知管道OutputParser 做的是把模型输出里的“散文”翻译成字段明确的 JSON而 ToolCall 让模型不只会说还能真的拿起工具帮你办事。三者一结合前面那个矛盾就解开了。1.2 SSE 和 WebSocketAI 场景下我为什么推荐 SSESSE 全称是 Server-Sent Events服务端推送事件协议非常简单就是 HTTP 响应里不断给你发一段text/event-stream文本。和 WebSocket 相比它既没有握手升级的复杂度也没有双向通信的状态机。AI 对话本来就是“服务端生成、客户端接收”的单向场景用 WebSocket 等于顺带实现了半个没用的反向通道白白增加维护成本。维度SSEWebSocket传输方向服务端到客户端单向双向底层协议HTTPtext/event-stream独立的 WebSocket 协议自动重连原生支持EventSource 自带需要自己实现消息格式文本按data:区分文本或二进制帧前端 API 复杂度极低需要处理连接状态、心跳、重连适用场景AI 聊天流、通知、进度推送实时协作、双向交互、游戏在 FastAPI 里返回 SSE 特别自然直接用StreamingResponse把 media_type 设成text/event-stream就行。客户端可以用原生EventSource但搞笑的是EventSource只支持 GET 请求。如果你要传用户消息、会话 ID、历史记录这些参数POST 更方便那就需要退回fetchReadableStream自己解析。后面我会专门写这个封装逻辑。1.3 结构化输出到底解决了谁的痛点很多人觉得模型能返回 JSON 就万事大吉。但如果你不主动约束输出结构模型随心情可能在 JSON 外面包一层 Markdown 代码块或者把字段名从orderId换成order_id再或者多一个换行符导致json.loads直接报错。结构化输出解决的核心问题是“下游可用性”。下游是数据库写入、审批流、工单系统它们不关心 LLM 的“文采”只认字段。OutputParser 在做的事就是把模型输出从“可能被我读懂”变成“一定被程序读懂”。2. LangChain 三大 OutputParser 实战与选型LangChain 生态里的 OutputParser 有很多种但日常实战真正高频出现的就是PydanticOutputParser、StructuredOutputParser和JsonOutputParser这三个。它们原理不同流式友好度也不同很多人一上来就乱用结果频繁报错。下面逐个拆。2.1 PydanticOutputParser离线批量场景的“严格考官”这个解析器是出镜率最高的一个核心思路是你在 Pydantic 模型里定义好字段和类型解析器把它变成一段格式说明塞进 Prompt要求模型严格按这个格式输出然后它拿到完整输出后用 Pydantic 校验并返回对象。from langchain_core.output_parsers import PydanticOutputParser from langchain_core.pydantic_v1 import BaseModel, Field from langchain_core.prompts import ChatPromptTemplate class OrderInfo(BaseModel): order_id: str Field(description订单号) status: str Field(description订单当前状态) amount: float Field(description订单金额) parser PydanticOutputParser(pydantic_objectOrderInfo) prompt ChatPromptTemplate.from_messages([ (system, 从用户描述中提取订单信息。严格遵循格式要求{format_instructions}), (human, {input}) ]).partial(format_instructionsparser.get_format_instructions()) chain prompt | llm | parser result chain.invoke({input: 我的订单10086现在是什么状态付款金额是多少}) print(result)它的优势非常明显类型校验严格字段缺失或类型不对时直接抛异常非常适合离线批量任务、数据清洗、定时抓取这类场景。但它有一个致命问题——非流式友好。因为它要求模型先输出完整内容再整体解析。如果你把它用在流式接口上前端会出现很长一段空白直到模型说完才能看到结构化结果。我的建议是如果响应要求在一秒内开始吐出内容就不要在流式链路里用 PydanticOutputParser它更适合做“终态校验器”。什么叫终态校验就是流式传输完成后把最后拼接好的完整文本再交给它做一次严格校验防止字段缺失。2.2 StructuredOutputParser老牌“映射器”适合简单 KVStructuredOutputParser是老项目里比较常见的解析器它和 Pydantic 版最大的区别是你不必定义 Pydantic 模型只需要声明一组ResponseSchema返回结果是普通 dict。from langchain.output_parsers import StructuredOutputParser, ResponseSchema schema [ ResponseSchema(nameorder_id, description订单号), ResponseSchema(namestatus, description订单状态), ] parser StructuredOutputParser.from_response_schemas(schema) format_instructions parser.get_format_instructions() chain prompt | llm | parser result chain.invoke({input: 查询10086}) print(result) # {order_id: 10086, status: ...}听起来很轻量但实战里它的问题也不少。返回结构完全依赖模型“照着格式说明编”没有强类型约束整数字段可能给你传字符串缺失字段也不提示。而且它的指令通常是“output a JSON object”有些模型会在外面套 Markdown 代码块还得再清理。现在新项目里我基本不用它了。但如果你是在维护 2023 年左右的老代码看到StructuredOutputParser时要知道它能干什么、弱在哪里别指望它做严格的类型校验。2.3 JsonOutputParser流式场景的“亲儿子”这才是接入 SSE 时最该用的解析器。JsonOutputParser内部实现了一个parse_partial_json它可以接收“尚未闭合的 JSON 片段”并尝试返回当前能识别的部分结构。这个特性让它在流式链路里非常顺手。from langchain_core.output_parsers import JsonOutputParser json_parser JsonOutputParser() # 模型流式吐出了不完整个 JSON 片段 partial {order_id: 10086, status: 已发 print(json_parser.parse_partial(partial)) # 多数情况下能返回 {order_id: 10086, status: 已发}注意并不是所有不完整片段都能解析成功比如 key 还没闭合的时候parse_partial会返回空或只返回已经完整解析的部分。但关键思路变了你是“边收边尝试”而不是“等满了再解析”。结合 LangChain 的异步流式接口可以这样用async def stream_structured(llm, prompt): async for chunk in llm.astream(prompt): content chunk.content if not content: continue parsed json_parser.parse_partial(content) if parsed: yield parsedJsonOutputParser在 SSE 场景还有个隐藏优势它不会因为输出里混入了其他字段而爆炸。它能容忍一定程度的噪音只要最终 JSON 核心结构能提取出来就行。不过它的校验能力弱你不能指望它把123转成float或者自动补全缺失字段。2.4 三个解析器到底怎么选一张表说清楚解析器返回类型流式友好校验强度最佳使用场景PydanticOutputParserPydantic 对象低强离线批量、严格校验、终态校验StructuredOutputParserdict低弱快速原型、简单 KV 提取JsonOutputParserdict高弱SSE 流式接口、部分 JSON 解析如果只记一条原则接口要走 SSE 且需要做到流式结构化展示直接选JsonOutputParser最后再用 Pydantic 模型兜底校验一次。这样性能和稳定性都有了。3. ToolCall 方案深入让模型真的“动手干活”3.1 ToolCall 和普通 Function Calling 是什么关系ToolCall是 LangChain 对模型工具调用能力的一层抽象。底层就是 OpenAI、Claude、Qwen 等模型都支持的 Function Calling但 LangChain 把它统一包装成了tool_call结构。模型在生成回复时不只给你文本还会在消息里附带一段字段明确的 JSON声明“我需要调用某个工具参数是什么”。一个典型的 tool_call 结构长这样{ name: get_order_status, args: { order_id: 10086 } }LangChain 拿到这个结构后会帮你找到对应的 Python 函数并执行然后把执行结果作为新的消息喂回模型让模型结合工具结果继续生成回复。这个循环就是 ToolCall 方案的核心。3.2 LangChain 工具绑定与执行循环定义一个工具很简单用tool装饰器函数的类型注解要写清楚因为类型注解会直接映射成工具描述里的 schema。from langchain_core.tools import tool tool def get_order_status(order_id: str) - str: 根据订单ID查询订单状态。 # 这里是真实场景一般会查数据库或调用外部接口 return 已发货预计明天送达 llm_with_tools llm.bind_tools([get_order_status]) response llm_with_tools.invoke(我的订单10086现在到哪了) # response.tool_calls 会包含调用信息 print(response.tool_calls)模型返回后你会看到类似这样的内容[{name: get_order_status, args: {order_id: 10086}, id: call_abc123}]接着你手动执行工具把结果拼进消息列表再让模型生成下一轮回复。如果不用 LangGraph你得自己写这个 while 循环while True: response llm_with_tools.invoke(messages) if not response.tool_calls: break messages.append(response) for call in response.tool_calls: tool_result get_order_status.invoke(call[args]) messages.append(ToolMessage(contenttool_result, tool_call_idcall[id]))写起来不难但状态一多循环、中断、重试都容易出问题。LangGraph 里已经帮你把这个流程编排好了用ToolNode接入即可。现在很多团队做 Agent 都直接跳过了手写循环这步转投 LangGraph。3.3 ToolCall 的流式处理别提前盼着完整参数这里有个大坑。模型在流式输出时tool_calls并不是一次性完整出现的它是一块块拼出来的。LangChain 会在每个流式块里返回tool_call_chunks里面是碎片化的参数片段需要自己累加。full_args tool_name async for chunk in llm.astream(messages): for tc_chunk in chunk.tool_call_chunks: if tc_chunk.name: tool_name tc_chunk.name if tc_chunk.args: full_args tc_chunk.args import json args json.loads(full_args)这里最忌讳的就是每收到一个 chunk 就去json.loads因为args很可能还没闭合直接解析会抛JSONDecodeError。你需要做的是先累加字符串等流结束或检测到完整闭合后再解析。如果实在要在流中展示工具调用进度可以用JsonOutputParser().parse_partial的同类思路做一个容错读取。3.4 在 LangGraph 里搭一个能“干活”的 AgentLangGraph 的核心价值是状态图和可编排。它把模型调用、工具执行、条件分支都变成图节点状态在节点间传递。FastAPI 只要把astream接口暴露出来就能把 LangGraph 里的每一步事件通过 SSE 推到前端。一个最简实现看起来是这样的from langgraph.graph import StateGraph, END from langgraph.prebuilt import ToolNode graph StateGraph(AgentState) def call_model(state): response llm_with_tools.invoke(state[messages]) return {messages: [response]} graph.add_node(model, call_model) graph.add_node(tools, ToolNode([get_order_status])) graph.add_edge(model, tools) graph.add_conditional_edges( model, lambda state: tools if state[messages][-1].tool_calls else END, )LangGraph 生态里现在还有 agent inbox 之类的方案专门处理人工介入环节。如果你的 Agent 要做审批、确认类操作可以在工具执行前先挂一个人工确认节点而不是让模型直接“擅自行动”。这些机制让“让 AI 真的下地干活”这句话从口号变成了可落地的工程实践。我在帮团队基于 DeerFlow 这类智能体框架做二次开发时见过最多的需求就是把底层的事物流式接口封装成统一 SSE 协议完成流式消息解析后再将结构化的 tool_call 和 content 事件暴露给上层。本质上你下面做的这套封装换到任何 Agent 框架上都是通用的。4. SSE 到结构化输出一套完整实战链路4.1 FastAPI 服务端怎么把事件流推出去服务端的关键是StreamingResponse每生成一条事件就往响应流里写一段 SSE 格式文本。from fastapi import FastAPI, Request from fastapi.responses import StreamingResponse import json app FastAPI() async def run_agent_with_events(messages): # 伪代码实际这里是 LangGraph 的 astream 逻辑 async for event in agent.astream_events({messages: messages}, versionv1): if event[event] on_chat_model_stream: delta event[data][chunk].content if delta: yield {type: content, delta: delta} elif event[event] on_tool_start: yield {type: tool_start, name: event[name], input: event[data].get(input)} elif event[event] on_tool_end: yield {type: tool_end, output: event[data].get(output)} app.post(/chat/stream) async def chat_stream(req: Request): payload await req.json() messages payload[messages] async def gen(): async for item in run_agent_with_events(messages): yield fdata: {json.dumps(item, ensure_asciiFalse)}\n\n return StreamingResponse(gen(), media_typetext/event-stream)这段代码最关键的是我统一了事件格式所有事件都有type字段。这样做的好处是前端拿到数据后不需要靠“猜测”来判断这是文本还是工具调用结果直接按类型分发即可。4.2 结构化输出在流式链路里的位置在实际链路里JsonOutputParser专门负责从局部 JSON 里提取可展示的结构化信息。比如模型在生成最终回复前可能先输出了一段 JSON 用来表达工具返回的订单信息。你想让前端订单卡片先渲染出来不需要等完整输出只需要在流式片段的 JSON 里解析到订单字段就立刻推送。async def run_agent_with_events(messages): json_parser JsonOutputParser() full_json async for event in agent.astream_events(...): if event[event] on_chat_model_stream: delta event[data][chunk].content yield {type: content, delta: delta} full_json delta parsed json_parser.parse_partial(full_json) if parsed and order_id in parsed: yield {type: order_card, data: parsed}这个order_card事件就是结构化输出的实战产物。前端收到后可以直接把订单号、状态、金额渲染成卡片而不需要在本地做复杂的字符串正则。4.3 前端 Vue3 fetch手动解析 SSE 的完整封装EventSource不支持 POST带消息历史的 AI 对话基本都是 POST所以我的习惯是用fetchReadableStream手动解析。核心思路是拿到 stream 后逐块读取按 SSE 协议的空行切分事件再逐条处理。async function streamChat(messages, onEvent) { const resp await fetch(/chat/stream, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ messages }), }); if (!resp.ok || !resp.body) throw new Error(stream failed); const reader resp.body.getReader(); const decoder new TextDecoder(utf-8); let buffer ; while (true) { const { value, done } await reader.read(); if (done) break; buffer decoder.decode(value, { stream: true }); const frames buffer.split(\n\n); buffer frames.pop(); for (const frame of frames) { if (!frame.startsWith(data:)) continue; const dataLine frame.replace(/^data:\s*/, ).trim(); if (!dataLine) continue; const event JSON.parse(dataLine); onEvent(event); } } }这个封装的几个细节值得注意TextDecoder一定要开{ stream: true }否则多字节字符比如中文正好被截断成两半时会出现乱码。buffer 处理不能少因为一次网络请求可能只返回半个事件帧必须积累到完整的\n\n才解析。每次解析完不要把 buffer 清空要保留最后一段不完整帧。4.4 实战案例一个带查询工具和结构化卡片的客服助手我用一个很常见的需求来跑通整条链路用户问“我的订单10086到哪了”Agent 调用get_order_status工具工具返回真实状态前端最终显示一段文本回复加一个订单卡片。事件流大致长这样event: {type: content, delta: 让我查一下订单状态} event: {type: tool_start, name: get_order_status, input: {order_id: 10086}} event: {type: tool_end, output: 已发货预计明天送达} event: {type: order_card, data: {order_id: 10086, status: 已发货, estimated: 明天送达}} event: {type: content, delta: 你的订单已经发货了}前端拿到order_card事件后渲染卡片拿到content事件后渲染流式文本。代码结构上可以做一个简单的分发器function handleEvent(event) { switch (event.type) { case content: appendStreamingText(event.delta); break; case order_card: renderOrderCard(event.data); break; case tool_start: showToolStatus(查询中...); break; case tool_end: showToolStatus(event.output); break; } }这个架构足够轻量但能覆盖绝大多数 AI 应用场景。你再怎么换模型、换框架只要 SSE 事件格式保持不变前端就永远不用改。5. 常见问题与排查实录5.1 SSE 断流“stream disconnected before completion: idle timeout waiting for SSE”这个报错是网关上最常见的报错看到它的第一反应不是查代码而是查网关的超时配置。问题本质是服务端和客户端之间有一个代理层该层设置了空闲超时。当 Agent 在处理工具调用或模型思考时短时间内没有任何数据通过连接代理就判断连接已空闲主动断开。解决办法有两个方向。第一在服务端生成器里发心跳包第二调大代理层读超时时间比如 Nginx 的proxy_read_timeout 600s。发心跳包是在没有实际业务数据时强行给连接制造“活动”同时不干扰前端解析。实现起来也很简单可以在生成器里做定时兜底跟异步队列配合使用import asyncio async def gen(): while True: try: event await queue.get(timeout15) if event is None: break yield fdata: {json.dumps(event, ensure_asciiFalse)}\n\n except asyncio.TimeoutError: # SSE 注释行不会被前端当业务数据处理但能有效避免 idle timeout yield : keep-alive\n\n那个: keep-alive不是数据是 SSE 的注释帧。它存在的意义就是告诉网关“我还活着”。这是我在排查断流问题时最有效的招数。5.2 OutputParser 在真实对话里不稳定十次有三五次报错这个现象我太熟了。模型输出里最常见的幺蛾子包括把 JSON 包在 Markdown 代码块json里、输出末尾多了逗号、字段不是合法 JSON key、直接输出注释文字。解决办法有两个。第一个办法是在解析前做清洗把多余代码块标记剥掉。第二个更推荐的办法是直接用with_structured_output它会利用模型的工具调用能力来生成结构化结果准确率远高于“用 Prompt 约束 正则提取”。LangChain 里这样用structured_llm llm.with_structured_output(OrderInfo) result structured_llm.invoke(我的订单10086现在是什么状态)with_structured_output本质上是把 Pydantic 模型转成工具 schema模型通过 tool_call 返回数据相当于走一条完全不同的输出通道比“让模型照着格式说明写 JSON 文本”稳定得多。你在流式场景里如果发现 JsonOutputParser 解析困难换个思路让模型直接bind_tools返回结构化结果再对 tool_call 做累加解析。5.3 ToolCall 参数解析失败永远不要早解析我在 3.3 节已经提醒过一次这里再强调一遍。流式下发 tool_call 时args是一点一点拼起来的。如果你在收到 100 个字符时就试图解析大概率会看到这样的报错Expecting property name enclosed in double quotes解决方案是要么等完整的tool_call_chunks收完再解析要么用容错解析工具。这里有一个工程经验在打印日志时一定要带full_args累积的字符串否则排查这类问题会非常痛苦。另外还有一个容易忽略的问题工具函数参数类型没写对。LangChain的工具 schema 依赖类型注解推断比如order_id: str和order_id: int生成的 schema 完全不同。如果不写类型注解模型会瞎猜参数类型错了它也不知道。5.4 前端拿到乱码或事件分割错乱乱码问题基本都在TextDecoder。网络流不是按字符边界切割的一个中文字符的 UTF-8 编码可能被拆到两个 chunk 里。如果直接decoder.decode(value)第一段就会出现奇怪的字符。正确做法是构造new TextDecoder(utf-8)并在每次调用decode(value, { stream: true })最后再调用一次decoder.decode()收尾。事件分割错乱则是我前面说过的 buffer 问题。SSE 协议的事件边界是连续的\n\n但你收到的网络字节流可能把一个事件的\n和下一个事件的开头拆断了。所以必须维护一个buffer先拼起来再按\n\n切分。还有一个前端容易踩的坑使用EventSource时它会自动重连而使用fetch时不会。如果你的场景需要断线重连就必须手动设计重试机制比如在流异常中断时间隔几秒重新发起请求同时利用Last-Event-ID之类的标识做续传。排查这些 SSE 问题我总结了一个速查清单可以直接对着看。症状可能原因排查重点idle timeout 断开代理层空闲超时是否发心跳、网关超时配置JSON 解析报错模型输出多余标记或非法字符是否用 with_structured_output、是否清洗tool_call 参数缺失流式 chunks 未累积完检查 full_args 累加逻辑前端乱码TextDecoder 未开 stream换 TextDecoder 参数事件被截断buffer 逻辑bug检查 split 后是否保留了不完整块连接断后无响应fetch 无自动重连实现手动重试逻辑6. 最后分享几个实战里拿命换来的经验说了这么多最后聊几句我做这类系统时最深的体会。第一不要在一开始就追求“全部结构化”。很多团队希望模型输出 100% 符合一个复杂的 Pydantic 模型结果 prompt 越长、约束越多模型越容易出错。我的做法是先流式传入主要内容等模型输出差不多完整了再用 Pydantic 做终态校验。真要修复字段缺失也来得及。第二SSE 事件流一定要从一开始就设计好type字段。我见过很多项目前端辛辛苦苦等流式接口结果服务端传出来的是一堆没有任何标识的 JSON前端只能靠猜。与其让前端去做语义推断不如后端多花五分钟设计事件类型。这个收益会一直持续到项目维护期。第三工具调用和结构化输出不是两套系统。它们本质上是同一条链路上的两个节点模型通过 tool_call 生成结构化数据前端通过 SSE 实时展示进度最终通过 OutputParser 完成严格校验。把这个闭环打通AI Agent 才算真正“下地干活”了。如果你正在做类似的项目遇到具体的坑按我上面那个排查速查表对号入座基本都能解决。最麻烦的从来不是某一个技术点而是你把 SSE、OutputParser、ToolCall 串起来那一刻的细节处理。希望这篇对你有点用。
返回列表