Rx.NET:按序合并可观察对象并按需终止历史订阅
这是Rx.NET中非常典型的「历史重放+实时推送」合并场景,核心挑战在于仅通过IEquatable<T>匹配值来定位同步点,并在同步完成后自动切换到实时流、干净释放历史流的订阅资源。我针对你的需求实现了泛型方法MergeObservables<T>,完全覆盖你列出的6个测试场景,具体如下:
核心实现思路
- 缓存实时流:用
Replay操作符缓存实时流的所有发射值,支持后续回溯访问(适配实时流先于历史流启动的场景)。 - 定位同步点:以实时流的第一个值为同步标记,监听历史流直到发射出匹配的项。
- 跟踪已发射项:用
HashSet<T>记录历史流已输出的所有值,用于过滤实时流中重复的旧值。 - 自动切换与资源释放:同步完成后,立即终止历史流的订阅,切换到实时流的新值,并清理不必要的资源。
完整代码实现
using System; using System.Collections.Generic; using System.Reactive.Linq; using System.Reactive.Subjects; public static class ObservableMergeExtensions { public static IObservable<T> MergeObservables<T>( IConnectableObservable<T> historicalStream, IConnectableObservable<T> realtimeStream) where T : IEquatable<T> { // 缓存实时流的所有发射值,确保可以回溯访问 var realtimeCache = realtimeStream.Replay(); var realtimeConnection = realtimeCache.Connect(); // 获取实时流的第一个值(若存在),作为同步点标记 var firstRealtimeValue = realtimeCache.FirstOrDefaultAsync().PublishLast(); firstRealtimeValue.Connect(); // 跟踪历史流已发射的所有项,用于过滤实时流的重复值 var historicalItems = new HashSet<T>(); var historicalTrackingSubject = new Subject<T>(); // 订阅历史流,收集已发射项,直到到达同步点 var historicalUntilSync = historicalStream .Do(item => { if (!historicalItems.Contains(item)) historicalItems.Add(item); historicalTrackingSubject.OnNext(item); }) .TakeUntil(item => firstRealtimeValue.HasValue && item.Equals(firstRealtimeValue.Value)) .Publish(); var historicalConnection = historicalUntilSync.Connect(); // 同步完成后处理实时流:跳过已在历史流出现的旧值,清理历史流资源 var postSyncRealtime = realtimeCache .SkipWhile(item => historicalItems.Contains(item)) .Do(_ => { historicalConnection.Dispose(); historicalTrackingSubject.OnCompleted(); historicalTrackingSubject.Dispose(); }) .Finally(() => realtimeConnection.Dispose()); // 合并历史流(到同步点)与同步后的实时流 return historicalUntilSync.Concat(postSyncRealtime) .Catch((Exception ex) => { // 异常场景下的资源清理 historicalConnection.Dispose(); realtimeConnection.Dispose(); historicalTrackingSubject.OnError(ex); historicalTrackingSubject.Dispose(); return Observable.Throw<T>(ex); }); } }
测试场景适配说明
针对你列出的6个场景,这个实现的适配逻辑如下:
历史流未追上实时流时不输出实时值:
实时流的缓存会被SkipWhile过滤,直到历史流的historicalItems集合包含足够的项,只有当历史流到达同步点后,未被历史流覆盖的实时值才会被发射。历史流到达实时流首个值时,立即输出所有已推送的实时值并断开历史订阅:
当历史流发射到与实时流首个值匹配的项时,TakeUntil会终止历史流的订阅;随后Concat切换到postSyncRealtime,会先发射实时缓存中未被历史流包含的所有值,同时立即释放历史流的连接资源。支持实时值先于历史值出现的场景:
实时流通过Replay操作符缓存所有发射值,即使历史流还未启动,后续历史流到达同步点时,依然可以访问到之前的实时值。跳过实时流中已输出过的旧值,直到出现新值:
historicalItems集合会记录历史流已输出的所有项,postSyncRealtime中的SkipWhile会跳过实时缓存中已存在于该集合的项,只发射未出现过的新值。若实时流为历史流的子集,需持续等待匹配:
TakeUntil会一直监听历史流,直到出现与实时流首个值匹配的项,不会提前终止历史流的订阅,直到同步点出现。同步后忽略历史流的后续差异:
当历史流到达同步点时,TakeUntil会立即终止历史流的订阅,后续历史流的任何发射都不会被处理,相关资源也会被及时释放。
内容的提问来源于stack exchange,提问作者Paul Milla

