Node.js下如何以Hot Observable方式实现rxjs WebSocket Subject自动重连?
实现RxJS WebSocket自动重连并保持多播(Hot Observable)
问题根源
直接对webSocket()返回的Subject使用retry无效,因为该Subject在连接关闭(complete/error)后会进入终止状态,无法再发射任何值,retry重新订阅终止的Subject不会触发新连接。
解决方案思路
- 用
defer包装WebSocket连接的创建逻辑,确保每次重连都生成新的WebSocket Subject。 - 结合
retry处理连接错误(如网络故障、服务器不可达),repeatWhen处理正常关闭后的重连。 - 用
share将流转为Hot Observable,让所有订阅者共享同一个连接实例,重连时自动切换到新连接。
完整代码实现
global.WebSocket = require("ws"); const { webSocket } = require("rxjs/webSocket"); const { defer, retry, repeatWhen, delay, tap, share, filter } = require("rxjs"); // 创建可重连的多播WebSocket Observable const createReconnectingWs = (wsUrl) => { return defer(() => { console.log(`正在连接到 ${wsUrl}`); return webSocket(wsUrl); }).pipe( // 处理连接错误(如服务器不可达、网络中断),延迟1.5秒重试 retry({ delay: (error, retryCount) => { console.log(`连接失败(第${retryCount+1}次重试):${error.message}`); return 1500; } }), // 处理连接正常关闭后的重连,延迟1.5秒 repeatWhen((completeNotifications) => { return completeNotifications.pipe( tap(() => console.log("连接已正常关闭,准备重连")), delay(1500) ); }), // 转为Hot Observable,所有订阅者共享同一连接,重连后保持多播状态 share({ resetOnError: false, resetOnComplete: false, resetOnRefCountZero: false }) ); }; // 使用示例 const port = 8080; const ws$ = createReconnectingWs(`ws://localhost:${port}`); // 多个订阅者共享同一个连接 ws$.subscribe(msg => console.log("订阅者1收到消息:", msg)); ws$.subscribe(msg => console.log("订阅者2收到消息:", msg)); // 基于主流创建派生Observable,依然共享连接 const connectedMsg$ = ws$.pipe( filter(msg => JSON.parse(msg) === "You are connected") ); connectedMsg$.subscribe(msg => console.log("连接确认消息:", msg));
关键细节说明
defer的作用:延迟创建WebSocket Subject,每次重连时都会生成新的实例,避免使用已终止的旧Subject。retryvsrepeatWhen:retry处理连接过程中的错误事件,repeatWhen处理连接正常关闭的complete事件,两者结合覆盖所有需要重连的场景。share的配置:通过resetOn*参数设置为false,确保重连后不会重置多播状态,所有现有订阅者自动切换到新连接。- 修正连接地址:原代码中
ws://localhost://${port}存在多余冒号,应改为ws://localhost:${port}。
内容的提问来源于stack exchange,提问作者reaver lover
相关产品推荐
相关产品推荐

