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

如何构造仅对首个订阅者生效的类ReplaySubject的Rx Observable?

实现方案

下面分别给出.NET(System.Reactive)和RxJS两种环境下的实现,核心逻辑完全对齐你的需求:仅首个订阅者能拿到订阅前的所有累计数据,后续所有订阅者的行为和普通Subject完全一致。

核心逻辑说明

  • 没有订阅者时,缓存所有下发的消息
  • 首个订阅者接入时,先把全量缓存数据推给它,之后立刻清空缓存,不再缓存任何后续消息
  • 从首个订阅者接入完成开始,所有新老订阅者都只能拿到订阅后下发的新消息,不会再拿到历史数据

.NET 实现

基于SubjectBase基类封装,和你之前实现DistinctSubject的思路完全兼容:

public class FirstSubscriberReplaySubject<T> : SubjectBase<T>
{
    private readonly Subject<T> _innerSubject = new();
    private readonly List<T> _cache = new();
    private bool _firstSubscribed = false;
    private bool _disposed = false;

    public override void OnNext(T value)
    {
        if (_disposed) throw new ObjectDisposedException(nameof(FirstSubscriberReplaySubject<T>));
        lock (_cache)
        {
            if (!_firstSubscribed)
            {
                _cache.Add(value);
            }
        }
        _innerSubject.OnNext(value);
    }

    public override void OnError(Exception error) => _innerSubject.OnError(error);

    public override void OnCompleted() => _innerSubject.OnCompleted();

    public override IDisposable Subscribe(IObserver<T> observer)
    {
        if (_disposed) throw new ObjectDisposedException(nameof(FirstSubscriberReplaySubject<T>));
        
        lock (_cache)
        {
            if (!_firstSubscribed)
            {
                // 给首个订阅者推送全量缓存
                foreach (var item in _cache)
                {
                    observer.OnNext(item);
                }
                _cache.Clear();
                _firstSubscribed = true;
            }
        }
        return _innerSubject.Subscribe(observer);
    }

    protected override void Dispose(bool disposing)
    {
        if (disposing)
        {
            _disposed = true;
            _innerSubject.Dispose();
            lock (_cache)
            {
                _cache.Clear();
            }
        }
        base.Dispose(disposing);
    }

    public override bool HasObservers => _innerSubject.HasObservers;
}

如果你的使用场景没有并发订阅的需求,可以去掉lock逻辑进一步提升性能。

RxJS 实现

直接继承RxJS原生Subject类,逻辑和.NET版本完全对齐:

import { Subject, Observer } from 'rxjs';

export class FirstSubscriberReplaySubject<T> extends Subject<T> {
  private cache: T[] = [];
  private firstSubscribed = false;

  next(value: T): void {
    if (!this.firstSubscribed) {
      this.cache.push(value);
    }
    super.next(value);
  }

  subscribe(observerOrNext?: Partial<Observer<T>> | ((value: T) => void)) {
    const subscription = super.subscribe(observerOrNext as any);
    
    if (!this.firstSubscribed) {
      // 给首个订阅者推送全量缓存
      this.cache.forEach(item => {
        if (typeof observerOrNext === 'function') {
          observerOrNext(item);
        } else if (observerOrNext?.next) {
          observerOrNext.next(item);
        }
      });
      this.cache = [];
      this.firstSubscribed = true;
    }

    return subscription;
  }
}

边界行为说明

  • 首个订阅者刚接入就取消订阅,也不会再给后续订阅者推送历史缓存
  • 首个订阅者接入前如果触发了OnError/OnCompleted,首个订阅者会先收到对应事件,再接收缓存数据,和普通Subject的事件优先级保持一致

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 13:27:01