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

如何在Rx.NET测试中同时观测唯一值与重复值?

如何同时观测Rx.NET中Distinct操作的去重流与重复值流

要同时获取Distinct()过滤后的结果和被剔除的重复值,核心是共享源Observable,然后分别处理两个分支——一个分支保留首次出现的元素,另一个分支筛选出重复出现的元素。你的思路方向是对的(用Publish共享源),但需要调整处理重复项的逻辑,而不是直接用CombineLatest。

解决方案:共享源+跟踪已见元素筛选重复项

这里有两种简洁的实现方式,都能精准捕获你需要的重复值(OnNext(200,"a")和OnNext(221,"a")):

方法1:用Publish共享源,结合HashSet跟踪重复项

我们可以用Publish把源转换成可连接的Observable,确保所有订阅者共享同一个数据流。然后一个分支走Distinct(),另一个分支通过维护一个已见元素的集合,筛选出重复的项:

public class Test : ReactiveTest
{
    [Fact]
    public void Observe_distinct_and_duplicates()
    {
        var scheduler = new TestScheduler();
        var source = scheduler.CreateHotObservable(
            OnNext(100, "a"),
            OnNext(110, "b"),
            OnNext(200, "a"),
            OnNext(220, "c"),
            OnNext(221, "a")
        );

        // 共享源Observable,避免多次订阅导致重复执行
        var sharedSource = source.Publish();

        // 去重后的结果流
        var distinctResults = scheduler.CreateObserver<string>();
        sharedSource.Distinct().Subscribe(distinctResults);

        // 重复值流:筛选出已经出现过的元素
        var duplicateResults = scheduler.CreateObserver<string>();
        var seenElements = new HashSet<string>();
        sharedSource
            .Where(item => !seenElements.Add(item)) // Add返回false表示元素已存在
            .Subscribe(duplicateResults);

        // 启动共享源的数据流
        sharedSource.Connect();
        scheduler.AdvanceBy(1000);

        // 验证去重结果
        distinctResults.Messages.AssertEqual(
            OnNext(100, "a"),
            OnNext(110, "b"),
            OnNext(220, "c")
        );

        // 验证重复值结果
        duplicateResults.Messages.AssertEqual(
            OnNext(200, "a"),
            OnNext(221, "a")
        );
    }
}

方法2:用Scan一次性拆分出两个流

如果你想更“Rx式”地处理(避免使用外部集合),可以用Scan操作符维护一个包含已见元素和当前元素类型(新元素/重复元素)的状态,然后拆分出两个流:

public class Test : ReactiveTest
{
    [Fact]
    public void Observe_distinct_and_duplicates_with_scan()
    {
        var scheduler = new TestScheduler();
        var source = scheduler.CreateHotObservable(
            OnNext(100, "a"),
            OnNext(110, "b"),
            OnNext(200, "a"),
            OnNext(220, "c"),
            OnNext(221, "a")
        );

        // 用Scan维护状态:已见元素集合 + 当前元素是否是重复项
        var statefulStream = source.Scan(
            (new HashSet<string>(), (string)null, false),
            (state, item) =>
            {
                var isDuplicate = !state.Item1.Add(item);
                return (state.Item1, item, isDuplicate);
            }
        );

        // 拆分去重流:筛选非重复项
        var distinctResults = scheduler.CreateObserver<string>();
        statefulStream
            .Where(tuple => !tuple.Item3)
            .Select(tuple => tuple.Item2)
            .Subscribe(distinctResults);

        // 拆分重复值流:筛选重复项
        var duplicateResults = scheduler.CreateObserver<string>();
        statefulStream
            .Where(tuple => tuple.Item3)
            .Select(tuple => tuple.Item2)
            .Subscribe(duplicateResults);

        scheduler.AdvanceBy(1000);

        // 验证结果
        distinctResults.Messages.AssertEqual(
            OnNext(100, "a"),
            OnNext(110, "b"),
            OnNext(220, "c")
        );

        duplicateResults.Messages.AssertEqual(
            OnNext(200, "a"),
            OnNext(221, "a")
        );
    }
}

为什么之前的CombineLatest效果不佳?

CombineLatest的逻辑是当任意一个流产生新值时,结合所有流的最新值,这和我们需要的“跟踪哪些元素是重复的”逻辑不匹配。我们需要的是对每个元素判断是否已出现过,而不是结合两个流的最新值,所以用上述的共享源+筛选/状态跟踪的方式更合适。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:05:06