百度一下吧头像
关注

SSE 流式会话实战:Node.js 与主流框架的生成式输出实现

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

文章来源转载

评论

赞0

评论列表

微信小程序
QQ小程序

关于作者

点赞数:0
关注数:0
粉丝:0
文章:0
关注标签:0
加入于:--