☰
FastAPI+LangChain AI Agent:SSE流式、OutputParser与ToolCall实战
2026/10/6 6:42:07 网站建设 项目流程

最近把一套基于 FastAPI + LangChain 的 AI Agent 服务接到了前端,SSE 流式传输、结构化输出、OutputParser 和 ToolCall 这四个词几乎缠了我一整个周末。第一版方案非常天真:让模型把回答一次性生成完,后端json.dumps成一个大 JSON 丢给前端。结果前端同事接二连三问"你什么时候告诉我这段 Markdown 应该高亮?什么时候渲染工具调用状态?报错了前端怎么知道?"——被问住了。

这篇文章就是当时排查和重构的记录。它适合被 LangChain 输出格式不稳定折腾过的人、需要自己封装 SSE 流式接口的开发者,以及想搞清楚 ToolCall 到底和 OutputParser 什么关系的人。我会直接给代码、给结论、给踩坑点,不铺垫概念。

1. 先理解"流式"和"结构化"为什么是一对矛盾:模型擅长说,程序要的是型

1.1 矛盾的根源在模型而不在工程

语言模型本质上是概率式文本生成器,它最擅长的是"给人讲故事"。但程序需要的是带 schema 的、字段完整、类型明确的"事实"。流式则是另一个维度的问题:模型生成一个 JSON 对象时,没法保证第一个 token 到达时,后续的 key 顺序、引号、缩进是完整的。前端如果直接尝试JSON.parse流里的每个片段,结果只可能是两种:语法错误,或者等到一个被截断的半拉子对象。

这不是 FastAPI 或 LangChain 的 bug,是语义层级不同。SSE 负责传输"时间的切片",OutputParser 负责约束"内容的形状",ToolCall 负责表达"动作的意图"。这三件事混在一个层级里处理,必然出事。

1.2 我的分层策略:先给事件定性,再谈内容

我在实际项目里把协议拆成了两层。

第一层是SSE 事件层的定性。每个事件帧都带event字段,只有几种固定值:token、tool_start、tool_end、done、error。data字段只放该事件类型对应的最小数据单元。

第二层是内容层的结构化。只有到达done之前,前端才会把收到的所有token片段拼接成完整文本交给 OutputParser 或者渲染。工具调用参数则在tool_end事件里作为一个完整的 JSON 对象下发,不在中间态里解析。

这样前端永远不需要处理"半个 JSON 对象"这种状态。它只负责按event分发,看到token就追加文本,看到tool_start就显示一个加载卡片,看到done才拼完整结果。后端才负责用 LangChain 的 OutputParser 和 ToolCall 机制去约束模型输出。职责分清楚之后,后面所有代码写起来都顺了。

2. FastAPI 的 SSE 端点从能跑到稳定:idle timeout 断连的完整排查

2.1 基础版本:StreamingResponse 直接上

FastAPI 做 SSE 其实不复杂,核心就是StreamingResponse,把媒体类型设成text/event-stream。一个最基础的可运行版本长这样:

import json from fastapi import FastAPI, Request from fastapi.responses import StreamingResponse app = FastAPI() @app.post("/chat") async def chat(request: Request): payload = await request.json() async def event_generator(): # your_agent.astream 是 LangChain/LangGraph 的异步流式入口 async for event in your_agent.astream(payload["messages"]): yield f"event: {event['type']}\ndata: {json.dumps(event['data'], ensure_ascii=False)}\n\n" 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 这类反向代理,这个响应不要缓冲,来一块发一块。第二,ensure_ascii=False保证中文不会被转成\uXXXX序列。如果省略,中文以转义形式到达前端,虽然功能没错,但阅读日志和调试时非常折磨。

2.2 线上报错:stream disconnected before completion: idle timeout waiting for sse

这个错误最初看到时我以为是代码写错了。网上一搜,错误文本里直接带着 "idle timeout waiting for sse",典型的中间层问题。实际排查链路是这样的:

模型 Agent 在执行工具调用时,经常会有 10 到 30 秒的"思考间隔"。模型在后台调用搜索或数据库工具时,不会向前端返回任何 token。如果这段时间超过反向代理的空闲超时阈值,代理层就认为连接已经死了,主动断开。浏览器端表现为:页面等了几十秒,突然抛出一个连接错误,抛出的时间点恰好和工具调用的思考期重合。

