ARTICLE DETAIL

资讯详情

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

SSE与LangChain实战:AI对话流式输出与结构化JSON解析

SSE与LangChain实战:AI对话流式输出与结构化JSON解析 做 AI 应用开发尤其是涉及到 LangChain 和流式对话的产品你会发现所有体验的根基都压在一个词上SSEServer-Sent Events。无论是打字机效果、实时日志推送还是结构化输出的增量解析背后都是这条看似简单的 HTTP 长连接在支撑。我最早接到一个需求——让大模型回答像 ChatGPT 那样逐字输出同时把返回结果里的关键字段抽成 JSON 给前端做结构化渲染。项目做完后回头看这条链路里能踩的坑几乎全踩了一遍SSE 格式写错导致前端无法解析、LangChain 结构化输出和流式模式互相打架、流式 JSON 在途中断开导致整段结果丢失。这篇文章把我验证过的完整方案整理出来覆盖服务端 FastAPI、LangChain 输出约束、前端打字机渲染和 JSON 解析适合正在搭建 AI Agent 对话产品或准备对智能体平台做二次开发的同学参考。1. 为什么 AI 对话应用绕不开 SSE流式协议的核心机制1.1 从一次请求-响应说起传统 HTTP 接口的工作方式是一条完成即交付的流水线浏览器发起请求后服务端把整段响应体攒齐一次性返回。这套模型做普通 API 没有任何问题但放到大模型对话场景里就非常别扭——一个稍微复杂一点的回答可能要生成几秒钟甚至几十秒如果让用户盯着一个 spinner 干等体验很糟糕。更关键的是大模型本身就是按照 token可以粗略理解为一个字或一个词逐个生成的。既然生成是流式的传输自然也应该跟着流式走。ChatGPT 之所以让人觉得它在思考、在打字并不是什么前端动画技巧而是服务端真的把 token 一个一个推送了下来。SSE 解决的就是这个服务端主动推送的问题。它基于普通的 HTTP 协议只是约定了一种特殊的数据格式服务端把响应的 Content-Type 设为text/event-stream然后持续向连接中写入事件流客户端一边收一边渲染。整个过程不需要额外的协议升级也没有 WebSocket 那么重的握手开销。我第一次实现时也以为这只是前端接一个接口的事真正做下去才发现这玩意儿有几个容易忽略的约束一是默认只支持服务端到客户端的单向推送双向通信得靠客户端另外发请求二是它建立在长连接上中间任何一层代理Nginx、网关、负载均衡都可能因为超时策略断开这条连接这就是后文会详聊的 idle timeout 问题。1.2 SSE 的数据格式与事件语义SSE 的格式比大多数人想象中简单。每条消息由若干字段行组成字段之间用\n分隔消息之间必须有一个空行也就是连续的\n\n作为结束标记。最基本的字段是datadata: {content: 你好} data: {content: }一个事件流里可以包含多种字段data:消息体内容可以有多行多行会被拼接成一个字符串event:事件类型默认值为message客户端可以用addEventListener监听不同事件id:事件 ID客户端断线重连时会带上这个 ID 请求增量数据retry:断线后客户端重连的等待毫秒数: comment以冒号开头的注释行通常用作心跳包客户端解析时会自动忽略字段开头不允许有空格冒号后面如果跟空格会被当作数据内容的一部分处理。这是很多新手容易犯的格式错误——我在本地调试时也遇到过自己手动拼的 SSE 字符串里data:后面多打了一个空格结果前端读到的 JSON 带了一个空格前缀JSON.parse直接报错。还有一个关键点事件流里传输的通常是一段文本而不是浏览器自动解析好的数据对象。也就是说前端拿到data字段后永远要自己做一次解析。这也引出了后面 JSON 解析的一系列问题。1.3 SSE 和 WebSocket 怎么选在技术选型上很多团队会纠结SSE 还是 WebSocket。我的建议很简单如果只是服务端单向推送数据比如大模型 token 流、日志流、通知推送优先选 SSE如果有双向交互需求比如聊天室、实时协作文档、在线游戏才需要考虑 WebSocket。原因也很实在。SSE 基于普通 HTTP可以直接复用现有的认证体系Cookie、Token、网关鉴权中间走 Nginx 也不需要额外做 Upgrade 配置。WebSocket 虽然也是 HTTP 升级而来但它会占用独立的连接通道Nginx 默认配置下 60 秒无数据也会断开反而更容易踩坑。SSE 还有一个 WebSocket 没有的优势自带断线重连机制。浏览器端的EventSource对象在连接断开后会根据服务端返回的retry字段自动重连并自动携带Last-Event-ID头。WebSocket 要自己实现心跳和重连逻辑复杂度高出一截。不过在 LangChain 对话场景里我们往往需要 POST 请求携带上下文参数EventSource 只支持 GET所以更常见的做法是使用fetch配合ReadableStream手动解析。这个方案我会在第三章详细展开。2. LangChain 结构化输出把自由文本变成可用的 JSON2.1 结构化输出到底解决什么问题如果你只是做一个纯聊天机器人模型的自然语言回复本身就是最终产物不需要额外加工。但实际业务里AI 的输出往往还要喂给下游系统从一段客服回复中抽取用户情绪标签和工单类型从一段分析报告中提取风险点和建议操作或者让 Agent 输出一个完整的任务计划供前端渲染成步骤卡片。这些都要求模型输出严格符合一个预定 JSON Schema。早期做法是在 prompt 里写请以 JSON 格式返回字段包括 xxx模型大部分时候能配合但只要字段一多、嵌套一深总会出现漏字段、多逗号、引号未闭合这类错误。与其在 prompt 层面碰运气不如用 LangChain 提供的能力把输出约束变成模型层面的强制行为。2.2 with_structured_output 的内部机制LangChain 的with_structured_output是对结构化输出能力的一层封装。用法很直白先定义一个 Pydantic 模型再把它传给链或模型from pydantic import BaseModel, Field from langchain_openai import ChatOpenAI class TaskPlan(BaseModel): title: str Field(description任务标题) steps: list[str] Field(description执行步骤列表) priority: int Field(description优先级1-5数值越高越优先) llm ChatOpenAI(modelgpt-4o, temperature0) structured_llm llm.with_structured_output(TaskPlan) result structured_llm.invoke(帮我规划一个周末学习计划) print(result.model_dump())这段代码背后发生了几件关键的事情第一LangChain 会把TaskPlan转换成 JSON Schema塞进模型接口的tools参数里。如果底层模型支持 function calling / tool callingOpenAI 系、Claude 系、Qwen 系都支持模型就会被强制按照这个 Schema 生成输出LangChain 拿到响应后自动解析成 Pydantic 对象。第二如果底层模型不支持 tool callingLangChain 会走一条降级路线把 JSON Schema 写进系统提示词再用输出解析器兜底。这条路线的稳定性完全取决于模型能力所以我会在选择基础模型时优先确认它是否原生支持 tool calling。第三带with_structured_output的调用默认是非流式的。模型要等所有 token 生成完、解析成对象后一次返回。这跟对话场景需要的打字机效果天然冲突。2.3 流式场景下结构化输出的取舍回到实战场景我需要打字机效果又需要结构化字段。这时候有两条路可以走。第一条路是双向奔赴——在同一套模型调用里既拿流式文本又拿结构化结果。LangChain 在高版本里对部分模型支持流式结构化输出但实际用下来兼容性参差不齐API 还在演进中。我在一个 Agent 项目里试过用astream配合with_structured_output有的模型能正常返回增量 JSON 片段有的模型直接报错说 stream mode 与 structured output 不兼容。这种不确定性在生产环境里很难接受。第二条路是拆分诉求——流式文本走一条链路结构化字段走另一条链路。具体做法是先用普通流式接口让模型逐 token 输出对话文本同时让模型在文本中用特殊标记包裹一个 JSON 块前端在流式渲染文本的同时从增量数据中提取 JSON 块做解析。更稳妥的设计会在第四章讲。这里有个值得记下的经验结构化输出和流式输出本质上是两个需求。结构化的核心诉求是严格可信流式的核心诉求是实时可见两者不要强行糅在一个接口里。与其纠结 LangChain 的流式结构化 API不如把数据协议设计成增量事件流 最终结果事件各取所需。这个思路在对接 LangGraph 或 DeerFlow 这类智能体框架时尤其管用。3. 打字机效果全链路FastAPI 流式接口 前端增量渲染3.1 服务端StreamingResponse 的正确姿势服务端我用的是 FastAPI配合 LangChain 的异步流式接口astream。一个最基础的 SSE 流式接口长这样import asyncio import json from fastapi import FastAPI from fastapi.responses import StreamingResponse from langchain_openai import ChatOpenAI app FastAPI() llm ChatOpenAI(modelgpt-4o, temperature0.7) async def event_generator(prompt: str): # 先发一个事件告知前端流开始 yield fevent: start\ndata: {json.dumps({message: 开始生成})}\n\n async for chunk in llm.astream(prompt): content chunk.content if content: # 这里把每个 token 包装成 SSE 事件 yield fdata: {json.dumps({delta: content})}\n\n # 流结束事件前端收到后可以做收尾 yield fevent: done\ndata: {json.dumps({message: 生成完成})}\n\n app.post(/chat/stream) async def chat_stream(payload: dict): prompt payload.get(prompt, ) return StreamingResponse( event_generator(prompt), media_typetext/event-stream, headers{ Cache-Control: no-cache, X-Accel-Buffering: no, Connection: keep-alive, }, )有几个细节必须交代清楚。X-Accel-Buffering: no是给 Nginx 看的。Nginx 默认会缓冲上游响应如果不关掉SSE 的数据会被攒在代理层前端收不到增量打字机效果直接失效。这个头的作用就是告诉 Nginx别缓冲往客户端直推。为什么要用event: start和event: done区分事件类型因为前端渲染逻辑需要知道边界——收到start时清空上一次的状态收到done时停止 loading 并尝试解析最终 JSON。如果所有数据都混在一个message事件里前端就不得不在数据内容里约定特殊标记容易产生歧义。还要注意别在流式生成器函数里做耗时初始化。StreamingResponse 一旦启动生成器内的代码就开始执行但如果你的生成器先花 10 秒去查数据库、调用外部认证服务那么客户端 TTFB首字节时间会很难看甚至触发前端的超时。正确做法是在生成器外先把需要的上下文、历史消息、链路 ID 都准备好生成器只负责等模型、透传 token。3.2 前端用 fetch 读取 text/event-stream浏览器端的EventSource只支持 GET 请求而聊天接口往往需要 POST 携带 prompt、角色、上下文等参数所以前端更通用的做法是直接用fetch读取流式响应async function streamChat(prompt) { const response await fetch(/chat/stream, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ prompt }), }); if (!response.ok) { throw new Error(HTTP ${response.status}); } const reader response.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 }); // SSE 消息以空行分隔按 \n\n 切割事件块 const chunks buffer.split(\n\n); buffer chunks.pop(); // 最后一段可能是不完整的留到下次拼接 for (const chunk of chunks) { handleSSEChunk(chunk); } } } function handleSSEChunk(chunk) { const lines chunk.split(\n); let eventType message; const dataLines []; for (const line of lines) { if (line.startsWith(event:)) { eventType line.slice(6).trim(); } else if (line.startsWith(data:)) { dataLines.push(line.slice(5).trimStart()); } } if (dataLines.length 0) return; const dataStr dataLines.join(\n); let data; try { data JSON.parse(dataStr); } catch (e) { console.error(SSE data 解析失败:, dataStr); return; } if (eventType start) { resetUI(); } else if (eventType done) { handleDone(data); } else if (data.delta) { appendDelta(data.delta); } }这段代码里最值得注意的是decoder.decode(value, { stream: true })。流式读取时一个 UTF-8 字符可能被拆分在两个 chunk 里尤其是中文、emoji 这类多字节字符如果不加stream: true边界上的字节被单独解码会变成乱码。另外buffer的暂存逻辑很重要——TCP 传输不保证每帧数据正好以\n\n结尾所以必须先把所有数据拼进缓冲区再按完整分隔符切割剩余未结束的部分留到下一轮。3.3 渲染层Markdown 增量渲染与性能优化拿到delta之后怎么渲染其实是个大坑。如果只是纯文本直接往 DOM 里 append 就行。但大多数对话产品用的是 Markdown 格式问题就来了Markdown 是块级语法一个代码块还没闭合、一个列表还没结束的时候如果每一帧都重新解析整段 Markdown会产生明显的闪烁和跳动如果完全不重渲染代码块的高亮又不会更新。我实测下来最稳的做法是双轨制。对话内容区用一个只读容器承载已经确认完整的 Markdown 渲染结果另用一个增量缓冲层积累最新的 1-2 个 token 原始文本。具体来说当缓冲区累计了足够长度比如 50 个字符或者遇到明显的语法边界比如换行符、代码块结束标记就把缓冲区内容拍平进完整文本重新渲染整块 Markdown在缓冲区还没拍平之前只在容器的末尾追加一小段纯文本占位避免频繁全量重渲染渲染时机上建议用requestAnimationFrame做节流。模型生成 token 的速度远高于浏览器的渲染帧率每收到一个 token 就同步更新 DOM 很容易造成页面卡顿。正确思路是把 token 先推入一个队列在下一帧统一批量提交这样 UI 的刷新节奏和浏览器渲染节奏保持一致。还有一个经验不要在流式过程中直接做代码高亮。高亮计算很消耗性能尤其是长回复。我的做法是先渲染纯 Markdown 结构流结束后再对代码块做一次高亮用户几乎察觉不到这个延迟。4. 流式 JSON 解析的实战方案从暴力等待到增量解析4.1 最简单的方案等全部输出完再解析如果聊天回复里末尾带一个 JSON 块比如以下是你需要的任务计划 { title: 学习计划, steps: [复习数学, 做物理题, 读英语] }新手最常见的写法是流式过程中只渲染文本等done事件到达后把完整文本里的 JSON 块用正则抠出来再JSON.parse。这个方案实现简单也不会出错但有两个明显问题。第一答案要等全部生成完才出现用户能明显感觉结构化内容是突然蹦出来的而不是跟随打字机效果逐步呈现。如果 JSON 块位于文本末尾用户要看完一大段文字才能看到结构化结果。第二对长输出来说流结束后一次性解析大 JSON 会有轻微卡顿体验不够顺滑。但必须承认对于一些边缘业务比如只需要最终结果、中间状态无所谓的后端任务流式展示 最终解析依然是最稳定、最值得推荐的方案。复杂的东西不一定就好稳定可靠才是第一位的。4.2 增量解析括号配对与可恢复 JSON 解析如果真的需要边接收边解析主动权其实在服务端手里。如果服务端能保证流式数据是一个正在生成中的合法 JSON那么前端可以做增量解析——每次拿到新数据后尝试解析当前缓冲区里的内容成功则更新界面失败则继续等待。这里的关键是判断当前缓冲区里的 JSON 是否已经完整。我写过一个简单的栈式判断器function isJsonLikelyComplete(str) { let depth 0; let inString false; let escaped false; for (let i 0; i str.length; i) { const ch str[i]; if (escaped) { escaped false; continue; } if (ch \\) { escaped true; continue; } if (ch ) { inString !inString; continue; } if (inString) continue; if (ch { || ch [) depth; else if (ch } || ch ]) depth--; } return depth 0 !inString str.trim().length 0; }然后每次收到新 token把增量追加到缓冲区先判断看起来完整再尝试JSON.parsefunction tryParseIncremental(accumulated) { if (!isJsonLikelyComplete(accumulated)) { return { status: pending, data: null }; } try { return { status: done, data: JSON.parse(accumulated) }; } catch (e) { // 括号配平了但 parse 失败说明中间有语法错误等待更多字符 return { status: error, data: null }; } }这个方案在实际项目中能用但要意识到模型输出的 JSON 可能存在各种脏数据转义的反斜杠、字符串里的换行、未闭合的 Unicode 字符。一旦isJsonLikelyComplete返回 true 但JSON.parse一直失败就不能死等要设置一个最大等待帧数或超时时间超时后放弃增量解析等done事件后用完整内容做最终解析兜底。4.3 工程推荐事件分离的流式协议设计做了几个项目之后我越来越倾向一种更工程化的设计——不在数据层面硬挤而是在协议层面把结构化 JSON和文本增量分离。服务端先发一个meta事件携带这次响应的整体结构说明比如有哪些字段、每个字段的展示类型然后流式发送delta事件只携带文本增量流结束后发一个done事件携带最终完整 JSON。前端的行为就变得清晰了delta只负责打字机渲染meta驱动结构化组件的骨架done驱动结构化组件的最终数据填充这种设计的最大好处是结构化解析不再依赖字符串里某个位置的 JSON 块这种脆弱的约定而是由事件类型明确保证。而且done事件里的 JSON 是服务端在拿到模型完整输出后自己解析好的保证合法性。前端即使对增量内容一无所知也能在流结束后正确展示结构化结果。我自己在给 Agent 平台做二次开发时就是按这个协议设计的工具调用信息、中间思考过程、最终回复三者的展示节奏完全不同只有事件分离才能让前端逻辑保持清爽。5. 踩坑实录idle timeout 断连问题的完整排障链路5.1 问题现象与根因分析我在项目中第一次遇到这个报错时印象很深前端控制台打出stream disconnected before completion: idle timeout waiting for sse症状是流式回复生成到一半突然中断后端看 LangChain 日志明明还在正常输出前端却已经收不到任何数据。根因其实分三层。第一层是代理层超时。几乎所有代理/网关都有空闲超时配置如果一段时间内连接上没有数据流动代理就认为这条连接已经闲置主动断开。Nginx 的proxy_read_timeout默认是 60 秒AWS 的 ALB 默认 idle timeout 也是 60 秒云厂商 API 网关的默认值普遍在 30-60 秒之间。第二层是数据生成节奏太慢。大模型生成两个 token 之间可能有明显的停顿尤其是开启 reasoning/thinking 能力时模型可能先思考几十秒再一次性输出。如果模型思考时间超过了代理的等待阈值连接就会被断开。我曾经排查过一个案例Agent 在调用工具时卡了 40 秒没输出前端直接断流用户体验就是AI 转圈转到一半突然报错。第三层是服务端代码里有同步阻塞调用。Python 的 FastAPI 是异步框架但如果你在流式生成器里用了requests.post、同步数据库查询这类阻塞操作整个事件循环会被卡住SSE 连接自然会看起来很空闲。5.2 代理层配置修复如果连接结构是客户端 - Nginx - FastAPI - LLM 服务优先处理 Nginx 这一层location /chat/stream { proxy_pass http://backend; proxy_http_version 1.1; proxy_set_header Connection ; proxy_buffering off; proxy_cache off; proxy_read_timeout 600s; proxy_send_timeout 600s; proxy_buffers 4 16k; proxy_buffer_size 8k; }proxy_buffering off是最关键的一项它会强制 Nginx 拿到上游数据立刻转发给客户端而不是攒在一起。proxy_read_timeout调大到 600 秒给 LLM 的思考停顿留足空间。如果你用的是云上负载均衡在控制台里找到 idle timeout 配置全部调大。数据库里这条经验我反复用所有链路节点——LB、CDN、Nginx、网关——都要检查一遍少一个都可能成为瓶颈。5.3 服务端保活与生成节奏优化代理层调完之后服务端自己也要做好两件事。第一件事是心跳保活。在生成器中周期性地发送注释行让连接始终有数据在流动async def event_generator(prompt: str): async def heartbeat(): while True: yield : ping\n\n await asyncio.sleep(15) # 用 asyncio 把模型生成流和心跳流合并这里用的是 SSE 的注释行做保活——客户端解析时会忽略以冒号开头的行但在网络层看来这条连接一直有数据流动就不会被判定为 idle。心跳间隔建议小于代理层超时时间的一半。第二件事更本质别让流式生成器出现长时间无输出。在接入 LangChain/Agent 时如果中间要调用工具或进行多轮推理应该在每个阶段切换时主动发一个status事件比如event: status\ndata: {stage: calling_tool}\n\n。这样前端能看到AI 正在调用工具的中间状态用户感觉系统是活跃的连接也因为有数据而保持存活。我在接入 LangGraph 时会把astream的多个事件流token、工具调用、状态更新分发给不同的 SSE 事件类型既能解决超时问题又天然实现了第四章说的事件分离协议。6. 进阶封装把 SSE 调用、解析、重连封装成可复用模块6.1 SSEStreamClient 的整体设计踩完了所有坑之后你会发现流式接入的代码其实高度重复值得封装成一个内部通用模块。我在项目中维护了一个轻量级的SSEStreamClient核心接口包括class SSEStreamClient: def __init__(self, base_url, on_eventNone, ...): self.base_url base_url self.on_event on_event self.last_event_id None self.retry_ms 3000 async def connect(self, path, payload): # 发起 POST 请求逐行解析 SSE ... async def _read_stream(self, response): # 解析缓冲区域、切割事件块、分发回调 ... async def _reconnect(self, path, payload): # 断线重连根据 Last-Event-ID 续传 ...关键设计点有三个。第一_read_stream内部封装了多字节解码、缓冲区切分、事件类型分发上层业务只注册on_event(event_type, data)回调完全不关心底层 HTTP 细节。第二重连逻辑要有指数退避。第一次断线等 1 秒第二次等 2 秒最多等 30 秒同时要规避服务端返回的retry字段覆盖这个节奏。重连时要带上Last-Event-ID头这样后端可以实现增量续传避免前端把整个长回复重新接收一遍。第三整个 client 要把网络异常和业务异常区分开。连接被断开是网络异常走重连服务端发了event: error且带错误码走业务错误回调此时重连没有意义。6.2 对接 LangGraph/DeerFlow Agent 的二次开发思路现在很多 Agent 平台包括 LangGraph 的 Agent 服务、DeerFlow 这类智能体框架已经自带了流式事件能力二次开发时最常做的事就是把这些内部事件转换成我们自定义的 SSE 协议。以 LangGraph 为例调用图时可以指定stream_modeevents graph.astream( {messages: [{role: user, content: prompt}]}, config{recursion_limit: 50}, stream_modemessages # 取 token 增量 )在这个流式迭代里每个节点执行时都会产生不同的事件。我通常的做法是在 FastAPI 层写一个适配器把 LangGraph 的事件流转换为标准 SSE 事件图开始运行时发送event: graph_start节点切换时发送event: node_status附带节点名和状态模型产出 token 时发送event: delta附带文本增量图运行结束发送event: done附带最终完整结果这样做的好处是前端只需要对接自己的 SSEStreamClient后续无论底层换成 LangGraph、DeerFlow 还是别的智能体框架前端代码几乎不用改。把框架输出和产品协议解耦是所有二次开发项目里最值得投入的一步。我对这套封装最满意的地方在于它把服务端的流式能力变成了一个统一出口所有业务模块共用同一条流式管道新接入一个 Agent 只需要写一段适配器不需要再碰前端。实际上这条经验已经帮我在多个项目里省掉了大量重复排障的时间——连接超时、缓冲区截断、重连失败这些常见问题在封装层就被兜住了业务侧基本不会感知到。
返回列表