ARTICLE DETAIL

资讯详情

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

基于SSE与BFF架构实现大模型流式输出:从原理到实践

基于SSE与BFF架构实现大模型流式输出:从原理到实践 1. 项目缘起为什么我们需要“三件套”来实现流式输出最近在做一个内部知识库问答的Demo核心需求就是模仿ChatGPT那种“一个字一个字往外蹦”的流式回答体验。一开始想得很简单不就是个HTTP请求吗前端发个问题后端调一下大模型API等结果全返回了再一次性渲染到页面上。但真做起来才发现这种“等全部完成再展示”的方式在等待几秒甚至十几秒的过程中用户面对一个空白的页面体验非常糟糕会反复怀疑“是不是卡住了”“我网断了”。这恰恰是流式输出要解决的核心痛点即时反馈。那么如何实现流式输出一个常见的答案是Server-Sent Events。但如果你直接让前端页面去请求大模型厂商的API比如OpenAI的接口立刻就会撞上两堵墙跨域和安全性。大模型的API密钥是最高机密绝不能暴露在前端代码里。所以我们需要一个中间层来代理这个请求这就是BFF的用武之地。而BFF层在向前端推送流式数据时又需要处理好SSE连接和跨域问题。于是“SSE BFF 跨域”这个组合拳就自然而然地浮出水面了。这个项目就是带你从零开始打通这条链路。我会假设你有一个Node.js环境并且对Express框架有基本了解。我们将构建一个极简但功能完整的系统让你彻底理解每一环是怎么工作的以及为什么必须这么工作。文末会提供完整的、可运行的代码仓库。2. 核心原理拆解SSE、BFF与跨域各自扮演什么角色在动手写代码之前我们必须先搞清楚这三个技术点分别解决了什么问题以及它们是如何协同工作的。这能让你在后续遇到任何怪异现象时都能快速定位到问题层。2.1 Server-Sent Events单向的、长时的数据流SSE是一种允许服务器主动向客户端推送数据的技术。它与WebSocket不同WebSocket是双向通信而SSE是服务器到客户端的单向通道。对于“接收模型流式输出”这种典型的“一发一收”场景SSE既简单又合适。它的工作原理基于普通的HTTP协议。客户端通过EventSourceAPI发起一个GET请求并在请求头中带上Accept: text/event-stream。服务器收到后会保持这个连接不关闭并通过返回一个Content-Type: text/event-stream的响应开始源源不断地发送数据。数据的格式有严格规范data: 这是一条消息\n\n每条消息以data:开头以两个换行符\n\n结束。服务器可以分多次发送data:行客户端会将其拼接直到遇到两个换行符才视为一条完整消息并触发事件。为什么选SSE而不是WebSocket对于只是接收文本流的场景SSE的实现成本更低。它基于HTTP无需额外的协议升级握手天然支持断线重连EventSource有内置机制并且浏览器兼容性良好。如果你的场景不需要客户端频繁向服务器发送数据例如聊天室那么SSE是更轻量、更直接的选择。2.2 BFF关键的安全代理与逻辑聚合层BFF即面向前端的后端。在这个架构里它的核心职责有两个隐藏敏感信息前端只知道BFF的地址而不知道大模型API的地址和密钥。所有对模型API的请求都由BFF发起密钥安全地存储在服务器环境变量中。协议转换与流代理大模型API例如OpenAI的Chat Completion返回的通常也是一个流stream: true。BFF需要接收这个流并将其“翻译”成符合SSE格式的流再推送给前端。同时BFF还可以在这里做一些额外工作比如请求参数格式化、错误处理统一、日志记录等。没有BFF行不行理论上如果你能解决跨域且不介意暴露API密钥可以让前端直接连模型API。但现实中这两点都是不可接受的。因此BFF是生产环境中的必选项。2.3 跨域处理为SSE铺平道路跨域问题是浏览器出于安全考虑施加的限制。当你的前端页面假设运行在http://localhost:3000试图直接请求BFF服务器假设运行在http://localhost:3001的SSE接口时浏览器会阻止这个请求。解决跨域主要是在BFF服务器上设置CORS响应头。对于SSE有两点需要特别注意Access-Control-Allow-Origin: 必须明确设置为前端的源地址如http://localhost:3000或使用*不推荐在生产环境使用且某些浏览器在携带凭证时禁用*。Access-Control-Allow-Headers: 如果需要前端传递自定义头比如认证Token需要在这里声明。此外由于EventSource请求默认不携带Cookie等凭证如果需要还须设置Access-Control-Allow-Credentials: true并且前端在创建EventSource时也要设置withCredentials。不过在我们的简单Demo里暂不涉及凭证。3. 环境准备与项目初始化接下来我们开始动手搭建。你需要确保电脑上安装了Node.js建议版本16和npm。首先创建一个项目目录并初始化mkdir chatgpt-stream-demo cd chatgpt-stream-demo npm init -y然后安装我们所需的依赖。核心依赖是express用于创建BFF服务器axios用于向大模型API发起流式请求cors用于方便地处理跨域。另外我们安装dotenv来管理环境变量。npm install express axios cors dotenv为了开发方便我们还可以安装nodemon作为开发依赖实现代码热更新。npm install --save-dev nodemon接着创建项目文件结构chatgpt-stream-demo/ ├── server/ │ ├── index.js # BFF服务器主入口 │ └── .env # 环境变量文件需自行创建不要提交到git ├── client/ │ └── index.html # 前端页面 ├── package.json └── README.md在server/.env文件中放入你的大模型API密钥和基地址。这里以OpenAI格式为例OPENAI_API_KEYsk-your-actual-api-key-here OPENAI_BASE_URLhttps://api.openai.com/v1重要提示.env文件务必添加到.gitignore中避免密钥泄露。最后修改package.json添加一个启动脚本{ scripts: { start: node server/index.js, dev: nodemon server/index.js } }4. BFF服务器实现构建流式代理中间件BFF服务器是整个系统的中枢。它的任务很明确提供一个SSE端点供前端连接当收到前端请求后去调用真正的大模型API将API返回的流实时转换为SSE格式推送给前端。4.1 基础服务器与SSE端点搭建首先在server/index.js中搭建一个基本的Express服务器并设置CORS。const express require(express); const axios require(axios); const cors require(cors); require(dotenv).config({ path: .env }); const app express(); const PORT 3001; // 配置CORS允许来自前端开发服务器的请求 app.use(cors({ origin: http://localhost:3000, // 你的前端地址 credentials: false // 本例不需要凭证 })); // 关键SSE端点 app.get(/api/chat/stream, async (req, res) { const { message } req.query; // 从前端获取用户消息 if (!message) { return res.status(400).json({ error: Message is required }); } // 1. 设置SSE相关的响应头 res.writeHead(200, { Content-Type: text/event-stream, Cache-Control: no-cache, Connection: keep-alive, // CORS头对于SSE同样重要 Access-Control-Allow-Origin: http://localhost:3000, }); // 2. 模拟或真实调用大模型API // 我们先写一个模拟函数确保SSE链路畅通 simulateLLMStream(res, message); // 注意我们不需要 res.end()连接需要保持打开以持续发送数据。 }); // 模拟流式生成函数 function simulateLLMStream(res, userMessage) { const responses [ 你好, 你问的问题是“${userMessage}”。, 这是一个模拟的流式响应。, 我正在一个字、一个字地返回给你。, 这样你就能看到类似ChatGPT的效果了 ]; let index 0; const intervalId setInterval(() { if (index responses.length) { // SSE格式 data: 内容\n\n res.write(data: ${JSON.stringify({ content: responses[index] })}\n\n); index; } else { // 发送结束信号 res.write(data: [DONE]\n\n); clearInterval(intervalId); // 在实际调用中不要主动结束连接由模型流结束或客户端关闭 // res.end(); } }, 300); // 每300毫秒发送一段 } app.listen(PORT, () { console.log(BFF server listening on http://localhost:${PORT}); });这段代码创建了一个/api/chat/stream的GET端点。它设置了正确的SSE响应头并使用一个simulateLLMStream函数来模拟大模型逐句返回数据的过程。每句数据都被包装成data: { content: 文本 }\n\n的格式发送。最后发送一个特殊的[DONE]事件通知前端流已结束。为什么用res.write而不是res.send或res.json因为SSE连接是一个持久化的连接我们需要多次向同一个响应对象写入数据。res.send或res.json会在调用后自动结束响应无法再写入。4.2 集成真实大模型API以OpenAI为例模拟成功之后我们来替换成真实的OpenAI API调用。我们需要使用axios并配置其responseType为stream来接收流式响应。首先在文件顶部配置axios实例const openai axios.create({ baseURL: process.env.OPENAI_BASE_URL, headers: { Authorization: Bearer ${process.env.OPENAI_API_KEY}, Content-Type: application/json, }, });然后将/api/chat/stream端点中的simulateLLMStream调用替换为真实的API调用逻辑app.get(/api/chat/stream, async (req, res) { const { message } req.query; if (!message) { return res.status(400).json({ error: Message is required }); } res.writeHead(200, { Content-Type: text/event-stream, Cache-Control: no-cache, Connection: keep-alive, Access-Control-Allow-Origin: http://localhost:3000, }); // 真实调用OpenAI API try { const response await openai.post(/chat/completions, { model: gpt-3.5-turbo, // 或你选择的模型 messages: [{ role: user, content: message }], stream: true, // 关键开启流式输出 temperature: 0.7, }, { responseType: stream, // 关键让axios返回一个流 }); // OpenAI的流式响应是多个SSE格式的数据块 const stream response.data; stream.on(data, (chunk) { // 每个chunk是一个Buffer需要转成字符串并处理 const lines chunk.toString().split(\n).filter(line line.trim() ! ); for (const line of lines) { if (line.startsWith(data: )) { const message line.replace(/^data: /, ); if (message [DONE]) { // 流结束 res.write(data: [DONE]\n\n); return; } try { const parsed JSON.parse(message); // OpenAI返回的choices[0].delta.content是增量内容 const content parsed.choices[0]?.delta?.content; if (content) { // 将内容包装成我们自己的SSE格式发送给前端 res.write(data: ${JSON.stringify({ content })}\n\n); } } catch (err) { // 忽略非JSON或解析错误的数据行 console.error(Error parsing SSE message:, err.message); } } } }); stream.on(end, () { console.log(Stream from OpenAI ended.); // 确保发送结束标记 res.write(data: [DONE]\n\n); }); stream.on(error, (err) { console.error(Stream error:, err); res.write(data: ${JSON.stringify({ error: Stream interrupted })}\n\n); res.write(data: [DONE]\n\n); }); // 当客户端断开连接时清理OpenAI的流 req.on(close, () { stream.destroy(); console.log(Client disconnected, stream destroyed.); }); } catch (error) { console.error(Error calling OpenAI API:, error.response?.data || error.message); res.writeHead(500, { Content-Type: application/json }); res.end(JSON.stringify({ error: Failed to call AI service })); } });这段代码有几个关键点stream: true和responseType: stream是让整个流程“流起来”的核心配置。OpenAI返回的流本身就是一种SSE格式每行以data:开头。我们需要解析这些行提取出delta.content再重新包装成我们自己的SSE事件发送给前端。这样做的好处是BFF层可以对数据做统一的格式化或过滤。错误处理至关重要。包括网络错误、API错误、解析错误以及客户端提前断开连接的情况req.on(close)。必须妥善处理这些情况避免资源泄漏如未关闭的流和服务器崩溃。我们最终发送给前端的格式是data: {content:xxx}\n\n这是一个JSON字符串。前端需要JSON.parse来获取content字段。你也可以直接发送纯文本data: xxx\n\n但JSON格式更易于扩展未来可以传递更多元数据如消息ID、是否结束等。5. 前端实现用EventSource接收并渲染流BFF准备好了现在我们来创建一个简单的前端页面来连接它。在client/index.html中编写如下代码!DOCTYPE html html langzh-CN head meta charsetUTF-8 meta nameviewport contentwidthdevice-width, initial-scale1.0 titleChatGPT 流式输出演示/title style body { font-family: sans-serif; max-width: 800px; margin: 40px auto; padding: 20px; } #chatBox { border: 1px solid #ccc; height: 400px; overflow-y: auto; padding: 10px; margin-bottom: 20px; } .message { margin-bottom: 10px; } .user { text-align: right; color: blue; } .assistant { text-align: left; color: green; } #inputArea { display: flex; } #userInput { flex-grow: 1; padding: 10px; font-size: 16px; } button { padding: 10px 20px; font-size: 16px; margin-left: 10px; } .cursor { display: inline-block; width: 8px; height: 1em; background-color: #333; animation: blink 1s infinite; margin-left: 2px; } keyframes blink { 50% { opacity: 0; } } /style /head body h1 流式对话演示/h1 div idchatBox/div div idinputArea input typetext iduserInput placeholder输入你的问题... / button onclicksendMessage()发送/button button onclickclearChat()清空/button /div script const chatBox document.getElementById(chatBox); const userInput document.getElementById(userInput); let currentEventSource null; let assistantMessageDiv null; let accumulatedText ; function appendMessage(role, text) { const div document.createElement(div); div.className message ${role}; div.textContent ${role user ? 你 : AI}: ${text}; chatBox.appendChild(div); chatBox.scrollTop chatBox.scrollHeight; // 自动滚动到底部 } function createStreamingMessage() { // 创建一个新的助手消息容器并添加闪烁光标 assistantMessageDiv document.createElement(div); assistantMessageDiv.className message assistant; const label document.createElement(span); label.textContent AI: ; const contentSpan document.createElement(span); contentSpan.id streamingContent; const cursorSpan document.createElement(span); cursorSpan.className cursor; assistantMessageDiv.appendChild(label); assistantMessageDiv.appendChild(contentSpan); assistantMessageDiv.appendChild(cursorSpan); chatBox.appendChild(assistantMessageDiv); chatBox.scrollTop chatBox.scrollHeight; accumulatedText ; // 重置累积文本 return contentSpan; } function sendMessage() { const message userInput.value.trim(); if (!message) return; // 显示用户消息 appendMessage(user, message); userInput.value ; // 创建流式消息容器 const contentSpan createStreamingMessage(); // 如果存在旧的连接先关闭 if (currentEventSource) { currentEventSource.close(); } // 创建新的EventSource连接 // 注意EventSource只支持GET请求所以我们将消息放在查询参数中 const url http://localhost:3001/api/chat/stream?message${encodeURIComponent(message)}; currentEventSource new EventSource(url); currentEventSource.onmessage function(event) { const data event.data; if (data [DONE]) { // 流结束移除光标 const cursor document.querySelector(#streamingContent .cursor); if (cursor) cursor.remove(); currentEventSource.close(); currentEventSource null; return; } try { const parsed JSON.parse(data); if (parsed.error) { contentSpan.textContent \n[错误: ${parsed.error}]; } else if (parsed.content) { // 累积并更新内容 accumulatedText parsed.content; contentSpan.textContent accumulatedText; // 保持滚动 chatBox.scrollTop chatBox.scrollHeight; } } catch (e) { console.error(解析SSE数据失败:, e); } }; currentEventSource.onerror function(err) { console.error(EventSource failed:, err); contentSpan.textContent \n[连接出错或已关闭]; const cursor document.querySelector(#streamingContent .cursor); if (cursor) cursor.remove(); currentEventSource.close(); currentEventSource null; }; } function clearChat() { chatBox.innerHTML ; if (currentEventSource) { currentEventSource.close(); currentEventSource null; } } // 支持按回车发送 userInput.addEventListener(keypress, function(e) { if (e.key Enter) { sendMessage(); } }); /script /body /html前端实现的核心逻辑建立连接使用new EventSource(url)连接到BFF的SSE端点。由于EventSource只支持GET请求我们将用户消息放在了URL的查询参数?message中。对于长消息这可能有问题更复杂的场景可以考虑先用POST发送消息再由服务器返回一个唯一的SSE连接URL。接收数据监听onmessage事件。每次服务器发送一个data: ...\n\n这个事件就会被触发。我们解析数据如果是[DONE]就结束流否则将content字段的内容累加到前一个DOM元素中实现逐字打印的效果。视觉反馈我们创建了一个闪烁的光标span classcursor在流式输出时显示在流结束时移除这能极大地提升用户体验明确指示AI正在“思考”或“打字”。错误处理与连接管理监听onerror事件处理连接错误。同时在发送新消息或清空聊天时记得用currentEventSource.close()关闭旧的连接防止连接数累积。6. 运行、测试与核心问题排查现在让我们把整个系统跑起来。启动BFF服务器在项目根目录下运行npm run dev。你应该看到BFF server listening on http://localhost:3001。启动前端服务器由于前端是纯HTML文件我们需要一个HTTP服务器来提供它。你可以使用任何静态服务器。一个快速的方法是使用Python在client目录下运行python3 -m http.server 3000。或者使用Node.js的serve工具npx serve client -p 3000。打开浏览器访问http://localhost:3000。测试在输入框中提问比如“介绍一下你自己”点击发送。你应该能看到你的问题先出现然后下方AI的回答开始逐字逐句地显示出来伴随着闪烁的光标。在这个过程中你可能会遇到一些典型问题下面是我的排查经验问题一前端控制台报错“Failed to load resource: net::ERR_FAILED”或跨域错误。检查确保BFF服务器的CORS配置中origin字段与前端页面的实际访问地址包括端口完全一致。浏览器控制台的Network标签页里查看SSE请求的Response Headers中是否包含Access-Control-Allow-Origin: http://localhost:3000。解决在BFF代码中确保SSE的响应头也设置了CORS头如我们代码中在res.writeHead里做的那样。有时候普通的中间件app.use(cors(...))对SSE连接可能不生效需要显式设置。问题二连接建立成功但收不到任何数据或者很快断开。检查首先在BFF服务器的控制台查看是否有请求进来是否有错误日志。然后在浏览器的Network标签页中找到那个SSE请求查看其“EventStream”标签页Chrome有此功能看是否能直接看到服务器推送的原始数据。解决数据格式确保服务器发送的数据严格遵守data: payload\n\n格式。多一个空格、少一个换行都可能导致EventSource无法正确解析。我常用res.write(data: ${JSON.stringify(payload)}\n\n)来确保格式正确。连接保持确保服务器没有在流结束前调用res.end()。SSE连接需要一直保持。防火墙/代理本地开发一般没问题但在服务器部署时确保中间没有代理或防火墙中断长连接。问题三流式输出不“流”而是一次性显示完整句子。检查这通常不是SSE或前端的问题而是大模型API返回的数据块本身就比较大。例如某些API可能一句话就是一个完整的delta。你可以尝试在BFF服务器中打印出每次收到的content看看它是不是已经是完整的句子。解决为了获得更细腻的“逐字”效果你可以在BFF层对content字符串进行进一步拆分比如按字符或按词对于英文循环发送。但这会增加复杂性和延迟需要权衡。大多数情况下以句子或短语为单位的流式体验已经足够好。问题四内存泄漏或连接数过多。检查如果用户频繁快速发送消息而旧连接没有正确关闭会导致服务器端积累大量僵尸连接。解决正如我们在前端代码中所做在创建新EventSource前先close()旧的。在服务器端监听req.on(close)事件并在其中销毁对应的模型API请求流stream.destroy()及时释放资源。7. 生产环境进阶考量与优化Demo跑通了但要用于实际项目还有几个关键点需要加固1. 认证与鉴权目前的端点对所有人开放。在生产环境中必须添加认证。一个常见的模式是用户登录后前端获取一个JWT Token。前端连接SSE时无法通过标准的EventSource设置Authorization头这是EventSource的一个限制。变通方案有两种方案A推荐将Token放在查询参数中如/api/chat/stream?tokenxxxmessageyyy。服务器端首先验证Token的有效性。注意Token出现在URL中可能被日志记录存在泄露风险需确保日志不记录完整URL并使用HTTPS。方案B放弃EventSource使用更灵活的fetchAPI来读取流它可以设置任意请求头。但你需要自己解析SSE格式的数据流实现会稍复杂一些。2. 错误处理与重试网络抖动EventSource有内置的重连机制但在断线时用户可能丢失正在接收的消息。对于关键应用可以考虑在BFF层实现一个简单的消息缓存或确认机制。模型API错误OpenAI API可能返回速率限制、模型过载等错误。BFF层需要捕获这些错误并将其转换为前端能理解的SSE错误事件例如data: {error: rate_limit}\n\n并优雅地结束流而不是直接抛出异常导致服务器500错误。3. 性能与可扩展性连接数Node.js的单个进程能保持的并发连接数有限。当用户量很大时需要考虑使用集群Cluster模式或将其部署为无状态服务通过负载均衡器分散连接。超时设置给SSE连接和上游模型API调用设置合理的超时时间。避免因为一个慢请求长期占用连接资源。4. 前端体验优化中止请求提供用户一个“停止生成”按钮。这需要前端能主动中止EventSource连接.close()并且BFF能相应地中止对模型API的请求。历史记录将对话历史存储在BFF的会话中或发送给模型API以实现多轮对话的上下文连贯。Markdown渲染如果模型返回Markdown格式的内容前端可以使用诸如marked的库进行实时渲染提升阅读体验。通过这个从零开始的“三件套”实践我们不仅实现了一个功能更打通了对流式传输、前后端分离架构中安全与通信的理解。每一层都有其不可替代的作用SSE提供了高效的推送通道BFF确保了安全和逻辑聚合而妥善的跨域处理则是让这一切在浏览器中顺利运行的桥梁。当你下次再看到“流式输出”这个词时希望你的脑海里能清晰地浮现出这条数据流的完整旅行路径。
返回列表