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_KEY=sk-your-actual-api-key-here OPENAI_BASE_URL=https://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 lang="zh-CN"> <head> <meta charset="UTF-8"> <meta name="viewport" content="width=device-width, initial-scale=1.0"> <title>ChatGPT 流式输出演示</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 id="chatBox"></div> <div id="inputArea"> <input type="text" id="userInput" placeholder="输入你的问题..." /> <button onclick="sendMessage()">发送</button> <button onclick="clearChat()">清空</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 class="cursor">,在流式输出时显示,在流结束时移除,这能极大地提升用户体验,明确指示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?token=xxx&message=yyy。服务器端首先验证Token的有效性。注意:Token出现在URL中可能被日志记录,存在泄露风险,需确保日志不记录完整URL,并使用HTTPS。 - 方案B:放弃
EventSource,使用更灵活的fetchAPI来读取流,它可以设置任意请求头。但你需要自己解析SSE格式的数据流,实现会稍复杂一些。
- 方案A(推荐):将Token放在查询参数中,如
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确保了安全和逻辑聚合,而妥善的跨域处理则是让这一切在浏览器中顺利运行的桥梁。当你下次再看到“流式输出”这个词时,希望你的脑海里能清晰地浮现出这条数据流的完整旅行路径。