ARTICLE DETAIL

资讯详情

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

SSE流式到结构化输出:LangChain解析器与ToolCall实战

SSE流式到结构化输出:LangChain解析器与ToolCall实战 开头先说实话干 LangChain 应用最烦的不是让模型说对一句话而是怎么让这句话以一种“前端看得懂、后端不炸锅”的形式到达用户手里。SSE 流式、OutputParser、ToolCall 这三个词几乎每个接近生产的 LangChain 项目都会碰到。标题写的是“从 SSE 流式到结构化输出”实际上问的是三件事模型输出怎么增量推送怎么把增量拼成结构怎么让工具调用在流式过程里不迷路。这不是一篇讲概念的帖子而是把我在 FastAPI Vue 项目里封装的整套 SSE 流式接口逻辑、流式消息解析、以及围绕 LangChain 三大 OutputParser 和 ToolCall 的实战记录整理出来。适合正在做 AI Agent 后端、或者被前端“流式输出”需求折磨过的人参考。1. 先看整条链路SSE 在 LangChain 应用里扮演什么角色1.1 为什么前端非要 SSE一次 HTTP 连接里拿到增量数据不少第一次做聊天式应用的人会问WebSocket 不是挺好吗为什么到了 LangChain 系项目里大家默认都在聊 SSE核心原因是简单且够用。WebSocket 是双向全双工协议但在绝大多数 LLM 应用场景下前端只需要“单向接收”后端源源不断推过来的 token。SSEServer-Sent Events本质上就是一次 HTTP 响应但服务端可以持续往连接里写数据客户端用 EventSource 或 fetch 流式读取即可。这意味着它可以复用现有的 HTTP 基础设施、Nginx、负载均衡不需要额外维护长连接状态。在 LangChain 项目里SSE 通常承担两件事一是把stream出来的 token 增量实时推给前端二是把 ToolCall 的开始、结束事件也推给前端渲染“正在调用工具”的状态。后者在 Agent 场景尤其重要否则用户看着界面几十秒没反应会怀疑服务挂了。我自己在项目里的链路是FastAPI 接收 POST 请求 → 创建 LangChain AgentExecutor → 通过astream_events或自定义 AsyncCallbackHandler 捕获事件 → 封装成 SSE 格式写入 Response。1.2 从 stream 事件目录说起LangChain 回调事件与 SSE 格式的映射LangChain 在流式场景下的关键 API 是astream_events。它会产出一系列事件对象每个事件带有event字段常见的有事件名含义前端建议处理方式on_chat_model_start模型开始生成清理输入框状态展示“思考中”on_chat_model_stream收到一个 token 增量追加到当前消息文本on_tool_start工具开始执行展示“正在调用搜索/计算”等提示on_tool_end工具执行结束隐藏工具提示继续等待模型回复on_chain_end整个链结束发送消息结束事件把这些事件映射成 SSE 的自定义数据结构我习惯用一个统一的 JSON 格式{type: token, content: 你} {type: tool_start, name: arxiv, input: LangChain agents} {type: tool_end, name: arxiv, output: Loaded 3 papers} {type: done, content: }前端拿到type字段就知道该往哪个区域渲染。这里有个容易踩坑的点不要把 LangChain 底层的事件原封不动推到前端。底层事件字段多、嵌套深前端处理成本高还会暴露内部逻辑。统一封装一层生成对前端友好的格式后面维护会轻松很多。2. 三大 OutputParser 逐一拆解2.1 CommaSeparatedListOutputParser最轻量的数组输出先说最省事的。如果模型只需要输出一个列表比如生成三个关键词、五条建议CommaSeparatedListOutputParser就够用。from langchain_core.output_parsers import CommaSeparatedListOutputParser parser CommaSeparatedListOutputParser() result await parser.aparse(苹果, 香蕉, 橙子) # [苹果, 香蕉, 橙子]它本质上就是按逗号拆分再去除空白。但实际模型输出往往不是规规矩矩的“item1, item2, item3”可能带编号、换行、甚至“中文逗号”。所以我用到这个解析器时都会先加一个 prompt 约束prompt ChatPromptTemplate.from_messages([ (system, 你只输出英文逗号分隔的列表不要编号不要多余换行。), (human, {input}) ])即便如此中文逗号还是防不住所以我封装了一层预处理把和、统一替换成,再交给解析器。这种小技巧在项目里比换更复杂的解析器更实用。2.2 PydanticOutputParserJSON Schema 校验与自动修复当输出结构稍微复杂一点比如要返回一个带字段的对象PydanticOutputParser是最常见的选择。from pydantic import BaseModel, Field from langchain_core.output_parsers import PydanticOutputParser class Article(BaseModel): title: str Field(description文章标题) tags: list[str] Field(description标签列表) summary: str Field(description100字以内摘要) parser PydanticOutputParser(pydantic_objectArticle) format_instructions parser.get_format_instructions()这个get_format_instructions()很关键它会把 Pydantic 模型转换成模型能理解的 JSON Schema 说明拼在 prompt 里模型输出时才会尽量贴合结构。然后解析parsed parser.parse(raw_model_output) # 返回 Article 实例注意PydanticOutputParser 强依赖模型输出合法 JSON。遇到 DeepSeek、Qwen 这类开源模型偶尔会输出 JSON 前后带 Markdown 注释或文字说明轻则需要预处理清理重则需要走“格式修复”分支。LangChain 有OutputFixingParser可以再喂一次模型让它修正 JSON但代价是延迟和额外 token。我一般只在关键路径上用。from langchain.output_parsers import OutputFixingParser from langchain_openai import ChatOpenAI fixing_parser OutputFixingParser.from_llm( parserparser, llmChatOpenAI(modelgpt-4o-mini) )2.3 StructuredOutputParserResponseSchema不引第三方模型也能出结构如果不想定义 Pydantic 类想用更轻的字典加描述方式StructuredOutputParser配合ResponseSchema也不错。from langchain_core.output_parsers import StructuredOutputParser, ResponseSchema schema [ ResponseSchema(nameanswer, description对用户问题的直接回答), ResponseSchema(nameconfidence, description置信度0到1之间的小数), ] parser StructuredOutputParser.from_response_schemas(schema) format_instructions parser.get_format_instructions()它跟 Pydantic 解析器相比省去类定义适合响应字段不固定、或者一场对话里有多种结构的场景。缺点是校验能力弱一些只保证能取出 JSON 字段不做类型强校验。实际项目中我推荐“Pydantic 管正式协议Structured 管临时兼容”的组合。2.4 三大解析器实战选型表维度CommaSeparatedListPydanticStructuredOutput输出复杂度一维数组嵌套对象字典结构类型校验无严格Pydantic 模型弱只有字段提取prompt 指令一句话约束自动生成 JSON Schema自动生成字段说明修复能力手动替换符号可配 OutputFixingParser有限适用场景标签/关键词/清单Agent 正式协议、持久化快速原型、多场景切换选解析器时不光看“能不能解析”还要看你后续拿这个结构干什么。如果结构要入库、要参与业务逻辑强烈建议上 Pydantic让类型错误尽量在解析阶段暴露。如果只是临时展示Structured 就够了。3. ToolCall让模型不只是“说出答案”而是“带着参数调用工具”3.1 ToolCall 是什么和普通文本输出有什么区别普通输出是模型生成一段文本ToolCall 则是模型在推理过程中生成一个结构化的指令“我要调用某个工具工具名叫 search参数是 { “query”: “LangChain” }”。在 LangChain 里ai_msg.tool_calls就是这样一个列表每个元素包含name、args、id等字段。ChatOpenAI 等模型在 API 层自带 function calling / tool calling 能力所以模型直接输出的是工具调用意图而不是普通文本。区别很实际普通文本靠解析器猜结构ToolCall 是模型按训练协议直接输出结构不需要额外的 OutputParser。但它也会带来流式处理的新问题工具调用的参数是分片到达的拼不完整就无法执行。3.2 在 FastAPI 回调里捕获 on_tool_start / on_tool_end我在封装 SSE 时用的是一个继承AsyncCallbackHandler的自定义处理器。当 Agent 开始执行工具时LangChain 会产生on_tool_start结束产生on_tool_end。from langchain_core.callbacks import AsyncCallbackHandler import json class AgentSSECallback(AsyncCallbackHandler): def __init__(self, queue): self.queue queue async def on_tool_start(self, serialized, input_str, *, run_id, parent_run_id, **kwargs): await self.queue.put({ type: tool_start, name: serialized.get(name), input: input_str }) async def on_tool_end(self, output, *, run_id, parent_run_id, **kwargs): await self.queue.put({ type: tool_end, output: str(output) })这里注意input_str有时是 JSON 字符串有时是普通字符串取决于工具。我在 push 给前端前会尝试json.loads失败就原样字符串。前端展示时再统一兜底防止渲染崩溃。3.3 ToolCall 的增量事件on_llm_new_delta 中的 tool_call_chunks如果你用astream_events监听模型增量会发现on_chat_model_stream携带的chunk.content往往为空但chunk.tool_call_chunks里有内容。这是很多新手卡住的地方工具调用参数不是通过 content 过来的而是通过tool_call_chunks分片到达。async def on_chat_model_stream(self, message, **kwargs): chunk message.chunk if chunk.tool_call_chunks: for tc in chunk.tool_call_chunks: # tc 里面有 index, name, args 等字段 # args 是 JSON 片段 await self.queue.put({ type: tool_call_chunk, index: tc.get(index), name: tc.get(name) or , args_fragment: tc.get(args) or })前端收到这些分片后可以按index分组把args_fragment字符串直接拼接最后JSON.parse。这一步比听起来容易出问题模型可能分多次吐出同一个工具的参数而每次返回的index要保持一致。如果前端状态管理没按index分组后一个工具的参数会拼到前一个上面。推荐前端这样维护const toolChunks reactive({}); // index - { name, argsText } function onToolCallChunk(chunk) { const cur toolChunks[chunk.index] || (toolChunks[chunk.index] { name: chunk.name, argsText: }); cur.name cur.name || chunk.name; cur.argsText chunk.args_fragment; }等整个执行流结束后再遍历toolChunks逐个JSON.parse(cur.argsText)就能拿到完整的工具调用参数。这一步就是“封装 SSE 流式接口调用逻辑、完成流式消息解析”的核心。4. 从 SSE 流式到结构化输出的完整实战4.1 后端FastAPI LangChain 实现一个带 ToolCall 的 SSE 接口我直接给一个最小可复用的 FastAPI 接口。它创建带工具调用的 Agent通过astream_events捕获事件转为 SSE 推送。from fastapi import FastAPI from fastapi.responses import StreamingResponse from langchain_openai import ChatOpenAI from langchain.agents import create_tool_calling_agent, AgentExecutor from langchain.tools import tool from langchain_core.prompts import ChatPromptTemplate import json, asyncio app FastAPI() tool def search_web(query: str) - str: 模拟搜索工具返回结果 return f关于 {query} 的搜索结果这是一条模拟数据。 async def event_generator(user_input: str): llm ChatOpenAI(modelgpt-4o-mini, temperature0) prompt ChatPromptTemplate.from_messages([ (system, 你是一个乐于助人的助手。), (human, {input}), (placeholder, {agent_scratchpad}), ]) agent create_tool_calling_agent(llm, [search_web], prompt) executor AgentExecutor(agentagent, tools[search_web]) async for event in executor.astream_events({input: user_input}, versionv2): if event[event] on_chat_model_stream: chunk event[data][chunk] if chunk.content: yield fdata: {json.dumps({type: token, content: chunk.content}, ensure_asciiFalse)}\n\n if chunk.tool_call_chunks: for tc in chunk.tool_call_chunks: yield fdata: {json.dumps({type: tool_call_chunk, index: tc.get(index), name: tc.get(name) or , args_fragment: tc.get(args) or }, ensure_asciiFalse)}\n\n elif event[event] on_tool_start: yield fdata: {json.dumps({type: tool_start, name: event[name]}, ensure_asciiFalse)}\n\n elif event[event] on_tool_end: yield fdata: {json.dumps({type: tool_end, name: event[name]}, ensure_asciiFalse)}\n\n yield fdata: {json.dumps({type: done}, ensure_asciiFalse)}\n\n app.post(/chat/stream) async def chat_stream(payload: dict): user_input payload.get(input, ) return StreamingResponse(event_generator(user_input), media_typetext/event-stream)注意几个要点versionv2必须显式传入LangChain 新版astream_events默认版本不同事件字段有差异。media_type必须是text/event-stream否则前端 EventSource 只认 HTTP 头Content-Type: text/event-stream。用ensure_asciiFalse保证中文直接以 UTF-8 输出否则推送的是\uXXXX前端看到的就是 Unicode 转义。4.2 前端Vue 里读取 SSE 流并组装 ToolCall前端我用 Vue 3 fetch 实现。为什么不直接用EventSource因为EventSource只支持 GET不支持 POST 带 body。很多聊天场景需要把上下文、参数 POST 给后端所以我用fetch加流式读取。async function streamChat(input: string) { const resp await fetch(/api/chat/stream, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ input }) }); const reader resp.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 }); // 按 SSE 规范事件之间用空行分隔 const lines buffer.split(\n\n); buffer lines.pop() ?? ; for (const part of lines) { const payload part.replace(/^data: /, ).trim(); if (!payload) continue; const msg JSON.parse(payload); handleSSEMessage(msg); } } } function handleSSEMessage(msg: any) { if (msg.type token) { currentContent.value msg.content; } else if (msg.type tool_call_chunk) { const chunk toolChunks[msg.index] || (toolChunks[msg.index] { name: msg.name, argsText: }); chunk.name chunk.name || msg.name; chunk.argsText msg.args_fragment; } else if (msg.type tool_start) { toolStatus.value 正在调用: ${msg.name}; } else if (msg.type tool_end) { toolStatus.value ; } else if (msg.type done) { isStreaming.value false; } }这段代码里有个大坑不要用resp.text()一次性读取整个流因为 SSE 是持续推送的text()会一直等到连接关闭才能拿到结果等于没流式。4.3 结构化输出如何兼容流式部分 JSON 与增量解析问题很多人的疑问是PydanticOutputParser 不是一次性解析完整输出吗那它跟流式怎么兼容实测下来方案有两种一种是前端攒完整后后端再解析也就是流式期间前端只展示文本等done事件后把完整文本交给后端解析成结构化 JSON。这个方案实现简单但体验上结构化字段出现得晚且多一次请求。另一种是前端做增量 JSON 解析。先让模型按 Pydantic JSON Schema 输出前端在流式过程中持续对 buffer 做JSON.parse尝试成功就提前渲染出结构化面板失败则继续等待。function tryParsePartial(buf: string) { try { const obj JSON.parse(buf); if (obj obj.title) { // 可以提前渲染标题了 } } catch { // JSON 不完整正常现象 } }这个方案有个附带问题频繁JSON.parse大文本会有性能开销。我的做法是节流每 500ms 或每 20 个 token 才尝试解析一次避免每次流式 chunk 都跑一次完整 JSON.parse。5. 实际踩坑SSE 断流、空闲超时与消息解析5.1 流在完成前断开idle timeout waiting for SSE我在生产环境最常碰到的报错就是标题里那句stream disconnected before completion: idle timeout waiting for SSE。表面上是“连接空闲超时”实际上分两种情况第一种是后端长时间没有输出任何数据。比如 Agent 在顿住思考、或者工具执行很慢SSE 连接没有新数据网关认为“空闲”就断开了。解决思路是加心跳即使没有 token 产出服务端也要定时推一个注释行或type: heartbeat事件。async def event_generator(user_input): async with asyncio.TaskGroup() as tg: task_deliver tg.create_task(deliver_events()) task_heartbeat tg.create_task(heartbeat_loop()) ... async def heartbeat_loop(): while True: await asyncio.sleep(15) yield : heartbeat\n\n注意 SSE 规范中以:开头的行是注释前端EventSource会自动忽略。如果用fetch流式读取也建议在后端跳过去。第二种是网关代理层有超时限制。常见的 Nginx 默认proxy_read_timeout是 60 秒如果后端超过 60 秒不发数据连接就会被切断。需要在配置里调大location /chat/stream { proxy_pass http://backend; proxy_buffering off; proxy_cache off; proxy_read_timeout 300s; proxy_send_timeout 300s; chunked_transfer_encoding on; }proxy_buffering off尤其重要否则 Nginx 会把 SSE 响应缓冲到一定大小才转发前端感受到的不是流式而是“一阵一阵地吐文字”。5.2 前端 EventSource 的自动重连与心跳兜底用EventSource的好处是自带重连机制。但它默认重连延迟是 3 秒且用 POST 的业务场景不能直接用所以如果用fetch就要自己实现断线重连。我在项目里的做法是发现reader.read()抛错或done提前触发就加指数退避重连最多重试 5 次。同时前端也做一层“最近消息超过 30 秒没有 token 就显示等待提示”避免用户以为卡死。let retryDelay 1000; async function connectWithRetry(input: string) { try { await streamChat(input); } catch (e) { if (retryDelay 32000) return; setTimeout(() connectWithRetry(input), retryDelay); retryDelay * 2; } }5.3 中文与 JSON 特殊字符导致的解析失败流式消息解析还有一个隐藏问题JSON 里的引号、换行、控制字符。模型输出 JSON 时如果某个字段包含换行标准 JSON 里应该\n但模型有时直接输出真实换行。前端JSON.parse就会直接报错。我后端处理这个问题的兜底方案import re def sanitize_json_text(text: str) - str: text re.sub(r[\x00-\x1f\x7f], lambda m: f\\u{ord(m.group()):04x}, text) return text这个函数会把所有控制字符转成\uXXXX保证 JSON.parse 不炸。但这个方案只适用于“已截取到完整 JSON 后”的清洗不能直接对半截 JSON 用因为半截就转义会导致结构错乱。5.4 故障排查速查表现象可能原因处理建议前端只收到部分 token 后断开网关 idle timeout加心跳、调大 proxy_read_timeout前端收到完整数据但 UI 不滚动未正确处理 buffer 分割检查split(\n\n)逻辑ToolCall 参数为空或乱未按 index 拼接分片前端按 index 分组拼 args 字符串结构化解析频繁报错模型输出带 Markdown、换行用 Pydantic 强约束 sanitize返回内容是一堆\uXXXX后端 ensure_ascii 太勤统一ensure_asciiFalse流式接口一直转圈Nginx buffering 开启关闭 proxy_buffering事件顺序错乱多个工具并发用 run_id 或 index 做排序5.5 关于 agent 二次开发的体会标题里提到的“基于 DeerFlow 智能体做二次开发”“LangChain agent-inbox”本质上是同一个诉求把流式接口和结构化输出沉淀成可复用的中间层而不是每个 Agent 接入时重写一遍前端解析逻辑。我现在的做法是单独维护一个sse_message.py/sse_parse.ts模块里面定义好 token、tool_call_chunk、tool_start、tool_end、done 这几类消息的协议。凡是新接的后端 Agent只要按这个协议推消息前端零改动。凡是新接的前端页面只要引入这个解析模块也能直接消费任意符合协议的 SSE 流。对 LangChain agent 的二次开发我自己的经验是优先把 ToolCall 相关事件作为一等公民处理。很多 Agent 调试困难就是因为只看了最终文本看不到中间工具调用链。把 tool_start / tool_end / tool_call_chunk 全部透出到前端 Debug 面板后Agent 的决策路径一目了然定位问题能快很多。项目做到后期你会发现“SSE 流式”和“结构化输出”并不是两个孤立主题它们是同一条流水线上的两个环节模型在流式地生成文字同时也流式地生成结构你需要一个能同时承载两者的通道让前端既不丢增量又能拼出可执行的 ToolCall。把这层想清楚LangChain 这套体系才算真正上手。
返回列表