【问题标题】:Is Reactive Extensions evaluating too many times?反应式扩展是否评估了太多次?
【发布时间】:2012-09-07 13:35:31
【问题描述】:

Reactive Extensions 应该评估其各种运算符多少次?

我有以下测试代码:

var seconds = Observable
    .Interval(TimeSpan.FromSeconds(5))
    .Do(_ => Console.WriteLine("{0} Generated Data", DateTime.Now.ToLongTimeString()));

var split = seconds
    .Do(_ => Console.WriteLine("{0}  Split/Branch received Data", DateTime.Now.ToLongTimeString()));

var merged = seconds
    .Merge(split)
    .Do(_ => Console.WriteLine("{0}   Received Merged data", DateTime.Now.ToLongTimeString()));

var pipeline = merged.Subscribe();

我希望它每五秒写入一次“生成的数据”。然后,它将数据传递给写入“Split/Branch received Data”的“split”流和写入“Received Merged data”的“merged”流。最后,因为“合并”流也从“拆分”流接收,所以它第二次接收数据并第二次写入“接收的合并数据”。 (它写其中一些的顺序并不特别相关)

但我得到的输出是这样的:

8:29:56 AM Generated Data
8:29:56 AM Generated Data
8:29:56 AM  Split/Branch received Data
8:29:56 AM   Received Merged data
8:29:56 AM   Received Merged data
8:30:01 AM Generated Data
8:30:01 AM Generated Data
8:30:01 AM  Split/Branch received Data
8:30:01 AM   Received Merged data
8:30:01 AM   Received Merged data

它正在写入两次“生成的数据”。据我了解,订阅“秒”IObservable 的下游观察者的数量不应该影响“生成的数据”写入的次数(应该是 ONCE),但确实如此。为什么?

注意 我在 .Net Framework 3.5 环境中使用反应式扩展的稳定版本 v1.0 SP1。