我用curl -N直接打到 FastAPI 端口复现不了,因为绕过 Nginx 就没这个错误。加上 Nginx 后又稳定复现。当时先怀疑是不是 Nginx 配置了proxy_read_timeout 60s,检查后确认是默认值导致。但我不建议直接把这个超时调到 3600 秒就完事,因为根因是"长时间无数据",正确做法是让连接保持"有数据可读"的状态。

2.3 解决方案:心跳帧让连接永远不空闲

SSE 协议支持注释行作为心跳。一个: ping开头的行,客户端会忽略且不会触发任何事件。我改造后的生成器长这样:

import asyncio import time async def event_generator(): queue: asyncio.Queue = asyncio.Queue(maxsize=100) async def agent_producer(): async for event in your_agent.astream(payload["messages"]): await queue.put(event) await queue.put(None) # 结束信号 producer_task = asyncio.create_task(agent_producer()) while True: try: event = await asyncio.wait_for(queue.get(), timeout=15) except asyncio.TimeoutError: # 15 秒没有任何数据,发心跳保活 yield ": ping\n\n" continue if event is None: break yield f"event: {event['type']}\ndata: {json.dumps(event['data'], ensure_ascii=False)}\n\n" producer_task.cancel()

核心是wait_for(queue.get(), timeout=15)。如果 Agent 流 15 秒内没有产出任何帧,就补一个心跳帧给代理层,告诉它连接还活着。这个方案比调大代理超时更优雅,它不掩盖"模型确实在思考"这件事,只向中间层证明连接健康状况良好。

2.4 其他必须一起改的配置

心跳能解决大部分 idle timeout,但如果你服务前面真的有 Nginx,下面三个配置建议一并检查:

location /chat { proxy_buffering off; proxy_read_timeout 300s; proxy_send_timeout 300s; }

proxy_buffering off很关键。如果不关,Nginx 会把后端发来的小包攒起来,等到一定大小或超时才一次性发给浏览器。表面现象是页面半天没反应,然后突然一大段文字全出来,完全失去流式体验。proxy_read_timeout是兜底配置,虽然心跳已经保证了连接活跃,但保险起见还是放宽。

网络上那些所有 SSE 流式项目几乎都会遇到的"connection closed"问题,八成逃不出这两种原因:要么缓冲没关,要么没有心跳。两个都处理掉之后,稳定性和体验会完全不同。

3. LangChain 三大 OutputParser 逐个上手:指令注入、解析、校验与补救

3.1 为什么不能直接 json.loads:模型输出不规整的现实

如果你只跟 OpenAI 的接口打过交道,可能觉得模型输出的 JSON 还算干净。但真实 Agent 场景里,模型通常会输出类似这样的东西:

好的,这是你要的信息: {"name": "张三", "age": 28,"address": "上海"} 希望帮助到你!

前面有语气词,后面有总结句,中间 JSON 的引号、逗号还很可能不标准。直接json.loads必崩。OutputParser 的定位就是解决这个问题的完整链条:先把"只输出 JSON,不要解释"的格式指令注入 prompt,再把模型的输出解析成程序对象。我在项目里最常拿来解决实际问题的三个 Parser 分别是JsonOutputParser、StructuredOutputParser和PydanticOutputParser。

3.2 JsonOutputParser:轻量级,适合快速原型

from langchain_core.output_parsers import JsonOutputParser json_parser = JsonOutputParser() # 注入 prompt 时使用 prompt = f"请抽取用户信息并输出。\n{json_parser.get_format_instructions()}" response = llm.invoke(prompt) parsed_data = json_parser.parse(response.content)

这个 Parser 的约束很弱,它只要求模型输出 JSON,然后尝试解析成 Pythondict。不会校验字段名、不会校验类型。所以它适合的场景是:我只想快速拿到一个非结构化的字典,字段本来就是动态的,或者下游能容忍脏数据。

