:SSE與OutputParser協(xié)同指南)
1. 為什么流式輸出和結構化返回總是打架做 GenAI 應用做了這幾年我最大的感受是流式輸出和結構化返回天然就是一對矛盾體。大模型默認吐出來的是自然語言你要它一段一段地流式返回給前端打字機效果又要它在最后給出一個能被程序直接校驗、入庫、調用工具的 JSON 結構這中間如果不做設計十有八九會翻車。先說一個最典型的場景你用 FastAPI 起了一個 Agent 服務前端用打字機效果展示大模型的回答同時后端還要把是否調用了某個工具工具參數(shù)是什么最終結果字段有哪些解析出來用于日志審計和業(yè)務判斷。如果只簡單地把模型輸出整段塞給前端前端拿到的是一堆夾雜著廢話的純文本如果只等全部生成完再一次性返回流式體驗就沒了。SSEServer-Sent Events恰好是這兩者之間的橋——它能保證文本一塊一塊地推給前端又能在每個事件塊里攜帶結構化元數(shù)據(jù)。這篇文章我就完整復盤一下怎么用 LangChain 的三大 OutputParser 配合 ToolCall在 SSE 流式場景下既保體驗、又保結構化。這篇文章適合誰看已經在用 LangChain 寫 Agent、但被流式解析折磨過的后端工程師或者正準備把 LLM 能力封裝成標準 API 服務、又不想丟掉打字機效果的前后端同學。我會把每一步的取舍和原理都講清楚不是讓你照抄代碼而是讓你下次遇到類似需求時能自己拍板選型。2. 先搞清楚 SSE 在這個場景里到底扮演什么角色2.1 SSE 和 WebSocket為什么聊天場景常選 SSE很多人一想到實時推送就默認 WebSocket但在大模型應用里SSE 往往是更務實的選擇。SSE 是單向的服務器往客戶端推數(shù)據(jù)客戶端不需要也不應該頻繁回傳。這和 LLM 生成的場景天然匹配——用戶提問之后剩下的就是模型一直說前端一直聽。協(xié)議層面SSE 就是普通的 HTTP 響應Content-Type 設為text/event-stream然后在響應體里按固定格式寫事件event: message data: {type: token, content: 你} event: message data: {type: token, content: 好} event: done data: {type: done, session_id: abc123}每個事件之間用空行隔開data字段是真正的載荷event是事件類型id可以用來做斷點續(xù)傳retry是客戶端自動重連的時間間隔。就這么簡單。WebSocket 呢它是全雙工要握手、要維護長連接心跳、要處理斷線重連邏輯在只需要服務器單向推送的場景里屬于殺雞用牛刀。更現(xiàn)實的問題很多企業(yè)內網的網關、Nginx 配置對 WebSocket 的升級請求支持不友好但對 SSE 這種普通 HTTP 長響應基本上零成本透明轉發(fā)。我在實際項目里用 SSE 遇到過的最大坑反而很簡單網關的超時時間設太短模型思考超過 60 秒連接直接被掐斷前端就收到一個stream disconnected before completion: idle timeout waiting for sse。這個問題后面在排查章節(jié)我會詳細說。2.2 從 LangChain 的流式機制到 SSE 事件拼裝LangChain 從Runnable體系開始把流式能力統(tǒng)一成了stream/astream兩個接口。chain.stream(query)會一塊一塊地yield輸出。需要注意的是這個一塊不一定是模型吐的一個 token而是 LangChain 每個Runnable步驟產出的一個完整單元。比如Retriever步驟吐出文檔列表LLM步驟吐出 token 片段。在 Agent 場景下中間可能還夾著Tool的執(zhí)行結果。所以設計 SSE 接口時我的做法是先約定一套內部事件協(xié)議而不是直接把 token 裸推出去。每個事件我用 JSON 串里面至少帶兩個字段type和payload。event: message data: {type: start, payload: {session_id: uuid}} event: message data: {type: token, payload: {content: 你好}} event: message data: {type: tool_call, payload: {name: search_news, args: {keyword: 人工智能}}} event: message data: {type: tool_result, payload: {result: ..., duration_ms: 1200}} event: message data: {type: structured, payload: {title: ..., summary: ...}} event: message data: {type: done, payload: {finish_reason: stop}}前端只需要根據(jù)type決定怎么渲染token追加到正文tool_call可以展示正在調用工具的動畫structured存到表單里做后續(xù)業(yè)務。這樣流式體驗和結構化數(shù)據(jù)就各歸其位了。后面我講的三大 OutputParser本質都是在最后那個結構化事件這個環(huán)節(jié)里保證你拿到的payload是干凈、可校驗的 JSON。3. 三大 OutputParser把模型的話轉成程序能用的結構OutputParser 在 LangChain 里的定位是從 LLM 的原始輸出里抽取出程序需要的結構。注意它通常不是魔法它靠的是提示詞約束 格式校驗兩步走。讓模型按指定格式輸出再用解析器校驗、糾錯。理解了這一點你就能明白為什么換個模型解析失敗中文環(huán)境下 JSON 不標準這類問題會反復出現(xiàn)了。3.1 PydanticOutputParser給 JSON 上一份類型保票PydanticOutputParser是項目里最常用的一個。它結合了 Pydantic 的BaseModel把輸出格式約束成明確的字段類型——字符串、整型、列表、嵌套對象類型不對直接校驗失敗。用法上核心三步定義模型類、創(chuàng)建解析器、把格式指令塞進提示詞。from typing import List from pydantic import BaseModel, Field from langchain.output_parsers import PydanticOutputParser class ArticleSummary(BaseModel): title: str Field(description生成的文章標題不超過20字) keywords: List[str] Field(description3-5個關鍵詞) summary: str Field(description100字以內的核心摘要) confidence: float Field(description模型對摘要質量的自信度0到1之間) parser PydanticOutputParser(pydantic_objectArticleSummary) prompt PromptTemplate( template請分析下面這段文本并嚴格按照格式要求輸出。\n文本{text}\n{format_instructions}\n, input_variables[text], partial_variables{format_instructions: parser.get_format_instructions()}, )parser.get_format_instructions()生成的那段指令本質是把你定義的字段、類型、約束翻譯成自然語言模板告訴模型必須輸出一個 JSONkey 有哪些value 是什么類型。它還會補充一句不要輸出其他內容之類的強調。實際效果上模型越強GPT-4 級別、Claude 3.5遵循度越高小模型經常把 JSON 包在 Markdown 代碼塊里或者多解釋一句。解析器的parse方法還內置了糾錯能力如果模型輸出的文本可以被eval成 JSON但類型不對它會嘗試用 LLM 自動修復。不過這個修復是有損的——它需要額外調一次模型速度和成本都要考慮。所以在流式場景里我通常不在最后階段用自動修復而是做兩段式先流式展示再后端靜默校驗校驗失敗才觸發(fā)修復。3.2 StructuredOutputParser輕量到不需要定義類StructuredOutputParser適合那種不想建 Pydantic 模型只要幾個簡單字段的場景。它通過ResponseSchema列表來聲明字段名、類型和描述使用起來比 Pydantic 版更輕。from langchain.output_parsers import StructuredOutputParser, ResponseSchema response_schemas [ ResponseSchema(nameanswer, description對問題的直接回答, typestring), ResponseSchema(namesource, description答案的參考來源如果沒有則為null, typestring), ] parser StructuredOutputParser.from_response_schemas(response_schemas)從源碼實現(xiàn)看StructuredOutputParser內部并沒有把輸出嚴格轉成 Pydantic 對象而是返回一個字典。它對字段順序、缺失字段的處理更寬松底層用的其實是類似正則 字典提取的簡易邏輯。所以它適合內部接口、日志記錄、字段不多的場景一旦你的下游真的要用強類型做入?yún)⑿r炦€是 Pydantic 版更省心。這里有個經驗如果返回結構里嵌套層級很深或者字段會因為模型輸出風格波動我建議直接用 Pydantic 版不要在 Structured 版上強行造輪子。輕量方案省下的代碼量會在排障時加倍還回去。3.3 JsonOutputParser最容易被低估的那個第三個其實是JsonOutputParser—— LangChain 里專門用來只要 JSON不要類定義的解析器。它和PydanticOutputParser最大的區(qū)別是不要求目標類型是 Pydantic 模型你用普通 dict 聲明一個期望的 JSON 結構模板它就能照這個模板校驗。from langchain.output_parsers import JsonOutputParser parser JsonOutputParser() prompt PromptTemplate( template輸出JSON格式結果字段包括: title(string), items(array of string)。\n{format_instructions}\n, input_variables[], partial_variables{format_instructions: parser.get_format_instructions()}, )它做的事情是拿到模型文本提取并解析出 JSON 對象。它不會去做嚴格類型強轉保持了 dict 的靈活性。實際項目中我經常把它作為流式過程中增量解析 JSON 的工具——不是等模型完整輸出后一次性解析而是配合字節(jié)流每次拿到新的 token 片段就去試著解析能解析出部分字段就先緩存。這個思路在長回答場景里特別有用你可以在模型還在生成正文時就把標題、關鍵列表等輪廓字段提前推給前端。總結一下三者的選擇邏輯要強類型校驗、下游要嚴格入庫選 Pydantic只要幾個字段、隨拿隨用選 Structured既要 JSON 又不想綁定模型定義、或需要在流中做增量解析選 Json。沒有絕對好壞只看約束強度。4. ToolCall 方案讓模型把工具意圖直接交出來4.1 為什么不用讓模型自己拼工具調用文本早期 LangChain 的 Agent 實現(xiàn)里模型是用純文本的方式假裝調用工具——輸出一行Action: search_news\nAction Input: 人工智能然后 AgentExecutor 去解析這段文本。這種方案在模型能力弱的時候還算勉強能用但問題很明顯模型一旦在文本里多加一句解釋、少寫一個換行整個解析就崩了。而且文本格式因模型而異換模型就要調解析規(guī)則?,F(xiàn)在主流方案是ToolCall函數(shù)調用。模型在生成時除了輸出自然語言還可以輸出一個結構化的工具調用意圖——方法名、參數(shù) JSON。OpenAI 的 function calling、Claude 的 tool use、國產模型不少也兼容這個協(xié)議。LangChain 里的做法是把工具定義綁定到模型上這一類模型能力稱為bind_tools。from langchain_openai import ChatOpenAI from langchain_core.tools import tool tool def search_news(keyword: str, limit: int 5) - list: 搜索新聞資訊keyword為關鍵詞limit為返回條數(shù)。 # 這里寫真實檢索邏輯 return [{title: 示例新聞, url: https://example.com}] llm ChatOpenAI(modelgpt-4o-mini, temperature0) llm_with_tools llm.bind_tools([search_news]) response llm_with_tools.invoke(幫我搜一下今天人工智能領域的新聞)關鍵在于response是一個AIMessage如果模型決定調用工具它的tool_calls屬性里會帶上結構化調用信息response.tool_calls # [{name: search_news, args: {keyword: 人工智能, limit: 5}, id: call_xxx}]這個id字段很重要尤其是做并發(fā)工具調用時它用來關聯(lián)工具結果和對應的調用請求。4.2 ToolCall 事件在 SSE 里怎么推既然AIMessage.tool_calls是結構化的那么從流式事件角度它也能流式地分片到達——模型先生成工具名再一點一點生成參數(shù) JSON。LangChain 的astream_events可以讓你捕獲on_chat_model_stream事件從而拿到 token 級的流。但我要提醒你工具參數(shù)這種 JSON前端完全沒必要做打字機效果。你只要在tool_call開始事件里推一條正在調用工具等工具結果出來再推一條結構化結果就行參數(shù) JSON 在中間過程可以直接攢在后端。這是我的實踐結論——不要一上來就把所有 token 都推給前端做逐字渲染那樣只會讓前端渲染邏輯又復雜又容易出錯。工具執(zhí)行完你拿到的結果同樣建議包裝成結構化事件推給前端。同時把工具結果作為新的上下文消息再喂回給模型讓它基于結果生成最終回答。這個模型→工具→模型的循環(huán)如果自己用for循環(huán)寫很容易在異常分支和超時控制上出問題——這也是為什么存在 LangGraph 這類帶狀態(tài)編排的框架。不過對于單輪工具調用場景手動循環(huán)完全可控不需要上重型框架。4.3 OutputParser 和 ToolCall 怎么配合這就是這個方案的精髓了ToolCall 解決模型要調用什么工具、參數(shù)是什么的結構化提取OutputParser 解決模型最終要返回給業(yè)務的最終結論的結構化提取。兩者是在一次請求的不同階段各司其職。class FinalAnswer(BaseModel): reply: str Field(description面向用戶的最終回答) used_tools: List[str] Field(description本次實際使用到的工具名稱列表) data_source: List[str] Field(description參考信息的來源列表) final_parser PydanticOutputParser(pydantic_objectFinalAnswer)流程大致是用戶提問 → 模型決定調用工具ToolCall 結構化→ 執(zhí)行工具 → 把工具結果拼進上下文 → 模型生成最終回答OutputParser 結構化→ 通過 SSE 推送給前端。中間環(huán)節(jié)的結構化靠tool_calls最后的業(yè)務結構靠 OutputParser。兩條線涇渭分明誰也不會干擾誰。5. FastAPI LangChain 完整落地一條 SSE 接口打通全流程5.1 服務端異步流式接口的分層設計我強烈建議把模型調用邏輯和HTTP 流式協(xié)議分開。模型調用邏輯是一個普通的異步生成器它只負責產出結構化事件字典HTTP 層只負責把事件字典按 SSE 協(xié)議編碼。這樣拆開你可以對模型邏輯單測不用每次起服務。from fastapi import FastAPI from fastapi.responses import StreamingResponse import json, asyncio from langchain_openai import ChatOpenAI from langchain_core.output_parsers import PydanticOutputParser from pydantic import BaseModel, Field from typing import List app FastAPI() llm ChatOpenAI(modelgpt-4o-mini, temperature0.3, streamingTrue) class FinalAnswer(BaseModel): reply: str Field(description面向用戶的最終回答) keywords: List[str] Field(description3-5個關鍵詞) parser PydanticOutputParser(pydantic_objectFinalAnswer) async def event_generator(prompt: str): # 1. 自動構建帶格式約束的提示詞 formatted_prompt ( 請回答用戶問題并輸出嚴格JSON。\n問題{q}\n{fmt}\n ).format(qprompt, fmtparser.get_format_instructions()) # 2. 第一個事件告知開始 yield { event: message, data: json.dumps({type: start, payload: {time: asyncio.time()}}, ensure_asciiFalse) } # 3. 流式輸出 token 事件 collected async for chunk in llm.astream(formatted_prompt): collected chunk.content yield { event: message, data: json.dumps({type: token, payload: {content: chunk.content}}, ensure_asciiFalse) } # 控制推送節(jié)奏避免瞬間把積壓的token全倒出去 await asyncio.sleep(0) # 4. 結構化解析并推給前端 try: parsed parser.parse(collected) yield { event: message, data: json.dumps({type: structured, payload: parsed.model_dump()}, ensure_asciiFalse) } except Exception as exc: yield { event: message, data: json.dumps({type: parse_error, payload: {error: str(exc)}}, ensure_asciiFalse) } # 5. 結束事件 yield { event: message, data: json.dumps({type: done, payload: {finish_reason: stop}}, ensure_asciiFalse) } 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, Connection: keep-alive, X-Accel-Buffering: no, # 重要禁止Nginx緩沖 }, )X-Accel-Buffering: no這個 header是很多人在 Nginx 反代下 SSE 不流式的元兇。Nginx 默認會緩沖響應攢滿 4KB 或者等連接結束才發(fā)給前端你明明在服務端yield了前端卻半天沒動靜。加了這個 header 就是明確告訴 Nginx 別緩沖。5.2 前端事件分發(fā)與渲染解耦前端用fetch配合ReadableStream解析 SSE 就夠了不一定非要引eventsource-parser這類庫但引了確實省事。核心邏輯是讀到一行data:解析 JSON按type走不同的渲染函數(shù)。async function streamChat(prompt) { const resp await fetch(/chat/stream, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ prompt }), }); const reader resp.body.getReader(); const decoder new TextDecoder(); let buffer ; while (true) { const { value, done } await reader.read(); if (done) break; buffer decoder.decode(value, { stream: true }); // SSE事件以空行分隔 let sepIndex; while ((sepIndex buffer.indexOf(\n\n)) ! -1) { const rawEvent buffer.slice(0, sepIndex); buffer buffer.slice(sepIndex 2); const dataLine rawEvent.split(\n).find(line line.startsWith(data:)); if (!dataLine) continue; const msg JSON.parse(dataLine.slice(5).trim()); handleEvent(msg); } } } function handleEvent(msg) { switch (msg.type) { case token: appendText(msg.payload.content); break; case structured: fillMetaPanel(msg.payload); break; case tool_call: showToolIndicator(msg.payload.name); break; case done: stopLoading(); break; } }注意千萬別直接用瀏覽器的原生EventSource—— 它只支持 GET 請求而我們往往需要 POST 傳遞 prompt。原生EventSource也沒有自定義 header 的能力鑒權都麻煩。用fetch流式讀取是最通用的方案。5.3 一邊流式一邊結構化增量解析的實踐前面提到的JsonOutputParser增量解析在實際項目中可以這樣用每收到一段新 token就把collected追加后嘗試parser.parse如果解析成功哪怕還不完整只要能出部分字段就把部分結果推給結構化預覽事件如果解析失敗因為 JSON 還沒閉合忽略即可不算錯誤。這個模式我用來解決一個具體的痛點用戶問幫我總結這份文檔并給出三個要點模型正文還沒寫完我希望前端右側欄已經先把要點標題渲染出來。雖然嚴格說最終結果要以最后完整解析為準但增量預覽的體驗提升非常明顯。代價是每次追加 token 都會觸發(fā)一次 JSON 解析token 非常長時會有少量 CPU 開銷。我實測下來對普通問答長度的文本這個開銷可以忽略。6. 常見問題與排查技巧實錄6.1 SSE 流中途斷開空閑超時是頭號殺手提示stream disconnected before completion: idle timeout waiting for sse這類報錯絕大多數(shù)不是代碼問題是鏈路中的代理/網關配置問題。我遇到過最典型的三層排查順序第一層本地測試。先用curl -N直接打服務接口觀察事件是不是正常持續(xù)輸出。curl -N能實時打印服務器推來的每個事件如果這一步正常問題就不在后端。第二層查反向代理。Nginx 的proxy_read_timeout默認 60 秒模型思考時間一旦超過代理直接斷連。調大或用proxy_read_timeout 300s可以緩解。另外確認proxy_buffering off;或X-Accel-Buffering: no已生效。第三層查云廠商網關。很多云負載均衡器對長連接也有空閑超時限制比如 60 秒、120 秒。盡量用 WebSocket 或 SSE 都能走的長連接配置同時后端在流式傳輸過程中即使沒有數(shù)據(jù)也要定期發(fā)一個: ping注釋行作為心跳。SSE 規(guī)范里以冒號開頭的行是注釋客戶端會忽略它但能刷新代理的空閑計時器。async def keepalive(): while True: yield : ping\n\n await asyncio.sleep(15)把這個生成器和主事件生成器用asyncio.gather合并就能在模型長時間思考時維持連接活性。6.2 OutputParser 拿到半截 JSON或 Markdown 代碼塊模型輸出里常見的臟格式有兩種一是把 JSON 藏在json代碼塊里二是前后夾帶解釋文字。PydanticOutputParser本身會嘗試從文本里提取 JSON 塊但并不可靠。我的兜底方案是寫一個基礎清洗函數(shù)在喂給解析器之前先做預處理。import re, json def extract_json_string(text: str) - str: text text.strip() # 去掉首尾的 markdown 代碼塊標記 code_block_pattern re.compile(r(?:json)?\s*(.*?)\s*, re.DOTALL) match code_block_pattern.search(text) if match: return match.group(1) # 嘗試從第一個 { 到最后一個 } 截取 start text.find({) end text.rfind(}) if start ! -1 and end ! -1 and end start: return text[start:end1] return text然后統(tǒng)一走parse。注意如果清洗后還是解析失敗再去觸發(fā) LLM 原文本修復。一定不要默認讓每個失敗都走修復不然成本和延遲都會失控。6.3 并發(fā)請求下事件錯亂上下文變量與隊列如果你的 FastAPI 服務同時處理多個 SSE 會話每個會話的生成器是獨立的理論上不會串。但我踩過一個實際的坑在生成器內部用了模塊級的全局變量緩存工具結果兩個用戶同時觸發(fā)同一個工具調用時A 用戶的結果可能被 B 用戶覆蓋。解決方案很簡單每個會話的事件生成器必須是自包含的所有狀態(tài)都放在生成器內部不要依賴模塊級可變對象。需要跨函數(shù)傳狀態(tài)就用contextvars或者干脆把 session_id 作為 key 放進一個字典管理隊列。我在項目里用的模式是每個會話一個asyncio.Queue生成器往隊列放事件SSE 層從隊列取事件編碼輸出。這個抽象能讓你在后續(xù)擴展多 Agent 編排時游刃有余。7. 一些我踩過坑之后的固定習慣先說工具聲明。LangChain 的tool裝飾器會讀取函數(shù)的 docstring 和類型注解來生成工具的 schema。docstring 里的描述、參數(shù)的類型提示、默認值都會成為傳給模型的 tool schema 的一部分。所以我在寫工具函數(shù)時會強制自己把每個參數(shù)的單位、取值范圍、邊界情況寫進 docstring這不是為了寫注釋好看而是直接決定模型能不能正確填參。一個只寫 keyword: 搜索關鍵詞 的工具和一個寫著 keyword: 搜索關鍵詞最長20字符不要帶引號 的工具在模型調參準確率上差很多。然后是解析器的temperature。做結構化輸出時模型溫度不建議設太高0 到 0.3 之間最穩(wěn)。溫度高了模型更容易發(fā)揮創(chuàng)造力去改格式、加注釋這對我們的 JSON 解析是災難。如果既要創(chuàng)意又要結構化我一般拆兩條鏈一條低溫度出結構化摘要一條高溫度潤色成自然語言回復最后再把兩段結果拼進 SSE 事件里。最后再分享一個小技巧給事件協(xié)議加一個trace_id字段。每個 SSE 會話生成一個trace_id在start事件里發(fā)給前端同時在服務端日志里打出來。前端報 bug 時直接甩這個 ID 給你你能在日志里把整條鏈路還原出來省掉大量你剛才問的什么來著的溝通成本。我在生產環(huán)境靠這個字段排查過很多偶發(fā)問題尤其是模型偶發(fā)出錯那種玄學問題有 trace_id 才能對上號。這套方案跑穩(wěn)之后你會發(fā)現(xiàn)流式體驗和結構化返回其實不是二選一關鍵是讓它們各走各的通道文本走 token 事件結構化走獨立事件中間用統(tǒng)一的 JSON 協(xié)議串聯(lián)。理解了這層設計后續(xù)換模型、加工具、上 LangGraph 編排都不會再被到底是文本還是 JSON這個問題卡住。