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

如何实现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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.08 08:12:43