【问题标题】:RX Observable.TakeWhile checks condition BEFORE each element but I need to perform the check afterRX Observable.TakeWhile 在每个元素之前检查条件,但我需要在之后执行检查
【发布时间】:2013-02-04 23:29:24
【问题描述】:

Observable.TakeWhile 允许您在条件为真时运行序列(使用委托,以便我们可以对实际序列对象执行计算),但它会在每个元素之前检查此条件。如何在每个元素之后执行相同的检查?

以下代码演示了问题

    void RunIt()
    {
        List<SomeCommand> listOfCommands = new List<SomeCommand>();
        listOfCommands.Add(new SomeCommand { CurrentIndex = 1, TotalCount = 3 });
        listOfCommands.Add(new SomeCommand { CurrentIndex = 2, TotalCount = 3 });
        listOfCommands.Add(new SomeCommand { CurrentIndex = 3, TotalCount = 3 });

        var obs = listOfCommands.ToObservable().TakeWhile(c => c.CurrentIndex != c.TotalCount);

        obs.Subscribe(x =>
        {
            Debug.WriteLine("{0} of {1}", x.CurrentIndex, x.TotalCount);
        });
    }

    class SomeCommand
    {
        public int CurrentIndex;
        public int TotalCount;
    }

这个输出

1 of 3
2 of 3

我无法获取第三个元素

看这个例子,你可能会认为我所要做的就是像这样改变我的条件 -

var obs = listOfCommands.ToObservable().TakeWhile(c => c.CurrentIndex <= c.TotalCount);

但是 observable 永远不会完成(因为在我的真实世界代码中,流不会在这三个命令之后结束)

