流式输出(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配合ReadableStreamEventSource虽然用起来简单,但它只能发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()。解决方案:确保writeHeadwrite之前。

报错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实例。解决方案:

  1. 使用会话亲和(Session Affinity),让同一客户端的请求始终路由到同一个Pod
  2. 使用Redis等外部存储共享连接状态

安全合规

  1. 认证鉴权:SSE连接通常需要携带Token。使用Fetch方式而非EventSource,可以在请求头中携带Authorization
  2. 内容过滤:如果数据来自大模型,服务端需要做内容安全审查,敏感内容不能推送给用户
  3. HTTPS:生产环境务必使用HTTPS,防止中间人攻击
  4. 连接数限制:每个SSE连接都会占用服务端资源,需要根据服务器规格设置合理的并发连接数上限

监控与告警

建议在生产环境对以下指标进行监控:

  • SSE连接数(活跃连接数)
  • 平均消息延迟(从生成到客户端收到的时间)
  • 连接断开率
  • 内存使用情况(流式数据可能积压)

WEB项目地址:演示地址
安卓APP下载地址:演示地址
以上就是从零开始实现流式输出的完整指南。从概念理解到代码落地,从开发调试到生产部署,希望这份指南能帮你少踩一些坑。

Logo

openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构

更多推荐