【问题讨论】:

    标签: c# system.reactive


    【解决方案1】:

    据推测,他们选择这种方法是为了让每个订阅者从他们的初始订阅开始以相同的时间间隔获取其值。考虑一下您的备用间隔将如何工作:

    0s - First subscriber subscribes
    5s - Value: 0
    8s - Second subscriber subscribes
    10s - Value: 1
    15s - Value: 2
    17s - Unsubscribe both
    

    你最终得到的是这样的:

    First  -----0----1----2-|
    Second         --1----2-|
    

    在这种情况下,两个观察者有明显不同的结果,具体取决于是否已经附加了任何其他观察者。在其实施过程中,Interval 为每个订阅者提供相同的体验,无论订单或过去的订阅者如何。

    综上所述,您可以通过在创建 seconds 可观察对象时添加 .Publish().RefCount() 将 Interval“转换”为您描述的行为。

    【讨论】:

    • .Publish().RefCount() 成功了,谢谢。我理解您的解释及其必要性,但我发现这种默认行为对于我想使用响应式扩展的大多数场景并不是特别有用。
    • 这种行为是精心设计的结果:-)。如果仅仅创建一个源就意味着它的激活(所谓的“热”模型),那么在创建时间和订阅时间之间可能会丢失消息。注意 Task 这不是问题,因为最多必须为“迟到的订阅者”缓存一个值(即 ContinueWith 回调)。对于事件流,不能缓存值;这样做是要求 OutOfMemoryException。 Asti 在下面的回答与 IEnumerable 进行了有用的类比。另一个类比是 .NET 事件。两个 += 调用可以导致两个“会话”。
    【解决方案2】:

    虽然有时如果在每一步都多播序列似乎会很好,但如果是这样,它就不允许你拥有 Rx 允许的丰富组合。

    换个角度想, IObservable 是 IEnumerable 的基于推送的对偶。 IEnumerable 具有惰性求值的属性 - 在您开始通过 Enumerator 之前不会计算值。 Rx 序列是惰性组合的,最后一个 Subscribe()(相当于 For-Each 的 Observable)实现了这个序列。

    通过这种方式,您只需从最后一个阶段取消订阅,即可在所有阶段停止管道,让您无需经历管理个人订阅的噩梦。

    【讨论】:

    • @Enigmativity 谢谢。这意味着你有很多东西。
    • IEnumerable 不使用惰性求值,LINQ 使用惰性求值。 IEnumerable 只是体现了遍历集合的概念,然后重新开始回到该集合的开头。 LINQ 以此为基础,因此可以通过 IEnumerable 随时重新开始其收集进程。相反,IObservable 只是提供了一种订阅数据的方式,并没有提供可以重新启动该数据的“合同”。能够从 IEnumerable 构建您可以让数据自行重新启动,这是一个巧合,这是一个未指定的巧合。
    • LINQ =/= IEnumerable。 LINQ 是一组语言扩展,编译器将其转换为 monad 的某些实现。 IEnumerable 根据定义是惰性的 - 在从序列中提取值之前,您不会计算值。
    • 说 IEnumerable 是惰性的就像说一个数组是惰性的,因为您只在需要时访问该数组。两者都只是提供访问机制;数组通过索引提供对数组的访问,枚举器提供只向前遍历(和重置)。 LINQ 选择了惰性,实际上某些运算符(如 OrderBy)实际上放弃了它们的惰性求值(一旦它们自己进行求值),因为它们必须立即求值整个枚举以计算其结果。
    【解决方案3】:

    在相关的注释中,这里有一个脑筋急转弯,说明 Asti 与惰性求值的可枚举序列的类比:

    private static Random s_rand = new Random();
    
    public static IEnumerable<int> Rand()
    {
        while (true)
            yield return s_rand.Next();
    }
    
    public static void Main()
    {
        var xs = Rand();
    
        var res = xs.Zip(xs, (l, r) => l == r).All(b => b);
    
        Console.WriteLine(res);
    }
    

    如果您使用自身压缩一个随机序列,您是否希望所有元素对都相同(即导致上面的代码永远运行)?或者,您是否希望代码因某种原因终止并打印 false?

    (创建类似的可观察代码留给读者作为练习。)

    【讨论】:

    • [剧透] All 运算符将在其谓词评估为 false 时中断。由于 xs 是惰性的,Zip 保证将两个不同的值配对,导致选择器总是返回 false - 因此导致枚举在第一个值上停止。为了在 Observables 中获得等效的多播,我们必须使用 xs = Memoize(Rand()) - 这将导致程序永远运行。 [/剧透]
    • Rx 等效项:var xs = Observable.Repeat(0).Select(_ =&gt; s_rand.Next()); ... res.Subscribe(Console.WriteLine);
    • 我觉得这体现了这样一个事实,即没有任何 Reactive Extensions 功能可以区分实时的、不可重新启动的数据和缓存的、可重新启动的数据 (IEnumerable)。您不知道何时使用运算符(即Do()),如果数据源是一个或另一个。然而默认的Subscribe() 实现似乎选择了后者,因此当你第一次遇到它时它会让你措手不及。我很感激开发团队意识到需要实时数据并至少添加了.Publish().RefCount()
    • @MattG 我添加了一个新答案,因为评论太大了。
    【解决方案4】:

    从面向对象的角度来看,通常认为流基于定义Observables/Enumerables 的接口。如果您可以忽略在枚举器上定义了一种称为重置的便捷方法这一事实 - 枚举在功能上是f -&gt; g -&gt; value?。可枚举本质上是您调用以获取枚举器的函数,该函数本质上是您一直调用的函数,直到没有更多值要返回为止。

    类似地,Observable 被简单地定义为f(g) -&gt; g(h) -&gt; h(value?) - 它是您提供的函数,当有值时您想要调用的函数。

    这就是为什么将可枚举或可观察的描述为除了以某种方式定义的一组函数之外的任何东西是没有意义的,因此它们可以组合 - 合同是为了确保组合计算的能力。

    无论它们是实时的、缓存的还是惰性的,都是可能在其他地方抽象出来的实现细节——虽然我当然不同意这些细节很重要,但更重要的是关注它的功能性质。

    作为数据库查询或目录列表的序列具有与预先计算的一组值(如数组)相同的IEnumerable 接口。这取决于最终使用序列来进行区分的代码。如果您能习惯它是一种组合高阶函数的方法,您会发现使用 Rx 或 Ix 建模问题会更容易。

    【讨论】:

    • 我同意。 Observables 和 Enumerables 只是调用函数(可能返回数据)的遍历机制。实现细节,即 LINQ 和响应式扩展,可以选择惰性、实时或缓存。 LINQ 选择允许惰性求值。 Reactive Extensions 选择假定 Observables 是缓存数据,并要求用户采取额外步骤 (.Publish().RefCount()) 将它们用作实时源。问题是 Observable 合约没有将其定义为缓存数据或实时数据,但 Rx 假定为前者。
    猜你喜欢
    • 2023-03-20
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2010-11-24
    • 2021-12-15
    相关资源
    最近更新 更多