如何实现CountLiveStreams?统计活跃并发内部流数量
实现
CountLiveStreams 扩展方法 这个需求在Rx场景里挺常见的——本质就是要追踪每个内部流的完整生命周期:新流激活时计数加一,流完成或出错时计数减一。咱们可以用Rx的操作符组合出优雅的实现:
using System.Reactive.Linq; using System.Reactive.Subjects; public static class ObservableExtensions { public static IObservable<int> CountLiveStreams(this IObservable<IObservable<Unit>> source) { return source // 把每个内部流转换成"激活信号+结束信号"的序列 .SelectMany(innerStream => Observable.Return(1) // 流激活时发送+1 .Concat( innerStream .IgnoreElements() // 不关心内部流的实际数据,只等它结束 .Catch((Exception _) => Observable.Empty<Unit>()) // 捕获错误,避免中断计数流 .Concat(Observable.Return(-1)) // 流结束(完成/出错)时发送-1 ) ) // 用Scan维护实时计数,初始值为0 .Scan(0, (currentCount, delta) => currentCount + delta) // 兜底:确保计数不会出现负数(理论上不会,但防极端情况) .Select(count => Math.Max(count, 0)); } }
关键逻辑拆解
- SelectMany处理单流生命周期:
- 每个内部流刚进来时,先发送
+1标记它已活跃; - 忽略内部流的所有
OnNext事件(我们只关心它的结束状态),用Catch吃掉内部流的错误(避免单个内部流的异常搞崩整个计数流),最后在流结束时发送-1标记它已终止。
- 每个内部流刚进来时,先发送
- Scan维护实时计数:
从0开始,每次收到+1或-1就更新当前计数,完美反映活跃流的实时数量。 - 兜底处理:
用Math.Max(count, 0)防止极端情况下出现负数计数,让逻辑更健壮。
测试示例
var source = new Subject<IObservable<Unit>>(); IObservable<int> count = source.CountLiveStreams(); // 订阅计数流,打印实时数值 count.Subscribe(currentCount => Console.WriteLine($"当前活跃流数量:{currentCount}")); // 模拟场景:发送第一个内部流 var inner1 = new Subject<Unit>(); source.OnNext(inner1); // 输出:当前活跃流数量:1 // 发送第二个内部流 var inner2 = new Subject<Unit>(); source.OnNext(inner2); // 输出:当前活跃流数量:2 // 结束第一个内部流 inner1.OnCompleted(); // 输出:当前活跃流数量:1 // 第二个内部流出错终止 inner2.OnError(new Exception("测试错误")); // 输出:当前活跃流数量:0
额外提示
如果你的内部流存在被主动取消订阅的场景(不是自然完成/出错),那可以结合TakeUntil和Disposable来追踪取消事件,确保取消时也能触发计数减一。不过大多数业务场景下,完成/出错已经是流结束的主要方式了。
内容的提问来源于stack exchange,提问作者James L
相关产品推荐
相关产品推荐

