【问题标题】:Does Observable.Range break The Observable Contract?Observable.Range 是否会破坏 Observable 合约?
【发布时间】:2016-12-30 23:36:59
【问题描述】:

在了解 Rx 的过程中,我遇到了一个经常重复的关于 Observables 的规则,在 The Observable Contract 中有详细说明。

在发出 OnCompleted 或 OnError 通知后,它可能不会再发出任何进一步的通知。

这对我来说很有意义,因为让 Observable 在完成后继续产生值会令人困惑,但是当我在 .NET 中测试 Observable.Range 方法时,我注意到它没有表现出这种行为,事实上很多 Observables 违反了这条规则。

var rangeObservable = Observable.Range(0, 5);

rangeObservable.Subscribe(Console.WriteLine, () => Console.WriteLine("Done first!"));
Console.ReadLine();

rangeObservable.Subscribe(Console.WriteLine, () => Console.WriteLine("Done second!"));
Console.ReadLine();

//Output:
//0
//1
//2
//3
//4
//Done first!

//0
//1
//2
//3
//4
//Done second!

显然rangeObservable 调用了两次OnComplete 并在第一次OnComplete 之后产生了值。这让我相信这不是关于 Observables 的规则,而是关于 Subscriptions 的规则。也就是说,一个 Observable 可以产生任意多的终止消息,甚至在它产生之后产生值,只要每个 Subscription 只接收一个终止消息并且没有收到之后的进一步消息。

当它说 Observable 时,它们实际上是指 Subscription 吗?它们真的是不同的东西吗?我对模型有根本的误解吗?

【问题讨论】:

    标签: c# system.reactive reactive-programming observable reactive


    【解决方案1】:

    observable 合约必须对任何被观察的 Observable 有效。 在 Observable 未被观察时是否发生任何事情都留给 observable 的实现。

    在 Enumerable 中考虑类比会有所帮助 - Observable 是 Enumerable 的对偶。在可枚举中,你会有 range = Enumerable.Range(0, 5),您将使用与上述类似的范围:

    range.ForEach(Console.WriteLine); //prints 0 - 4
    
    range.ForEach(Console.WriteLine); //prints 0 - 4 again
    

    并发现这是完全可以接受的行为,因为只有在调用 GetEnumerator 时才会创建实际的数字生成器。同样,在 Observable 中,等价的方法是Subscribe

    范围的实现是这样的:

            static IObservable<int> Range(int start, int count)
            {
                return Observable.Create<int>(observer =>
                {
                    for (int i = 0; i < count; i++)
                        observer.OnNext(start + i);
    
                    observer.OnCompleted();
    
                    return Disposable.Empty;
                });
            }
    

    这里,每次有订阅时都会调用observer =&gt; {...} 函数。工作在 subscribe 方法中完成。您可以很容易地看到它 (1) 为每个观察者推送相同的序列,(2) 每个观察者只完成一次。

    这些只有在你观察它们时才会发生某些事情的可观察对象称为冷可观察对象。 Here's an article 描述这个概念。

    注意

    Range 是一个非常幼稚的实现,仅用于说明目的。该方法在完成之前不会返回一次性用品 - 所以Disposable.Empty 是可以接受的。正确的实现将在调度程序上运行工作,并使用已检查的一次性来查看订阅是否已被释放,然后再继续循环。

    要点是手动实现可观察合约是困难,这就是存在 Rx 库的原因 - 通过组合构建功能。

    【讨论】:

    • 完美。让我理解这一点的思想转变是意识到 Observable 不应被视为“对象”,而应被视为“订阅方法”。只是每次想要订阅时调用的方法。那以及我真的不知道如果没有这个模型将如何实现 Range 运算符这一事实。你打算做什么,只是让每个范围运算符在构造后立即开始?没有人能够从中获得任何价值!
    • @gitbox 哈哈。物有所值。
    • @Asti - 除了你对 Range 的实现之外,我喜欢你对所有问题的回答 - 如果你发现自己在 .Create 方法中写 return Disposable.Empty;,你正在 这样做出了点问题。
    • @Enigmativity 你完全正确。这是一个非常幼稚的实现,只是为了说明。该方法在完成之前不会返回一次性用品 - 所以我推断它返回什么并不重要。正确的实现将在调度程序上运行工作,并可能使用已检查的一次性,但这可能会使解释更加混乱。
    • @Asti - 最好将这个解释添加到答案中。我讨厌有人在生产代码中使用该实现。
    【解决方案2】:

    Observable.Range 返回一个 cold 可观察对象,这意味着它为每个订阅者“重播”它的行为。由于“OnNext* OnComplete|OnError”合约仅适用于订阅,这完全没问题。

    有关热/冷 observables 的更多信息,请参阅my answer on "IConnectableObservables in Rx"

    【讨论】:

    • 我越想越觉得这个答案很明显!当我们谈论一些 Observable 生成值时,它指的是单个订阅的值。我没有想到的情况是当一个 Observable 有多个观察者时,它会因此发送大量 OnCompleted 消息,尽管是在“同一时间”。由于这是一种常见的情况,Observable 合约永远不会禁止它。哇!
    • 我熟悉冷和热可观察对象的概念。我的问题可以表述为“所有冷的 observables 都违反了 Observable Contract 吗?”因为他们都会以我之前的理解。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-01-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-09-19
    • 2012-05-26
    • 1970-01-01
    相关资源
    最近更新 更多