实测要注意,get_format_instructions()生成的内容会告诉模型"输出应为 JSON 对象",但如果你不额外在 system prompt 里强调"不要 Markdown 代码块、不要解释文字",它偶尔还是会包一层 ```json 围栏。真要规避,可以在 prompt 末尾加一句:直接输出 JSON,不要使用代码块。

3.3 StructuredOutputParser:字段级约束,约束适中

from langchain.output_parsers import StructuredOutputParser, ResponseSchema response_schemas = [ ResponseSchema(name="name", description="用户姓名"), ResponseSchema(name="age", description="用户年龄", type="integer"), ] parser = StructuredOutputParser.from_response_schemas(response_schemas) format_instructions = parser.get_format_instructions()

它的价值在于把字段清单和描述交给模型,让模型严格按字段列表组织 JSON。相对于JsonOutputParser,它多了一层字段存在性约束:模型不太可能漏掉age,因为格式指令里明确列了。但是它的校验也就到此为止,它不会告诉你age是不是真的 int,更不支持嵌套模型。一个极其平常但容易踩的坑是解析结果默认是dict,你拿到后还得自己做类型转换。

3.4 PydanticOutputParser:强校验、嵌套模型,真正的类型安全

from langchain_core.output_parsers import PydanticOutputParser from pydantic import BaseModel, Field from typing import Optional class UserInfo(BaseModel): name: str = Field(description="用户姓名") age: int = Field(description="用户年龄") address: Optional[str] = None parser = PydanticOutputParser(pydantic_object=UserInfo) format_instructions = parser.get_format_instructions() response = llm.invoke(prompt + format_instructions) user_info = parser.parse(response.content) # user_info.name, user_info.age 的类型是确定的

这是我在生产环境最常用的一个。优势有三个。第一,格式指令直接从 Pydantic 模型生成,字段描述来自Field(description=...),字段必填性来自类型注解,模型不太容易跑偏。第二,解析结果是真正的 Pydantic 对象,字段访问、类型转换、model_dump()都可靠,下游如果要把参数传给工具函数或写数据库,能省掉一堆防御式代码。第三,支持嵌套模型,结构化程度可以做得比较高。

代价是它比前两个更"重"。模型输出只要有一处不符合 schema,解析就会抛OutputParserException。在流式场景里,这种报错如果不预判,用户只会看到对话戛然而止。

3.5 三类 Parser 选型对比

Parser约束强度校验能力输出类型推荐场景
JsonOutputParser低仅 JSON 合法性dict快速原型、字段动态
StructuredOutputParser中字段存在性dict简单平面字段
PydanticOutputParser高类型、必填、嵌套Pydantic 对象强类型、生产环境

一个实用的建议:不要把 Parser 绑定在某个固定的llm.invoke()上。我习惯把get_format_instructions()的产物打印出来看一遍,确认模型到底收到了什么。很多时候你以为自己在约束模型,实际格式指令已经随版本变了。

3.6 解析失败的补救:OutputFixingParser 与重试的边界

再强的格式指令也防不住模型偶尔发疯。LangChain 提供了两个补救工具,我都实际用过:

from langchain.output_parsers import OutputFixingParser fixing_parser = OutputFixingParser.from_llm(parser=parser, llm=llm) result = fixing_parser.parse(original_output)

OutputFixingParser的原理是把原始错误信息和解析失败的文本一起丢给 LLM,让模型帮忙修复,再重新解析。RetryOutputParser则是在保持 prompt 上下文的情况下让模型重新生成一次。两者都会额外消耗 token,而且有进入死循环的风险。我的经验是:在流式 Agent 流程里,解析失败不要盲目重试超过两次。第三次还失败,就把它降级成一个普通的error事件发给前端,同时在日志里记录原始文本和错误详情。一次失败的代价远小于一次无限循环的代价。

4. ToolCall 方案实战:让 Agent 的"调用意图"以结构化事件流过链路

4.1 ToolCall 和"模型说一段话"的本质区别

传统思路是让模型在文本里提一句"我要查询北京的天气了",然后代码去字符串里正则匹配。这在 demo 里能跑,生产环境就是灾难。ToolCall 的本质区别在于:模型输出的不是自然语言,而是一个结构化的调用意图声明,类似{"name": "get_weather", "arguments": {"city": "北京"}}。这个声明由代码来校验、执行,执行结果作为新的消息继续喂给模型。

LangChain 里定义工具非常简单,基于类型注解自动生成 schema:

from langchain_core.tools import tool @tool def get_weather(city: str) -> str: """获取指定城市当前的天气情况。 Args: city: 城市名称,例如“北京”“上海”。 """ # 这里可以调真实天气 API return "晴,26 摄氏度"

然后绑定给模型:

llm = ChatOpenAI(model="gpt-4o", temperature=0) llm_with_tools = llm.bind_tools([get_weather])

工具的描述和参数说明特别重要。模型就是靠这些文本理解什么时候该调用什么工具的。city参数如果只写"城市名",模型可能把"北京市"和"北京"混着传;描述里明确写了"例如北京、上海"之后,参数的规范性明显提高。

4.2 流式场景下 ToolCall 事件的聚合

真正麻烦的是流式。OpenAI 兼容接口在 stream 模式下,工具调用参数是一个一个片段返回的,而不是一个完整的 JSON 一起到达。AIMessageChunk的tool_call_chunks字段就是这个片段集合。我在代码里做了这样一段聚合:

from langchain_core.messages import AIMessageChunk tool_calls_draft = {} async for chunk in llm_with_tools.astream(messages): if chunk.tool_call_chunks: for tc in chunk.tool_call_chunks: index = tc["index"] draft = tool_calls_draft.setdefault( index, {"name": "", "args": ""} ) draft["name"] += tc.get("name") or "" draft["args"] += tc.get("args") or "" # 如果有普通文本内容,也走 token 事件 if chunk.content: yield {"type": "token", "data": {"text": chunk.content}}

流式工具调用的参数片段可能是零散的 JSON token,凑在一起才能成为可解析的完整 JSON。这段代码的关键是把同 index 的 chunk 按顺序拼接起来,而不是每到一个 chunk 就去json.loads。

在生产链路里,我会把这些聚合完的tool_calls_draft再交给 Pydantic 做一次参数校验,确认city是字符串、天数没传成负数。这一步是双保险:模型生成的参数虽然是 JSON,但类型错误和缺字段在真实场景经常发生。

4.3 工具执行结果如何回流到 Agent

这里有个容易漏掉的闭环:工具执行完以后,一定要把结果包成ToolMessage喂回给模型,而且错误也要回流。

from langchain_core.messages import ToolMessage tool_result = None try: tool_result = get_weather.invoke(arguments) status = "success" except Exception as e: tool_result = f"执行失败:{str(e)}" status = "error" messages.append(ToolMessage(content=str(tool_result), tool_call_id=tool_call_id)) yield {"type": "tool_end", "data": {"name": tool_name, "status": status, "result": tool_result}}

这里有个很反直觉的教训:get_weather.invoke(arguments)接收的是一个 dict,这个 dict 来自上一步聚合后的解析结果。如果参数解析失败,你只发一个标准error事件给前端而不同步喂回模型,模型就会处于"工具调用了但没有任何反馈"的悬空状态。正确做法是把错误内容放进ToolMessage,让模型知道刚才失败了,这样它可以自主决定换个参数重试或者承认失败。Agent 的自我纠错能力完全依赖这个回流设计。

4.4 OutputParser 和 ToolCall 到底是互补关系

经常有人问:有了 ToolCall 主动输出结构化参数,还需要 OutputParser 吗?我的答案是需要,但职责不同。

ToolCall 的价值在参数级别:模型声明要调用哪个工具、传什么参数,这是厂商协议层面的结构化。但它不负责业务模型级别的强约束,不保证UserInfo模型所有字段完整。所以实际链路是:ToolCall 先把模型的意图结构化,拿到参数 JSON 后,再交给 PydanticOutputParser 校验和转换成具体业务对象。前者解决"模型想干什么",后者解决"下游代码安全接收什么"。两者是上下游配合,不是二选一。

LangChain 目前还推荐with_structured_output(),它内部会利用 provider 原生的 tool calling 或 JSON mode 来实现结构化输出,在模型支持的情况下比 prompt 式 Parser 更稳定。你可以把它理解成"给模型加载原生结构化协议",而 Parser 是给没有这种能力的模型兜底的通用方案。

5. Vue 端 SSE 解析封装:把事件流映射成可渲染的业务状态

5.1 用 fetch 还是 EventSource

前端对接 SSE 有一道选择题。原生EventSource非常简单,自带断线重连,但它只支持 GET,不能自定义请求头,也没法通过 POST 带复杂参数。我的接口是POST /chat,还要带 Authorization 头,所以直接淘汰了 EventSource。

兼容库方面,我推荐@microsoft/fetch-event-source。它是基于 fetch 的封装,支持 POST、自定义 header、自定义事件名解析,且事件解析逻辑已经处理好。如果你不想引依赖,我也实现过一个手写解析器,代码量不大但坑不少。

5.2 手写 SSE 帧解析器:注意 UTF-8 截断

核心是维护一个 buffer,按空行切分事件帧。下面是一个最小实现:

export type SSEEvent = { event?: string; data?: string; }; export function createSSEParser(onEvent: (e: SSEEvent) => void) { let buffer = ''; const decoder = new TextDecoder('utf-8'); return (chunk: Uint8Array) => { buffer += decoder.decode(chunk, { stream: true }); const frames = buffer.split(/\r?\n\r?\n/); buffer = frames.pop() ?? ''; for (const frame of frames) { const event = parseFrame(frame); if (event.data !== undefined) { onEvent(event); } } }; } function parseFrame(frame: string): SSEEvent { let event: SSEEvent = {}; for (const line of frame.split(/\r?\n/)) { const colonIndex = line.indexOf(':'); if (colonIndex < 0) continue; const field = line.slice(0, colonIndex); const value = line.slice(colonIndex + 1).replace(/^ /, ''); if (field === 'event') event.event = value; if (field === 'data') event.data = value; } return event; }

这里最容易翻车的地方是TextDecoder。如果你不开启stream: true,当中文字符被 TCP 分包切成两半时,解码结果会出现替换字符,后续 JSON.parse 必然失败。开启 stream 模式后,解码器会缓存不完整的字节序列,等下一段数据来了再完成拼接。这个坑藏得很深,不是流式项目真的很难碰到。

5.3 把流式事件映射成前端业务状态

在 Vue 里我封装了一个简单的组合式函数,对外只暴露渲染状态,不暴露底层协议:

export function useChatSSE() { const messageText = ref(''); const toolCalls = ref<Array<{ name: string; status: string; result?: string }>>([]); const status = ref<'idle' | 'streaming' | 'done' | 'error'>('idle'); function handleSSEChunk(chunk: Uint8Array) { parser(chunk, (event) => { if (event.event === 'token' && event.data) { const payload = JSON.parse(event.data); messageText.value += payload.text; } else if (event.event === 'tool_start') { toolCalls.value.push(JSON.parse(event.data)); } else if (event.event === 'tool_end') { // 按 id 更新对应工具卡片状态 updateToolStatus(JSON.parse(event.data)); } else if (event.event === 'done') { status.value = 'done'; } else if (event.event === 'error') { status.value = 'error'; } }); } return { messageText, toolCalls, status, send }; }

业务侧只关心文本内容、工具调用状态和整体状态。token只是拼接,tool_start往卡片列表里加一项,tool_end更新卡片内容。前端永远不做 JSON 的整体解析,不关心模型中间态,这是它能稳定工作的关键。

5.4 断连重连与指数退避

当初纯手写实现时,最痛的就是 Markdown 渲染和工具调用顺序错乱。后来发现不是前端的问题,是断连后重连没有恢复上下文导致的。SSE 是无状态的,断线后你要自己决定"重连后要不要把已收到的消息历史重新发给服务端"。

我目前的策略:断线先判断是否已经收到过任何done事件,收到就不再重连。没收到就保留页面上已渲染的部分文本,然后指数退避重试:0.5 秒、1 秒、2 秒、4 秒,最多 8 秒。重连成功后,会把已确认消费的消息列表作为上下文重新发送,避免模型重复生成已渲染的内容。同时页面上提示用户"网络不稳定,正在尝试重连"。这一步不能说解决了所有问题,但至少不会出现整段文本重复渲染的尴尬。

调 SS E 流式遇到断连和乱序问题时,还有一个简单好用的调试手法:浏览器里打开 Network 面板,看响应是不是真的逐 chunk 到达。如果看到请求在长时间等待后一次性返回全部内容,那就是缓冲层的问题;如果确认是分块到达但页面依然空白,那问题大概率在前端解析逻辑上。

这篇文章写下来,最大的体会是:SSE 流式、OutputParser、ToolCall 这三样东西单独拿出来都不难,难的是把它们放在一条链路上串起来时不互相打架。我把 TextDecoder stream、心跳帧、tool_call 参数聚合、解析失败降级这几件事在代码里都固定下来之后,这套链路才真正变得可靠。对我个人来说,最值钱的经验是:不要试图让模型在任何流式中间态输出可解析的完整结构,先把事件协议定稳,再用 Parser 和 ToolCall 管好内容的形状。后续如果要把这套方案扩展成多 Agent 协作,也只需要在事件类型里增加agent_update这一类新事件,链路本身不用动。

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

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

立即咨询