【问题标题】:FirstOrLast IObservable ExtensionFirstOrLast IObservable 扩展
【发布时间】:2012-09-08 00:41:03
【问题描述】:

我想通过IObservable<T> 查找与谓词匹配的元素,如果找不到,则返回IObservable<T> 的最后一个元素。我不想存储IObservable<T>的全部内容,也不想循环IObservable两次,所以我设置了一个扩展方法

public static class ObservableExtensions
{
    public static IObservable<T> FirstOrLastAsync<T>(this IObservable<T> source, Func<T, bool> pred)
    {
        return Observable.Create<T>(o =>
        {
            var hot = source.Publish();
            var store = new AsyncSubject<T>();
            var d1 = hot.Subscribe(store);
            var d2 = hot.FirstAsync(x => pred(x)).Amb(store).Subscribe(o);
            var d3 = hot.Connect();
            return new CompositeDisposable(d1, d2, d3);
        });
    }

    public static T FirstOrLast<T>(this IObservable<T> source, Func<T, bool> pred)
    {
        return source.FirstOrLastAsync(pred).Wait();
    }
}

Async 方法从传入的可能冷的 observable 创建一个热的 observable。它订阅一个 AsyncSubject&lt;T&gt; 来记住最后一个元素,并订阅一个 IObservable&lt;T&gt; 来查找该元素。然后它从IObservable&lt;T&gt;s 中的任何一个中获取第一个元素,该元素首先通过.Amb 返回一个值(AsyncSubject&lt;T&gt; 在收到.OnCompleted 消息之前不会返回一个值)。

我的问题如下:

  • 可以用不同的 Observable 方法写得更好或更简洁吗?
  • 是否所有这些一次性用品都需要包含在 CompositeDisposable 中?
  • 当 hot observable 在没有找到匹配元素的情况下完成时,FirstAsync 抛出异常和 AsyncSubject 传播其值之间是否存在竞争条件?
  • 如果是这样,我是否需要将行更改为:

var d2 = hot.Where(x =&gt; pred(x)).Take(1).Amb(store).Subscribe(o);

我对 RX 很陌生,这是我在 IObservable 上的第一个扩展。

编辑

我最终选择了

public static class ObservableExtensions
{
    public static IObservable<T> FirstOrLastAsync<T>(this IObservable<T> source, Func<T, bool> pred)
    {
        var hot = source.Publish().RefCount();
        return hot.TakeLast(1).Amb(hot.Where(pred).Take(1).Concat(Observable.Never<T>()));
    }

    public static T FirstOrLast<T>(this IObservable<T> source, Func<T, bool> pred)
    {
        return source.FirstOrLastAsync(pred).First();
    }
}

【问题讨论】:

    标签: c# system.reactive


    【解决方案1】:

    您可以将您想要的两个案例放在一起。 如果你的源 observable 是冷的,你可以做一个Publish|Refcount

        public static IObservable<T> FirstOrLast<T>(this IObservable<T> source, Func<T, bool> predicate)
        {
            return source.TakeLast(1).Amb(source.Where(predicate).Take(1));
        }
    

    测试:

            var source = Observable.Interval(TimeSpan.FromSeconds(0.1))
                                   .Take(10)
                                   .Publish()
                                   .RefCount();
    
            FirstOrLast(source, i => i == 5).Subscribe(Console.WriteLine); //5
            FirstOrLast(source, i => i == 11).Subscribe(Console.WriteLine); //9
    

    【讨论】:

    • 当没有任何值与谓词匹配时,这似乎不会返回最后一个值。
    • @Enigmativity 怎么样?即使没有一个值匹配,TakeLast 流仍然会产生一个值。
    • 当我尝试这个时,我得到一个序列不包含任何元素异常,但如果它有效的话,这绝对有资格更简单。
    • 这可行,但我不确定为什么需要添加位。 return Observable.Amb(source.Where(predicate).Take(1).Concat(Observable.Never&lt;T&gt;()), source.TakeLast(1));
    • @MartinNeal 这可能与 Amb 对通知的反应而不是对值的反应有关。只有当序列完成时,TakeLast 才会被转发,所以如果完成通知通过另一个流到达 Amb,它可能会断开连接。我已经编辑了答案以反映这一点。
    【解决方案2】:

    我试图生成一个“更简单”的有效查询,但到目前为止没有。

    如果我坚持你的基本结构,我可以提供一点改进。试试这个:

    public static IObservable<T> FirstOrLastAsync<T>(
        this IObservable<T> source, Func<T, bool> pred)
    {
        return Observable.Create<T>(o =>
        {
            var hot = source.Publish();
            var store = new AsyncSubject<T>();
            var d1 = hot.Subscribe(store);
            var d2 =
                hot
                    .Where(x => pred(x))
                    .Concat(store)
                    .Take(1)
                    .Subscribe(o);
            var d3 = hot.Connect();
            return new CompositeDisposable(d1, d2, d3);
        });
    }
    

    它并没有好太多,但我比使用Amb 更喜欢它。我认为它只是更清洁一点。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多