流式输出(Streaming)实现与前端对接实战指南
流式输出(Streaming)实现与前端对接实战指南
① 流式传输核心概念与生活化类比解析
先从一个日常场景说起。
你点了一份外卖,有两种等待方式:一种是商家把所有菜做好、打包完毕,再一次性送到你手上——这期间你只能干等,啥也看不到;另一种是厨师边做边让骑手先送过来,做好一道送一道,你就能看着菜一盘盘上桌,心里有底,也不着急。
流式输出就是后者。
传统的HTTP请求就像第一种方式——客户端发一个请求,服务器闷头把活干完,把完整结果一次性甩回来。而流式输出允许服务器一边生成数据一边往客户端推送,客户端收到一点就渲染一点。
具体到技术层面,最常用的实现方案叫 SSE(Server-Sent Events) ,全称“服务器发送事件”。你可以把它想象成你关注了一个新闻App的“突发新闻”推送——你只需要在App里点一次“允许通知”(这就是建立连接),之后只要有大新闻,服务器就会主动把消息推送到你手机上,你不用一遍遍去刷新。
SSE和WebSocket的区别在哪?打个比方:WebSocket像微信电话——你和服务器都能随时说话,是双向的;SSE像新闻推送——只有服务器能“说话”,你只管听,是单向的。在AI对话场景里,用户问完问题后只需要静静看AI把答案一个字一个字“说”出来就够了,不需要中途再插话。所以更轻量的SSE是更合适的选择。
SSE基于HTTP协议,使用text/event-stream作为MIME类型,有三个核心优势:单向实时通信、基于HTTP无需复杂握手、自带自动重连机制。
② 后端流式接口搭建与环境快速部署
理解了概念,咱们直接上手写代码。
Node.js(Express)版本
const express = require('express');
const app = express();
app.get('/stream', (req, res) => {
// 三个响应头是SSE正常工作的关键
res.writeHead(200, {
'Content-Type': 'text/event-stream; charset=utf-8',
'Cache-Control': 'no-cache',
'Connection': 'keep-alive'
});
// 这句很重要:告诉Nginx等代理不要缓冲
res.flushHeaders();
let count = 0;
const interval = setInterval(() => {
count++;
// SSE数据格式:data: 内容\n\n
res.write(`data: 第${count}条消息,时间:${new Date().toLocaleTimeString()}\n\n`);
}, 1000);
// 客户端断开连接时清理定时器
req.on('close', () => {
clearInterval(interval);
res.end();
});
});
app.listen(3000, () => console.log('服务已启动: http://localhost:3000'));
关键点在于响应头必须在写响应体之前设置,否则会报错。text/event-stream声明了这是一个事件流,no-cache禁用缓存保证实时性,keep-alive维持长连接。
Go + Gin 版本
package main
import (
"fmt"
"time"
"github.com/gin-gonic/gin"
)
func main() {
r := gin.Default()
r.GET("/events", func(c *gin.Context) {
c.Header("Content-Type", "text/event-stream")
c.Header("Cache-Control", "no-cache")
c.Header("Connection", "keep-alive")
clientClosed := c.Writer.CloseNotify()
ticker := time.NewTicker(1 * time.Second)
defer ticker.Stop()
for {
select {
case <-clientClosed:
return
case t := <-ticker.C:
event := fmt.Sprintf("data: %s\n\n", t.Format("2006-01-02 15:04:05"))
c.SSEvent("message", event)
c.Writer.Flush() // 立即刷新缓冲区
}
}
})
r.Run(":8080")
}
这里用了Flush()强制刷新缓冲区——如果不调用,数据可能被积压在缓冲区里,客户端就收不到实时效果。
③ 服务端数据分块发送代码实现详解
在实际业务中,数据通常不是定时器生成的,而是来自大模型、数据库查询或文件读取。下面展示如何把大模型的流式响应转发给前端。
// 以调用大模型API为例
app.post('/chat', async (req, res) => {
res.writeHead(200, {
'Content-Type': 'text/event-stream; charset=utf-8',
'Cache-Control': 'no-cache',
'Connection': 'keep-alive'
});
res.flushHeaders();
const response = await fetch('https://api.example.com/llm/chat', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
messages: req.body.messages,
stream: true // 开启流式
})
});
const reader = response.body.getReader();
const decoder = new TextDecoder();
while (true) {
const { done, value } = await reader.read();
if (done) break;
const chunk = decoder.decode(value, { stream: true });
// 直接把数据块转发给前端
res.write(`data: ${chunk}\n\n`);
}
res.write('data: [DONE]\n\n');
res.end();
});
SSE数据格式有固定要求:每组数据以\n\n结束,组内不同字段用\n分隔。比如要同时传递id和data:
id: 1
event: message
data: hello world
id: 2
event: custom
data: hello
data: world
④ 前端 Fetch API 接收流式响应基础写法
前端接收流式数据,最直接的方式就是用Fetch API配合ReadableStream。EventSource虽然用起来简单,但它只能发GET请求,不能自定义请求头,在需要鉴权(比如带Token)的场景下就捉襟见肘了。而Fetch可以自由控制请求方法、头信息和请求体。
基础写法:
async function fetchStream(url) {
const response = await fetch(url);
const reader = response.body.getReader();
const decoder = new TextDecoder();
while (true) {
const { done, value } = await reader.read();
if (done) break;
// value是Uint8Array,需要解码成字符串
console.log(decoder.decode(value));
}
}
response.body返回的就是一个ReadableStream,通过getReader()拿到读取器,然后循环调用read()方法逐块读取数据。
⑤ 使用 ReadableStream 逐块解析数据流
上面的基础写法有个问题——如果直接用decoder.decode(value),当一个多字节字符(比如中文的“你”,UTF-8编码占3个字节)被切分到两个数据块里时,解码就会出乱码。
正确的做法是传入{ stream: true }参数:
const decoder = new TextDecoder('utf-8');
let buffer = '';
while (true) {
const { done, value } = await reader.read();
if (done) break;
// stream: true 告诉解码器这是流式数据,会保留未完成的字节序列
const chunk = decoder.decode(value, { stream: true });
buffer += chunk;
// 按行分割处理SSE格式
const lines = buffer.split('\n');
buffer = lines.pop() || ''; // 最后一行可能不完整,留着下次处理
for (const line of lines) {
if (line.startsWith('data: ')) {
const data = line.slice(6);
if (data === '[DONE]') {
// 流结束
return;
}
// 处理收到的数据
console.log(data);
}
}
}
{ stream: true }是关键——它告诉TextDecoder这是一个持续的数据流,如果一个多字节字符被切分到了两个chunk里,解码器会妥善处理。
⑥ 实时渲染效果:从原始字节到页面动态展示
现在把接收到的数据渲染到页面上。
<div id="output" style="min-height:100px;padding:16px;border:1px solid #ddd;"></div>
<script>
async function chatWithStream() {
const output = document.getElementById('output');
output.textContent = '';
const response = await fetch('/chat', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ messages: [{ role: 'user', content: '你好' }] })
});
const reader = response.body.getReader();
const decoder = new TextDecoder();
let buffer = '';
let fullText = '';
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 data = line.slice(6);
if (data === '[DONE]') return;
try {
const json = JSON.parse(data);
const content = json.choices?.[0]?.delta?.content || '';
if (content) {
fullText += content;
output.textContent = fullText; // 实时更新DOM
}
} catch (e) {
// 忽略非JSON数据
}
}
}
}
}
</script>
每次收到新的内容片段,就拼接到fullText里,然后更新DOM。页面上的文字就会像打字机一样逐字出现。
⑦ 完整案例:构建一个打字机效果的对话组件
用Vue3来实现一个完整的对话组件:
<template>
<div class="chat-container">
<div class="messages">
<div v-for="msg in messages" :key="msg.id" class="message">
<span class="role">{{ msg.role === 'user' ? '👤' : '🤖' }}</span>
<span class="content">{{ msg.content }}</span>
</div>
<div v-if="isStreaming" class="message streaming">
<span class="role">🤖</span>
<span class="content">{{ streamingContent }}</span>
<span class="cursor">|</span>
</div>
</div>
<div class="input-area">
<input v-model="inputText" @keyup.enter="sendMessage" placeholder="输入消息..." />
<button @click="sendMessage" :disabled="isStreaming">发送</button>
</div>
</div>
</template>
<script setup>
import { ref, nextTick } from 'vue';
const messages = ref([]);
const inputText = ref('');
const isStreaming = ref(false);
const streamingContent = ref('');
const sendMessage = async () => {
if (!inputText.value.trim() || isStreaming.value) return;
const userMsg = { id: Date.now(), role: 'user', content: inputText.value };
messages.value.push(userMsg);
const question = inputText.value;
inputText.value = '';
isStreaming.value = true;
streamingContent.value = '';
try {
const response = await fetch('/chat', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ messages: [{ role: 'user', content: question }] })
});
const reader = response.body.getReader();
const decoder = new TextDecoder();
let buffer = '';
let fullText = '';
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 data = line.slice(6);
if (data === '[DONE]') {
messages.value.push({
id: Date.now(),
role: 'assistant',
content: fullText
});
isStreaming.value = false;
streamingContent.value = '';
return;
}
try {
const json = JSON.parse(data);
const content = json.choices?.[0]?.delta?.content || '';
if (content) {
fullText += content;
streamingContent.value = fullText;
await nextTick(); // 确保DOM更新
}
} catch (e) { /* 忽略非JSON */ }
}
}
}
} catch (error) {
console.error('流式请求失败:', error);
isStreaming.value = false;
}
};
</script>
这个组件实现了完整的对话流式输出:用户发送消息后,AI的回复会逐字显示,还有一个闪烁的光标模拟打字效果。
⑧ 常见报错排查:连接中断与数据格式异常处理
报错1:ERR_INVALID_CHUNKED_ENCODING
原因:在设置响应头之前就调用了res.write()。解决方案:确保writeHead在write之前。
报错2:前端收不到数据,但服务端日志显示已发送
大概率是Nginx等反向代理在缓冲响应。需要在Nginx配置中关闭缓冲:
location /stream {
proxy_pass http://backend;
proxy_buffering off; # 关闭缓冲
proxy_cache off; # 关闭缓存
proxy_set_header Connection '';
proxy_http_version 1.1;
chunked_transfer_encoding off;
}
SSE连接需要长时间保持,Nginx的proxy_read_timeout建议设置1小时以上。
报错3:中文乱码
确保服务端响应头指定了charset=utf-8,前端TextDecoder也指定'utf-8',并且解码时使用了{ stream: true }。
报错4:连接意外断开
SSE协议自带自动重连机制,但如果使用Fetch方式,需要手动实现重连逻辑:
async function fetchWithRetry(url, maxRetries = 3) {
for (let i = 0; i < maxRetries; i++) {
try {
await fetchStream(url);
break;
} catch (e) {
if (i === maxRetries - 1) throw e;
await new Promise(r => setTimeout(r, 1000 * (i + 1)));
}
}
}
⑨ 性能优化技巧:降低延迟与控制背压策略
优化1:减少首字延迟
首字延迟(Time to First Token)是流式体验的核心指标。服务端要尽早flush数据,不要等攒够了再发。每生成一小块数据就立即调用flush()。
优化2:控制背压(Backpressure)
当服务端推送数据的速度超过前端处理速度时,数据会在内存中堆积,可能导致内存溢出。ReadableStream提供了desiredSize属性来判断内部队列是否已满:
const reader = response.body.getReader();
// 可以通过 reader.closed 等状态判断流的状态
// 如果前端处理不过来,read()方法会自动产生背压,减慢数据读取速度
更精细的控制可以使用TransformStream手动管理队列。
优化3:减少DOM操作频率
每次收到数据都直接操作DOM会引发频繁的重排重绘。可以用requestAnimationFrame来批量更新:
let pendingUpdate = false;
let displayText = '';
function updateDisplay(text) {
displayText = text;
if (!pendingUpdate) {
pendingUpdate = true;
requestAnimationFrame(() => {
output.textContent = displayText;
pendingUpdate = false;
});
}
}
优化4:服务端避免缓冲
很多Web框架默认会缓冲响应。确保在服务端代码中显式调用flush(),并且在框架层面关闭响应缓冲。
⑩ 生产环境注意事项:超时设置与安全合规建议
超时设置
SSE连接是长连接,需要合理配置各层的超时时间:
- 应用层:Spring的
SseEmitter可以设置超时时间,如new SseEmitter(30000L)表示30秒无数据则超时 - 反向代理层:Nginx的
proxy_read_timeout要设得足够长,建议1小时以上 - 网关层:如果使用了API网关,Request Timeout要大于事件之间的空闲间隔
如果连接长期空闲,可以定期发送心跳(比如每15秒发一个data: ping\n\n)来保持连接活跃。
多实例部署的会话一致性
在Kubernetes等多Pod环境中,负载均衡器可能把SSE连接和后续请求分发到不同的Pod实例。解决方案:
- 使用会话亲和(Session Affinity),让同一客户端的请求始终路由到同一个Pod
- 使用Redis等外部存储共享连接状态
安全合规
- 认证鉴权:SSE连接通常需要携带Token。使用Fetch方式而非
EventSource,可以在请求头中携带Authorization - 内容过滤:如果数据来自大模型,服务端需要做内容安全审查,敏感内容不能推送给用户
- HTTPS:生产环境务必使用HTTPS,防止中间人攻击
- 连接数限制:每个SSE连接都会占用服务端资源,需要根据服务器规格设置合理的并发连接数上限
监控与告警
建议在生产环境对以下指标进行监控:
- SSE连接数(活跃连接数)
- 平均消息延迟(从生成到客户端收到的时间)
- 连接断开率
- 内存使用情况(流式数据可能积压)
WEB项目地址:演示地址
安卓APP下载地址:演示地址
以上就是从零开始实现流式输出的完整指南。从概念理解到代码落地,从开发调试到生产部署,希望这份指南能帮你少踩一些坑。
openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构
更多推荐

所有评论(0)