【问题讨论】:

    标签: c# system.reactive


    【解决方案1】:

    没有内置运算符可以满足您的要求,但这里有一个使用 Publish 运行两个查询,同时只订阅底层 observable 一次:

    // Emits matching values, but includes the value that failed the filter
    public static IObservable<T> TakeWhileInclusive<T>(
        this IObservable<T> source, Func<T, bool> predicate)
    {
        return source.Publish(co => co.TakeWhile(predicate)
            .Merge(co.SkipWhile(predicate).Take(1)));
    }
    

    然后:

    var obs = listOfCommands.ToObservable()
        .TakeWhileInclusive(c.CurrentIndex != c.TotalCount);
    

    【讨论】:

    • 是的,你的还是更好:单一合并订阅
    • 我认为当它被称为 TakeUntil 而不是 TakeWhileInclusive 时,这种用法可以更好地阅读。例如 Status.TakeUntil(s => s == Status.Completed) 比 Status.TakeWhileInclusive(s => s != Status.Completed) 读得更好。
    【解决方案2】:

    最终编辑:

    我的解决方案基于 Sergey 在此线程中的 TakeWhileInclusive 实现 - How to complete a Rx Observable depending on a condition in a event

    public static IObservable<TSource> TakeUntil<TSource>(
            this IObservable<TSource> source, Func<TSource, bool> predicate)
    {
        return Observable
            .Create<TSource>(o => source.Subscribe(x =>
            {
                o.OnNext(x);
                if (predicate(x))
                    o.OnCompleted();
            },
            o.OnError,
            o.OnCompleted
        ));
    }
    

    【讨论】:

    • 仅供参考,您永远不会取消订阅source。如果它不能自然地自行完成,您可能会遇到内存泄漏。
    • 关于如何做到这一点的任何建议?
    • 我会创建一个SerialDisposable 并将Subscribe 返回分配给它的Disposable 属性。然后在 Complete/Error 上处理 SerialDisposable 实例。出于兴趣,为什么基于发布的实现不适合?
    • 如果我没记错的话,这也是我不能使用 Alex G 的答案的原因。当进一步减少它时,它会丢弃最后一个元素——这是我们想要得到的元素。 (请参阅我对上述 Alex G 回答的评论)
    • 不知道为什么会这样。我刚刚重现了您的问题,并使用我的TakeWhileInclusive 作为替代品,按预期返回了所有三个结果。
    【解决方案3】:

    您可以使用TakeUntil 运算符来获取每个项目,直到第二个来源产生一个值;在这种情况下,我们可以将第二个流指定为谓词通过后的第一个值:

    public static IObservable<TSource> TakeWhileInclusive<TSource>(
        this IObservable<TSource> source,
        Func<TSource, bool> predicate)
    {
        return source.TakeUntil(source.SkipWhile(x => predicate(x)).Skip(1));
    }
    

    【讨论】:

    • 订阅TakeWhileInclusive() 再进行reducing时不起作用。即.. TakeWhileInclusive().Timeout().. 试试看,它会丢弃最后一个元素
    • 抱歉,由于我的测试出现错误,我进行了修改,更改了代码 - 确保存在 .Skip(1) 操作!我现在已经用.Timeout() 操作进行了测试,它似乎执行得很好,你确定不是你的进一步操作正在切断序列的结尾吗?
    • 经过更多测试,很明显使用此运算符,OnComplete 直到处理完序列中的下一个元素后才会调用,因此与某些运算符结合使用时将无法正常工作。我会尝试想出一个巧妙的解决方案,但我不确定是否有。
    • 如果它有效,那么它就有效:) 虽然我仍然希望 C# 扩展 yield 关键字以适用于 Observables - 会让这变得更容易!
    • 此示例假定源是热的,即它不发布源以防止来自多个订阅的多个副作用。请参阅@Richard Szalay 的安全实施答案。
    【解决方案4】:

    我认为你在关注TakeWhile,而不是TakeUntil

    var list = (new List<int>(){1,2,3,4,5,6,7,8,9,10});
    var takeWhile = list
            .ToObservable()
            .Select((_, i) => Tuple.Create(i, _))
            .TakeWhile(tup => tup.Item1 < list.Count)
            .Do(_ => Console.WriteLine("Outputting {0}", _.Item2));
    

    好的,您想要的东西不存在开箱即用,至少我不知道具有这种特定语法的东西。也就是说,你可以很容易地把它拼凑起来(而且它不是讨厌):

    var fakeCmds = Enumerable
        .Range(1, 100)
        .Select(i => new SomeCommand() {CurrentIndex = i, TotalCount = 10})
        .ToObservable();
    
    var beforeMatch = fakeCmds
        .TakeWhile(c => c.CurrentIndex != c.TotalCount);
    var theMatch = fakeCmds
        .SkipWhile(c => c.CurrentIndex != c.TotalCount)
        .TakeWhile(c => c.CurrentIndex == c.TotalCount);
    var upToAndIncluding = Observable.Concat(beforeMatch, theMatch);
    

    【讨论】:

    • 有人已经使用 TakeWhile 给出了答案(然后将其删除)。 TakeWhile 不采用最终元素,因为它似乎是在每个元素之前而不是之后检查条件。我需要 RX 等效的 Do While 循环,而不是简单的 While 循环。
    • 嗯,有DoWhile...这里,我举个例子。
    • DoWhile 不像 TakeWhile 那样接受委托。我知道我可以使用不太优雅的代码来完成这项工作,但重点是尽可能以最整洁的方式做事。
    • 实际上,等等:我这里的例子打印了Outputting 1...10,所以它显然在处理最后一个元素。你到底看到了什么?
    • 我不喜欢你的例子,它太复杂了。我添加了一个显示问题的答案。
    【解决方案5】:

    组合,使用新的SkipUntilTakeUntil

    SkipUntil return source.Publish(s => s.SkipUntil(s.Where(predicate)));

    TakeUntil (含) return source.Publish(s => s.TakeUntil(s.SkipUntil(predicate)));

    完整来源https://gist.github.com/GeorgeTsiokos/a4985b812c4048c428a981468a965a86

    【讨论】:

      【解决方案6】:

      也许以下方式对某人有用。 您必须使用“Do”方法和空的“Subscribe”方法。

          listOfCommands.ToObservable()
          .Do(x =>
          {
              Debug.WriteLine("{0} of {1}", x.CurrentIndex, x.TotalCount);
          })
          .TakeWhile(c => c.CurrentIndex != c.TotalCount)
          .Subscribe();
      

      这样您无需编写自己的扩展即可获得结果。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2012-10-18
        • 2013-04-25
        • 1970-01-01
        • 2022-10-13
        • 2019-11-24
        • 2013-03-28
        • 2021-04-21
        • 2021-03-09
        相关资源
        最近更新 更多