如何在Express中将WebSocket数据转发至Server-Sent Event?
实现Express通过SSE转发外部WebSocket消息到客户端
问题核心
需要搭建一条实时数据转发链路:客户端 ←SSE→ Express服务器 ←WebSocket→ 外部服务器,解决异步连接间的数据实时传递问题。
1. 改造SocketClient类
原类存在事件绑定错误、未处理连接就绪状态的问题,咱们把它改成事件驱动结构,方便Express路由实时接收消息:
const { EventEmitter } = require('events'); const WebSocket = require('ws'); // Node环境用ws包,浏览器环境直接用原生WebSocket class SocketClient extends EventEmitter { #socket; #isConnected = false; constructor() { super(); this.#socket = new WebSocket("ws://example.com"); // 连接成功后再发送初始消息 this.#socket.addEventListener("open", () => { this.#isConnected = true; this.#socket.send("something"); }); // 收到WebSocket消息时,触发自定义message事件 this.#socket.addEventListener("message", (event) => { this.emit("message", event.data); }); // 连接关闭时通知外部 this.#socket.addEventListener("close", () => { this.#isConnected = false; this.emit("close"); }); // 错误处理 this.#socket.addEventListener("error", (err) => { this.emit("error", err); }); } get isConnected() { return this.#isConnected; } // 主动关闭WebSocket连接 close() { if (this.#isConnected) { this.#socket.close(); } } }
2. Express路由实现SSE转发
设置正确的SSE响应头,通过事件监听实时转发WebSocket消息,同时处理客户端断开的情况:
const express = require('express'); const app = express(); app.get('/', (req, res) => { // 配置SSE响应头,确保连接保持活跃且不缓存 res.writeHead(200, { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache', 'Connection': 'keep-alive' }); const socketClient = new SocketClient(); // 转发WebSocket消息为SSE格式 const handleMessage = (data) => { // SSE规范要求每条消息以`data: `开头,结尾用`\n\n`分隔 res.write(`data: ${JSON.stringify(data)}\n\n`); }; // WebSocket连接关闭时,结束SSE响应 const handleClose = () => { res.end(); }; // WebSocket出错时,发送错误信息并结束连接 const handleError = (err) => { res.write(`data: ${JSON.stringify({ error: err.message })}\n\n`); res.end(); }; // 绑定事件监听 socketClient.on('message', handleMessage); socketClient.on('close', handleClose); socketClient.on('error', handleError); // 客户端断开请求时,清理WebSocket连接和事件监听 req.on('close', () => { socketClient.off('message', handleMessage); socketClient.off('close', handleClose); socketClient.off('error', handleError); socketClient.close(); }); }); app.listen(3000, () => { console.log('服务器运行在3000端口'); });
3. 为什么不用isDataAvailable和getNextDataChunk?
原代码里的while await循环模式不适合实时流场景:
- 这种同步轮询会阻塞Node.js事件循环,导致服务器无法处理其他请求
- WebSocket消息是异步触发的,轮询方式无法实时响应新消息
如果偏好使用异步迭代器(支持for await...of),可以给SocketClient添加异步迭代器实现:
class SocketClient extends EventEmitter { // ... 保留之前的代码 ... #messageQueue = []; #resolve; constructor() { super(); // ... 原有的事件监听 ... this.#socket.addEventListener("message", (event) => { if (this.#resolve) { this.#resolve(event.data); this.#resolve = null; } else { this.#messageQueue.push(event.data); } }); } async #getNextMessage() { if (this.#messageQueue.length > 0) { return this.#messageQueue.shift(); } return new Promise(resolve => { this.#resolve = resolve; }); } // 实现异步迭代器接口 [Symbol.asyncIterator]() { return { async next() { const data = await this.#getNextMessage(); return { value: data, done: false }; } }; } }
对应的Express路由可以这样写:
app.get('/', async (req, res) => { res.writeHead(200, { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache', 'Connection': 'keep-alive' }); const socketClient = new SocketClient(); let isClientClosed = false; // 监听客户端断开 req.on('close', () => { isClientClosed = true; socketClient.close(); }); try { // 用for await...of遍历WebSocket消息 for await (const data of socketClient) { if (isClientClosed) break; res.write(`data: ${JSON.stringify(data)}\n\n`); } } catch (err) { res.write(`data: ${JSON.stringify({ error: err.message })}\n\n`); } finally { res.end(); } });
内容的提问来源于stack exchange,提问作者goose_lake
相关产品推荐
相关产品推荐

