You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Node.js下如何以Hot Observable方式实现rxjs WebSocket Subject自动重连?

实现RxJS WebSocket自动重连并保持多播(Hot Observable)

问题根源

直接对webSocket()返回的Subject使用retry无效,因为该Subject在连接关闭(complete/error)后会进入终止状态,无法再发射任何值,retry重新订阅终止的Subject不会触发新连接。

解决方案思路

  1. 用defer包装WebSocket连接的创建逻辑,确保每次重连都生成新的WebSocket Subject。
  2. 结合retry处理连接错误(如网络故障、服务器不可达),repeatWhen处理正常关闭后的重连。
  3. 用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。
  • retry vs repeatWhen:retry处理连接过程中的错误事件,repeatWhen处理连接正常关闭的complete事件,两者结合覆盖所有需要重连的场景。
  • share的配置:通过resetOn*参数设置为false,确保重连后不会重置多播状态,所有现有订阅者自动切换到新连接。
  • 修正连接地址:原代码中ws://localhost://${port}存在多余冒号,应改为ws://localhost:${port}。

内容的提问来源于stack exchange,提问作者reaver lover

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.24 11:24:10