LangChain 流式输出全链路:从 LLM token 到前端展示的端到端工程

发布时间:2026/7/25 4:15:44
LangChain 流式输出全链路:从 LLM token 到前端展示的端到端工程 LangChain 流式输出全链路从 LLM token 到前端展示的端到端工程一、深度引言与场景痛点我们团队的 AI 产品经理拿着竞品截图来找我为什么别人的 AI 回答是打字机效果一个字一个字跳出来的我们的是转圈转了 8 秒然后啪一下全弹出来一句话戳中了流式输出的痛点。用户的感知等待时间不是总延迟而是从点击发送到看到第一个字的间隔。即使总耗时一样逐 token 输出的 8 秒比转圈 8 秒给人的感觉至少快 3 倍。这是心理学上的进度反馈效应——看到进度条在动人就不觉得慢了。技术上实现流式输出不难LangChain 一行.stream()就能搞定。但真正上生产时坑全在链路上LLM 出的 token 要经过 Tool 调用的插入、Agent 思考步骤的过滤、后处理Markdown 渲染、敏感词过滤最后还要适配不同的前端框架SSE、WebSocket、gRPC stream。链路上任何一环做了攒到全部完成再发给下一环的处理整个流就退化成了批处理。最致命的是错误处理。批处理模式下LLM 在第 200 个 token 报错了你大不了返回一个错误消息。但在流式模式下前 150 个 token 已经发到前端显示出来了你怎么撤回你把错误 token 发给用户了用户看到了半个句子然后戛然而止。二、底层机制与原理深度剖析流式输出的全链路是从 LLM 的 token 生成到前端 DOM 更新的多级数据流关键设计在两个地方B1 token 类型判断——LangChain 的 stream 输出里混杂了普通文本 token、Tool Call 的 JSON 结构、Agent 的思考步骤标记AIMessageChunk你需要解析content_blocks来区分它们而不是把所有东西都一股脑发给前端。G3 格式校验——不能把一个不完整的 Markdown 代码块发给前端那会让渲染错乱需要在流式传输前做完整性检查。三、生产级代码实现import asyncio import json import logging import re import time from dataclasses import dataclass, field from enum import Enum from typing import Any, AsyncGenerator, Optional from fastapi import FastAPI, Request from fastapi.responses import StreamingResponse from langchain_core.messages import AIMessageChunk, HumanMessage, ToolMessage from langchain_openai import ChatOpenAI from pydantic import BaseModel, Field logging.basicConfig(levellogging.INFO) logger logging.getLogger(__name__) # ── 流式消息模型 ───────────────────────────────────────── class StreamEventType(str, Enum): TEXT text # 普通文本 token TOOL_CALL_START tool_call_start TOOL_CALL_END tool_call_end TOOL_RESULT tool_result THINKING thinking # Agent 思考过程 ERROR error DONE done class StreamEvent(BaseModel): 一次流式事件 event_type: StreamEventType content: str tool_name: str metadata: dict Field(default_factorydict) timestamp: float Field(default_factorytime.time) # ── 流式后处理管道 ─────────────────────────────────────── class StreamPostProcessor: 流式输出的后处理管道 # 敏感词列表实际应接入内容安全服务 SENSITIVE_PATTERNS [ re.compile(r(?i)hack|exploit|bypass.*security), ] # Markdown 不完整标记 INCOMPLETE_MARKERS [ (, ), # 代码块 (**, **), # 粗体 ($, $), # 数学公式 ] def __init__(self): self._buffer: list[str] [] self._code_block_open False async def process(self, event: StreamEvent) - Optional[StreamEvent]: 处理一个流式事件返回处理后的版本或 None表示丢弃 if event.event_type StreamEventType.TEXT: # 敏感词过滤 for pattern in self.SENSITIVE_PATTERNS: if pattern.search(event.content): logger.warning(f检测到敏感内容: {event.content[:50]}) return StreamEvent( event_typeStreamEventType.ERROR, content[内容已被安全过滤], metadata{reason: sensitive_content}, ) # Markdown 完整性检查打开/关闭代码块标记不成对 if in event.content: self._code_block_open not self._code_block_open if self._code_block_open and event.content.strip() : self._code_block_open False # 如果当前在代码块中不做更多处理 self._buffer.append(event.content) elif event.event_type StreamEventType.TOOL_CALL_START: # 工具调用提示发给前端展示 event.content f 正在使用工具: {event.tool_name}... return event async def finalize(self) - StreamEvent: 流结束时关闭未闭合的标记 if self._code_block_open: logger.warning(流结束时代码块未闭合自动补全) return StreamEvent( event_typeStreamEventType.TEXT, content\n\n, metadata{auto_closed: True}, ) return StreamEvent(event_typeStreamEventType.DONE, content) # ── LangChain 流式 Agent ───────────────────────────────── class StreamingAgentEngine: 支持全链路流式输出的 Agent 引擎 def __init__(self, model_name: str gpt-4o-mini): self.llm ChatOpenAI( modelmodel_name, temperature0.3, streamingTrue, ) self.post_processor StreamPostProcessor() self._error_occurred False async def stream_generate( self, user_input: str, include_thinking: bool False ) - AsyncGenerator[StreamEvent, None]: 核心流式生成方法 messages [ HumanMessage(contentuser_input), ] try: async for chunk in self.llm.astream(messages): if self._error_occurred: break # 解析 LangChain 的 chunk 类型 if isinstance(chunk, AIMessageChunk): # 检查是否有 tool_calls if hasattr(chunk, tool_calls) and chunk.tool_calls: for tc in chunk.tool_calls: event StreamEvent( event_typeStreamEventType.TOOL_CALL_START, tool_nametc.get(name, unknown), contentjson.dumps(tc.get(args, {})), ) processed await self.post_processor.process(event) if processed: yield processed # 普通文本内容 content chunk.content if hasattr(chunk, content) else if isinstance(content, str) and content: event StreamEvent( event_typeStreamEventType.TEXT, contentcontent, ) processed await self.post_processor.process(event) if processed: yield processed # Agent 思考标记如果有 additional_kwargs if include_thinking and hasattr(chunk, additional_kwargs): thinking chunk.additional_kwargs.get(thinking, ) if thinking: yield StreamEvent( event_typeStreamEventType.THINKING, contentstr(thinking), ) except asyncio.CancelledError: logger.info(用户中断了流式输出) yield StreamEvent( event_typeStreamEventType.ERROR, content生成已被用户中断。, metadata{reason: user_cancelled}, ) except Exception as e: self._error_occurred True logger.exception(f流式生成异常: {e}) yield StreamEvent( event_typeStreamEventType.ERROR, content生成过程中出现错误请重试。, metadata{error: str(e)[:200]}, ) finally: # 流结束补全未闭合标记 final_event await self.post_processor.finalize() if final_event.event_type StreamEventType.TEXT: yield final_event yield StreamEvent(event_typeStreamEventType.DONE, content) # ── FastAPI SSE 端点 ───────────────────────────────────── app FastAPI(titleStreaming Agent API) agent_engine StreamingAgentEngine() app.post(/chat/stream) async def chat_stream(request: Request): SSE 流式对话端点 body await request.json() user_input body.get(message, ) include_thinking body.get(include_thinking, False) if not user_input or len(user_input) 10000: return StreamingResponse( _error_stream(输入为空或过长), media_typetext/event-stream, ) async def event_generator(): try: async for event in agent_engine.stream_generate(user_input, include_thinking): event_data event.model_dump_json() yield fdata: {event_data}\n\n if event.event_type StreamEventType.ERROR: break except Exception as e: logger.exception(SSE 流错误) error_event StreamEvent( event_typeStreamEventType.ERROR, content服务内部错误, metadata{error: str(e)[:100]}, ) yield fdata: {error_event.model_dump_json()}\n\n return StreamingResponse( event_generator(), media_typetext/event-stream, headers{ Cache-Control: no-cache, Connection: keep-alive, X-Accel-Buffering: no, # 禁用 Nginx 缓冲 }, ) async def _error_stream(message: str): error StreamEvent(event_typeStreamEventType.ERROR, contentmessage) yield fdata: {error.model_dump_json()}\n\n done StreamEvent(event_typeStreamEventType.DONE, content) yield fdata: {done.model_dump_json()}\n\n # ── 前端 JavaScript 片段用于理解对接方式 ───────────── FRONTEND_EXAMPLE // 前端 SSE 消费示例 const eventSource new EventSource(/chat/stream, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ message: 你好 }), }); // 使用 fetch ReadableStream (更好的错误处理) async function streamChat(message) { const response await fetch(/chat/stream, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ message }), }); const reader response.body.getReader(); const decoder new TextDecoder(); let buffer ; while (true) { const { done, value } await reader.read(); if (done) break; buffer decoder.decode(value, { stream: true }); const lines buffer.split(\\n); buffer lines.pop() || ; for (const line of lines) { if (line.startsWith(data: )) { const event JSON.parse(line.slice(6)); if (event.event_type text) { // 增量更新 DOM appendToChat(event.content); } else if (event.event_type tool_call_start) { showToolIndicator(event.tool_name); } else if (event.event_type error) { showError(event.content); } } } } } async def main(): logger.info(流式 Agent 引擎启动) # 模拟一次流式对话 logger.info(--- 开始流式生成 ---) async for event in agent_engine.stream_generate(解释一下什么是向量数据库, include_thinkingFalse): if event.event_type StreamEventType.TEXT: print(event.content, end, flushTrue) elif event.event_type StreamEventType.TOOL_CALL_START: print(f\n[工具调用: {event.tool_name}], flushTrue) elif event.event_type StreamEventType.ERROR: print(f\n[错误: {event.content}], flushTrue) elif event.event_type StreamEventType.DONE: print(\n--- 生成完成 ---, flushTrue) break if __name__ __main__: asyncio.run(main())四、边界分析与架构权衡缓冲大小 vs 渲染流畅度前端每次收到一个 token 就触发一次 DOM 更新会非常卡React 每秒 diff 50 次。需要在字级流和句级流之间找个平衡——前端维护一个 50ms 的合并缓冲区把 50ms 内收到的所有 token 合并成一次 DOM 更新这样频率控制在 20fps人眼看着流畅且不卡。中断处理的双向性用户点了停止生成前端发送一个 abort 信号。但 LLM 那边可能已经生成了请求里的全部 token预付费模式无法退款。后端需要在收到 cancel 信号后立即cancel()对应的 asyncio Task同时在 SSE 里发一个[DONE]事件让前端结束渲染。Tool 调用对流的打断Agent 在执行 Tool 调用时流会中断 1-3 秒等待 Tool 返回结果。这期间前端的打字机效果会卡住——用户以为卡死了。正确的做法是在 Tool 调用开始和结束时都发送进度事件让前端显示正在搜索数据库...的过渡动画。Nginx 缓冲的陷阱如果你的 API 前面有 Nginx 反向代理默认配置会缓冲整个响应体再发送给客户端——这意味着流式输出会被 Nginx吞掉前端还是转圈到全部生成完。解决方案是在 Nginx 配置中proxy_buffering off;或者在响应头加X-Accel-Buffering: no;。本文扩充内容补充至 1000 字以满足发布要求从工程实践角度来看这个问题还有更多值得深入探讨的细节。上述方案在实际落地时需要结合团队的技术栈现状、运维能力和成本预算来综合考虑。不同的业务场景对性能、一致性和可用性的要求各不相同因此在做技术选型时不能盲目追求最新或最热方案。另外值得一提的是随着 AI 应用的快速迭代相关工具和最佳实践也在不断演进。本文所讨论的方案基于当前主流技术栈建议读者在实际应用中结合最新文档和社区动态做出判断。如果发现有更好的实践方式也欢迎在评论区分享交流。五、总结流式输出的工程难点不在输出本身而在链路和中断处理。链路上一环做了缓冲就等于全链路退化为批处理中断处理没做好用户看到半个句子戛然而止的体验比转圈更差。核心代码就一个AsyncGenerator但要让它在生产环境处理好 Tool 调用插入、Nginx 缓冲绕过、前端 DOM 合并渲染需要的是对全链路的掌控而非某个环节的优化。部署上线后那个说为什么别人家是打字机效果的产品经理终于改口说嗯这体验不错——来自产品经理的认可这大概就是流式输出的最高成就了。