1. 引言
随着大语言模型应用的普及,流式输出已成为提升用户体验的关键能力。SSE(Server-Sent Events)作为一种基于 HTTP 的单向服务端推送协议,凭借其轻量、易用、天然支持文本流的特点,成为实现生成式 AI 对话输出的主流方案。本文将以 Node.js 生态为背景,分别介绍原生 Node.js、Express、Koa、Egg.js 和 Fastify 五种后端技术栈下的 SSE 流式会话实战实现,并补充 Next.js、Angular 和 React Native 三种前端框架的接收案例。
2. SSE 基础概念
SSE 是 HTML5 规范中定义的一种服务器推送技术,客户端通过 EventSource 接口建立长连接,服务端可以持续向客户端推送文本数据。与 WebSocket 不同,SSE 是单向通信,仅支持服务端向客户端推送,但实现简单、自动重连、基于 HTTP 协议,非常适合流式文本生成场景。
SSE 的核心响应格式为 text/event-stream,每条消息以 data: 开头,以两个换行符结束。客户端通过 EventSource 监听 onmessage 事件即可实时接收数据。
// 客户端基础用法
const eventSource = new EventSource('/api/chat/stream');
eventSource.onmessage = (event) => {
console.log('收到数据:', event.data);
};
3. 原生 Node.js 实现
原生 Node.js 不依赖任何框架,通过 http 模块即可实现 SSE 流式输出。核心思路是设置正确的响应头,然后分块写入数据。
const http = require('http');
const server = http.createServer((req, res) => {
if (req.url === '/api/chat/stream') {
res.writeHead(200, {
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
'Connection': 'keep-alive',
'Access-Control-Allow-Origin': '*'
});
// 模拟生成式输出
const chunks = ['你好', ',我是', 'AI 助手', ',很高兴', '认识你!'];
let index = 0;
const timer = setInterval(() => {
if (index < chunks.length) {
res.write(`data: ${JSON.stringify({ content: chunks[index++] })}\n\n`);
} else {
res.write('data: [DONE]\n\n');
clearInterval(timer);
res.end();
}
}, 200);
} else {
res.writeHead(404);
res.end();
}
});
server.listen(3000, () => {
console.log('SSE 服务已启动:http://localhost:3000');
});
上述代码中,服务端每 200 毫秒推送一个数据块,客户端通过 EventSource 即可实时接收。当所有数据推送完毕后,发送 [DONE] 标记并关闭连接。
4. Express 框架实现
Express 是 Node.js 最流行的 Web 框架,基于原生 HTTP 模块封装,实现 SSE 同样简单直接。通过路由处理函数直接操作响应对象即可。
const express = require('express');
const app = express();
app.get('/api/chat/stream', (req, res) => {
res.setHeader('Content-Type', 'text/event-stream');
res.setHeader('Cache-Control', 'no-cache');
res.setHeader('Connection', 'keep-alive');
res.flushHeaders();
const chunks = ['Express', '流式', '输出', '实战'];
let index = 0;
const timer = setInterval(() => {
if (index < chunks.length) {
res.write(`data: ${JSON.stringify({ content: chunks[index++] })}\n\n`);
} else {
res.write('data: [DONE]\n\n');
clearInterval(timer);
res.end();
}
}, 150);
});
app.listen(3000, () => {
console.log('Express SSE 服务已启动');
});
Express 实现 SSE 的关键在于 res.flushHeaders(),该方法会立即发送响应头,确保客户端能第一时间建立连接并开始接收数据。
5. Koa 框架实现
Koa 采用洋葱模型中间件机制,通过 ctx.res 访问底层响应对象。由于 Koa 默认将响应体作为整体返回,实现 SSE 时需要直接操作底层 Node 响应对象。
const Koa = require('koa');
const app = new Koa();
app.use(async (ctx) => {
if (ctx.path === '/api/chat/stream') {
ctx.set({
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
'Connection': 'keep-alive'
});
ctx.status = 200;
const res = ctx.res;
res.flushHeaders();
const chunks = ['Koa', '流式', '输出', '示例'];
let index = 0;
const timer = setInterval(() => {
if (index < chunks.length) {
res.write(`data: ${JSON.stringify({ content: chunks[index++] })}\n\n`);
} else {
res.write('data: [DONE]\n\n');
clearInterval(timer);
res.end();
}
}, 150);
} else {
ctx.status = 404;
ctx.body = 'Not Found';
}
});
app.listen(3000, () => {
console.log('Koa SSE 服务已启动');
});
Koa 中需要注意,一旦开始流式写入,就不能再通过 ctx.body 设置响应体,否则会与已写入的流数据冲突。
6. Egg.js 框架实现
Egg.js 是阿里开源的企业级 Node.js 框架,基于 Koa 封装,提供了完整的目录结构和插件机制。在 Egg.js 中实现 SSE,通常在 Controller 层直接操作响应对象。
// app/controller/chat.js
'use strict';
const Controller = require('egg').Controller;
class ChatController extends Controller {
async stream() {
const { ctx } = this;
ctx.set({
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
'Connection': 'keep-alive'
});
ctx.status = 200;
const res = ctx.res;
res.flushHeaders();
const chunks = ['Egg.js', '流式', '输出', '实战'];
let index = 0;
const timer = setInterval(() => {
if (index < chunks.length) {
res.write(`data: ${JSON.stringify({ content: chunks[index++] })}\n\n`);
} else {
res.write('data: [DONE]\n\n');
clearInterval(timer);
res.end();
}
}, 150);
}
}
module.exports = ChatController;
// app/router.js
'use strict';
module.exports = app => {
const { router, controller } = app;
router.get('/api/chat/stream', controller.chat.stream);
};
Egg.js 的 Controller 继承自 Koa 的上下文对象,因此实现方式与 Koa 基本一致,只需将逻辑放入 Controller 方法中,并通过路由配置暴露接口。
7. Fastify 框架实现
Fastify 以高性能著称,内置了完整的插件体系和 Schema 校验。Fastify 实现 SSE 时,通过 reply.raw 访问底层响应对象进行流式写入。
const fastify = require('fastify')();
fastify.get('/api/chat/stream', (request, reply) => {
reply.raw.writeHead(200, {
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
'Connection': 'keep-alive'
});
const chunks = ['Fastify', '流式', '输出', '示例'];
let index = 0;
const timer = setInterval(() => {
if (index < chunks.length) {
reply.raw.write(`data: ${JSON.stringify({ content: chunks[index++] })}\n\n`);
} else {
reply.raw.write('data: [DONE]\n\n');
clearInterval(timer);
reply.raw.end();
}
}, 150);
});
fastify.listen({ port: 3000 }, (err) => {
if (err) {
console.error(err);
process.exit(1);
}
console.log('Fastify SSE 服务已启动');
});
Fastify 中通过 reply.raw 获取 Node 原生响应对象,写入方式与其他框架一致。需要注意的是,使用 reply.raw 后应避免再调用 Fastify 的 reply.send() 方法。
8. 前端统一接收方案
无论后端使用哪种框架,前端都可以使用统一的 EventSource 接口接收流式数据。以下是一个通用的前端接收示例。
// 前端统一接收 SSE 数据
const eventSource = new EventSource('/api/chat/stream');
eventSource.onopen = () => {
console.log('SSE 连接已建立');
};
eventSource.onmessage = (event) => {
if (event.data === '[DONE]') {
console.log('流式输出结束');
eventSource.close();
return;
}
const data = JSON.parse(event.data);
console.log('收到内容:', data.content);
};
eventSource.onerror = (error) => {
console.error('SSE 连接异常:', error);
eventSource.close();
};
前端通过 EventSource 建立连接后,服务端推送的每条数据都会触发 onmessage 回调。当收到 [DONE] 标记时,主动关闭连接,完成整个流式会话。
9. React 实战案例
在 React 中接收 SSE 流式数据,推荐使用 useEffect 管理 EventSource 的生命周期,并结合 useState 实时更新界面。以下是一个完整的 React 组件示例。
import React, { useState, useEffect, useRef } from 'react';
function ChatStream() {
const [messages, setMessages] = useState([]);
const [status, setStatus] = useState('connecting');
const eventSourceRef = useRef(null);
useEffect(() => {
const eventSource = new EventSource('/api/chat/stream');
eventSourceRef.current = eventSource;
eventSource.onopen = () => {
setStatus('connected');
};
eventSource.onmessage = (event) => {
if (event.data === '[DONE]') {
setStatus('done');
eventSource.close();
return;
}
const data = JSON.parse(event.data);
setMessages((prev) => [...prev, data.content]);
};
eventSource.onerror = (error) => {
console.error('SSE 连接异常:', error);
setStatus('error');
eventSource.close();
};
return () => {
eventSource.close();
};
}, []);
return (
<div>
<p>连接状态:{status}</p>
<div>
{messages.map((msg, index) => (
<span key={index}>{msg}</span>
))}
</div>
</div>
);
}
export default ChatStream;
上述组件在挂载时建立 SSE 连接,收到数据后追加到 messages 数组,组件卸载时自动关闭连接,避免内存泄漏。使用 useRef 保存 EventSource 实例,便于在清理函数中关闭。
10. Vue 实战案例
Vue 3 中可以使用组合式 API 封装 SSE 接收逻辑,通过 ref 管理响应式数据,并在 onUnmounted 中关闭连接。以下是一个完整的 Vue 组件示例。
<template>
<div>
<p>连接状态:{{ status }}</p>
<div>
<span v-for="(msg, index) in messages" :key="index">{{ msg }}</span>
</div>
</div>
</template>
<script setup>
import { ref, onMounted, onUnmounted } from 'vue';
const messages = ref([]);
const status = ref('connecting');
let eventSource = null;
onMounted(() => {
eventSource = new EventSource('/api/chat/stream');
eventSource.onopen = () => {
status.value = 'connected';
};
eventSource.onmessage = (event) => {
if (event.data === '[DONE]') {
status.value = 'done';
eventSource.close();
return;
}
const data = JSON.parse(event.data);
messages.value.push(data.content);
};
eventSource.onerror = (error) => {
console.error('SSE 连接异常:', error);
status.value = 'error';
eventSource.close();
};
});
onUnmounted(() => {
if (eventSource) {
eventSource.close();
}
});
</script>
Vue 组件通过 onMounted 建立连接,onUnmounted 清理资源,确保组件销毁时不会遗留未关闭的连接。响应式数组 messages 会自动驱动视图更新,实现流式文本的实时渲染。
11. Next.js 实战案例
Next.js 作为 React 的全栈框架,在 App Router 模式下可以通过客户端组件配合 EventSource 接收 SSE 流式数据。以下是一个完整的 Next.js 客户端组件示例。
'use client';
import { useState, useEffect, useRef } from 'react';
export default function ChatStream() {
const [messages, setMessages] = useState([]);
const [status, setStatus] = useState('connecting');
const eventSourceRef = useRef(null);
useEffect(() => {
const eventSource = new EventSource('/api/chat/stream');
eventSourceRef.current = eventSource;
eventSource.onopen = () => {
setStatus('connected');
};
eventSource.onmessage = (event) => {
if (event.data === '[DONE]') {
setStatus('done');
eventSource.close();
return;
}
const data = JSON.parse(event.data);
setMessages((prev) => [...prev, data.content]);
};
eventSource.onerror = (error) => {
console.error('SSE 连接异常:', error);
setStatus('error');
eventSource.close();
};
return () => {
eventSource.close();
};
}, []);
return (
<div>
<p>连接状态:{status}</p>
<div>
{messages.map((msg, index) => (
<span key={index}>{msg}</span>
))}
</div>
</div>
);
}
Next.js 组件与 React 组件写法基本一致,关键区别在于文件顶部需要添加 'use client' 指令,表明这是一个客户端组件。在 App Router 中,SSE 连接逻辑应放在客户端组件中执行,服务端组件无法建立长连接。
12. Angular 实战案例
Angular 中可以通过 HttpClient 的 text/event-stream 响应类型接收 SSE 数据,也可以直接使用浏览器原生 EventSource。以下是一个使用 EventSource 的 Angular 服务示例。
import { Injectable } from '@angular/core';
import { Observable } from 'rxjs';
@Injectable({ providedIn: 'root' })
export class SseService {
connect(url: string): Observable<string> {
return new Observable((subscriber) => {
const eventSource = new EventSource(url);
eventSource.onopen = () => {
console.log('SSE 连接已建立');
};
eventSource.onmessage = (event) => {
if (event.data === '[DONE]') {
subscriber.complete();
eventSource.close();
return;
}
const data = JSON.parse(event.data);
subscriber.next(data.content);
};
eventSource.onerror = (error) => {
console.error('SSE 连接异常:', error);
subscriber.error(error);
eventSource.close();
};
return () => {
eventSource.close();
};
});
}
}
在 Angular 组件中,可以通过 Subscription 订阅该服务返回的 Observable,并在组件销毁时取消订阅,避免内存泄漏。使用 RxJS 的 Observable 封装 SSE 接收逻辑,可以更好地与 Angular 的响应式编程模型结合。
13. React Native 实战案例
React Native 中同样支持 EventSource API,可以用于接收 SSE 流式数据。以下是一个完整的 React Native 组件示例,使用 useEffect 管理连接生命周期。
import React, { useState, useEffect, useRef } from 'react';
import { View, Text, StyleSheet } from 'react-native';
function ChatStream() {
const [messages, setMessages] = useState([]);
const [status, setStatus] = useState('connecting');
const eventSourceRef = useRef(null);
useEffect(() => {
const eventSource = new EventSource('/api/chat/stream');
eventSourceRef.current = eventSource;
eventSource.onopen = () => {
setStatus('connected');
};
eventSource.onmessage = (event) => {
if (event.data === '[DONE]') {
setStatus('done');
eventSource.close();
return;
}
const data = JSON.parse(event.data);
setMessages((prev) => [...prev, data.content]);
};
eventSource.onerror = (error) => {
console.error('SSE 连接异常:', error);
setStatus('error');
eventSource.close();
};
return () => {
eventSource.close();
};
}, []);
return (
<View style={styles.container}>
<Text>连接状态:{status}</Text>
<View>
{messages.map((msg, index) => (
<Text key={index}>{msg}</Text>
))}
</View>
</View>
);
}
const styles = StyleSheet.create({
container: {
flex: 1,
padding: 16
}
});
export default ChatStream;
React Native 的 SSE 实现与 Web 端 React 基本一致,主要区别在于使用 View、Text 等原生组件替代 HTML 标签,并通过 StyleSheet 管理样式。需要注意的是,在移动端网络环境不稳定时,应结合 EventSource 的自动重连机制做好状态提示。
14. 实战对比与选型建议
八种技术栈实现 SSE 的核心思路完全一致,差异主要体现在框架的 API 风格和项目组织方式上。以下是各方案的对比总结。
| 框架 | 实现难度 | 性能表现 | 适用场景 |
|---|---|---|---|
| 原生 Node.js | 中等 | 高 | 轻量服务、学习原理 |
| Express | 低 | 较高 | 中小型项目、快速开发 |
| Koa | 低 | 较高 | 中间件丰富的项目 |
| Egg.js | 低 | 高 | 企业级应用、团队协作 |
| Fastify | 低 | 极高 | 高性能 API 服务 |
| Next.js | 低 | 高 | 全栈应用、服务端渲染 |
| Angular | 中 | 较高 | 大型企业级前端应用 |
| React Native | 中 | 较高 | 移动端跨平台应用 |
在实际项目中,如果追求开发效率和生态成熟度,Express 和 Koa 是不错的选择;如果是企业级应用,需要完整的工程化支持,Egg.js 或 Angular 更为合适;如果对性能有极致要求,Fastify 是首选;如果需要全栈能力,Next.js 可以同时承担前后端职责;如果是移动端场景,React Native 则是最佳选择。
15. 总结
本文通过八个完整的代码示例,详细介绍了 Node.js 生态下使用原生 Node.js、Express、Koa、Egg.js 和 Fastify 实现 SSE 流式会话的方法,并补充了 Next.js、Angular 和 React Native 三种前端框架的实战案例。所有方案都遵循相同的协议规范,核心在于正确设置响应头并分块写入数据。前端使用统一的 EventSource 接口即可无缝对接。在实际开发中,可以根据项目规模、团队技术栈和性能要求灵活选择合适的技术方案。
转载自 CSDN-专业IT技术社区
原文链接:https://blog.csdn.net/qq_45393395/article/details/167360600



