如何在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
相关产品推荐
相关产品推荐

