如何构造仅对首个订阅者生效的类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
相关产品推荐
相关产品推荐

