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

如何用RxJS实现Server Side Event的按需订阅与自动断开?

实现多订阅者共享Server-Sent Events(SSE)连接的RxJS方案

核心实现逻辑

通过维护订阅计数器结合RxJS的多播能力,精准控制SSE连接的生命周期:

  • 首次订阅时创建SSE连接并启动监听
  • 多订阅者复用同一连接,仅维护一个活跃的SSE实例
  • 最后一个订阅者取消订阅时,主动关闭SSE连接
  • 新订阅触发时自动重建连接

完整代码实现

import { Observable, Subject } from 'rxjs';
import { multicast, refCount } from 'rxjs/operators';

function createSharedSSEObservable(url: string): Observable<MessageEvent> {
  let subscriberCount = 0;
  let eventSource: EventSource | null = null;

  return new Observable<MessageEvent>((subscriber) => {
    subscriberCount++;

    // 首次订阅初始化SSE连接
    if (subscriberCount === 1) {
      eventSource = new EventSource(url);

      // 转发SSE消息到所有订阅者
      eventSource.onmessage = (event) => subscriber.next(event);

      // 错误处理:通知订阅者并清理连接
      eventSource.onerror = (error) => {
        subscriber.error(error);
        cleanup();
      };
    }

    // 订阅取消时的清理逻辑
    return () => {
      subscriberCount--;
      // 最后一个订阅者取消,关闭SSE连接
      if (subscriberCount === 0) {
        cleanup();
      }
    };

    // 统一关闭SSE连接的函数
    function cleanup() {
      eventSource?.close();
      eventSource = null;
    }
  }).pipe(
    // 用Subject多播数据流,refCount自动管理连接的建立/销毁
    multicast(() => new Subject<MessageEvent>()),
    refCount()
  );
}

代码关键点解析

  • 订阅计数器subscriberCount:精准控制SSE连接的创建与销毁时机,避免不必要的连接开销
  • multicast + refCount:实现数据流共享,refCount会在订阅数从0变为1时自动激活Observable,从1变为0时自动终止
  • cleanup函数:统一处理SSE连接关闭,防止内存泄漏

使用示例

// 创建共享的SSE数据流实例
const sseStream$ = createSharedSSEObservable('/api/stream');

// 订阅者X:首次订阅,触发SSE连接建立
const subX = sseStream$.subscribe({
  next: e => console.log('X接收:', e.data),
  error: err => console.error('X错误:', err)
});

// 订阅者Y:复用已有的SSE连接
const subY = sseStream$.subscribe({
  next: e => console.log('Y接收:', e.data),
  error: err => console.error('Y错误:', err)
});

// X取消订阅:仍有Y活跃,SSE连接保持
subX.unsubscribe();

// Y取消订阅:最后一个订阅者,SSE连接关闭
subY.unsubscribe();

// 订阅者Z:触发新的SSE连接建立
const subZ = sseStream$.subscribe({
  next: e => console.log('Z接收:', e.data),
  error: err => console.error('Z错误:', err)
});

额外注意事项

  • 自定义重连:默认EventSource会自动重连,若需自定义重连间隔或失败策略,可在onerror回调中添加延迟重试逻辑
  • 组件生命周期适配:在框架中使用时(如React/Angular),需确保组件销毁时正确取消订阅,避免内存泄漏
  • 类型安全:TypeScript环境下可为MessageEvent指定泛型,明确返回数据的结构

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 07:15:50