如何基于首个元素的判断条件转换Rx.NET中的Observable
Rx.NET 操作符实现方案
实现逻辑说明
你需要的操作符可以通过两种方式实现,以下是可直接使用的代码:
方案1:基于现有Rx操作符拼接实现
无需自己处理订阅逻辑,复用内置操作符保证线程安全和Rx规范兼容性:
using System.Reactive.Linq; public static class ObservableExtensions { public static IObservable<string> ForwardIfFirstIsA(this IObservable<string> source) { // 使用Publish避免多次订阅源Observable return source.Publish(published => // 取第一个元素判断是否等于"a" published.Take(1) .Where(first => first == "a") // 符合条件则拼接后续所有元素 .Concat(published.Skip(1)) ); } }
方案2:通用自定义操作符
如果后续需要修改首元素判断规则、适配其他元素类型,可以用通用版本:
using System.Reactive; using System.Reactive.Linq; public static class ObservableExtensions { // 通用扩展:首元素符合指定条件则转发所有元素,否则仅发完成信号 public static IObservable<T> ForwardIfFirstMatch<T>(this IObservable<T> source, Func<T, bool> firstElementPredicate) { return Observable.Create<T>(observer => { bool hasReceivedFirst = false; bool allowForward = false; return source.Subscribe( onNext: value => { if (!hasReceivedFirst) { hasReceivedFirst = true; allowForward = firstElementPredicate(value); if (allowForward) { observer.OnNext(value); } return; } if (allowForward) { observer.OnNext(value); } }, onError: observer.OnError, onCompleted: observer.OnCompleted ); }); } // 针对你的需求封装的专用方法 public static IObservable<string> ForwardIfFirstIsA(this IObservable<string> source) { return source.ForwardIfFirstMatch(str => str == "a"); } }
效果验证
和你给出的示例完全匹配:
- 输入:
-a-b-c-d-|-→ 输出:-a-b-c-d-|-,所有元素原样转发 - 输入:
-b-c-d-|-→ 输出:-|-,无任何元素转发,直接发送完成信号
内容的提问来源于stack exchange,提问作者jack
相关产品推荐
相关产品推荐

