【问题标题】:Observable.Retry() is not working as expected with Hot ObservablesObservable.Retry() 在 Hot Observables 中没有按预期工作
【发布时间】:2018-05-10 04:41:05
【问题描述】:

我有以下可观察序列

int num = 0;            

var o = Observable.Create<int>(observer => Task.Run(() =>
{
    var rnd = new Random((int)DateTime.Now.Ticks);
    Console.WriteLine($"Starting subscription loop # {++num}");
    for (int i=0;i<100;i++)
    {
        Thread.Sleep(200);

        if (i == 3)
        {
            observer.OnError(new ApplicationException("test exception"));
            break;
        }

        observer.OnNext(rnd.Next(0, 50));
    }
})).Publish().RefCount();

以及以下通知处理程序

o
    .Retry()
    .Subscribe(Console.WriteLine, ex => Console.WriteLine($"Exception occurred: {ex.Message}"), () => Console.WriteLine("Completed"));

这是我的输出

Starting subscription loop # 1
47
27
12
Starting subscription loop # 2
Starting subscription loop # 3
Starting subscription loop # 4
Starting subscription loop # 5
Starting subscription loop # 6
Starting subscription loop # 7
Starting subscription loop # 8
Starting subscription loop # 9
...

我在Lee Campbell's IntroToRx book 中阅读了以下内容

如果您预计您的序列会遇到可预测的问题, 您可能只是想重试。一个这样的例子,当你想 重试是在执行 I/O 时(例如 Web 请求或磁盘访问)。输入输出 因间歇性故障而臭名昭著。重试扩展方法 提供重试失败指定次数的能力,或 直到成功。

我在样本中注意到的行为与坎贝尔指出的行为不符,也不符合他的样本。我错过了什么?

如果我不Publish().RefCount(),它可以正常工作。

【问题讨论】:

  • new Random((int)DateTime.Now.Ticks) 是一种浪费,因为 new Random() 在幕后就是这样做的。你最好创建一个[ThreadStatic] 静态字段变量。
  • @Enigmativity No it does not 它使用的逻辑涉及两个 Random 对象一个循环一个维护。
  • 有趣。我想知道什么时候改变了。尽管如此,使用new Random((int)DateTime.Now.Ticks) 而不是new Random() 并没有什么不同。它仍然存在同样的坏种子问题。
  • 相关:How to fix the inconsistency of the Publish().RefCount() behavior? TL;DR,您观察到的问题行为无法可靠修复。您Publish 的原因是为了确保多个订阅者会收到来自同一序列的通知。这只能通过有状态的Subject(记住源的完成状态)备份的Publish 来保证。可悲的后果是Publish 不可重复使用。只能连接一次。

标签: c# system.reactive


【解决方案1】:

当一个 observable 出错时,它就死了并且结束了。不再有通知流出。在您的情况下,o 出错了,并且由于.Publish().Refcount()Retry 正在尝试重新订阅相同的 observable(已死且已完成)。这就是 Publish 所做的:它不是创建新的 observable,而是为多个客户端订阅相同的 observable。

如果您删除 .Publish().Refcount(),您会看到它尝试重新订阅一个新的 observable。

【讨论】:

  • Retry 应该取消订阅死的 observable 应该让RefCount 层知道没有更多的订阅者,因此在Publish() 返回的IConnectableObservable 上调用Disconnect()Retry 应该尝试另一个订阅。 Retry 是否在重新订阅之前不处理订阅?
  • 这基本上是一个对你不利的实现细节。无论是订阅然后取消订阅,还是反之亦然,都不是您想要依赖的。
  • @fahadash - OnError 正在取消订阅 - .Retry 正在重新订阅。但正如 Shlomo 指出的那样,.Publish().RefCount() 组合是一次只能观察到的,因此.Retry 在这种情况下永远不会成功。
猜你喜欢
  • 2021-10-22
  • 2014-01-26
  • 2021-10-19
  • 2020-03-18
  • 2012-06-14
  • 2014-11-15
  • 1970-01-01
  • 2012-07-02
  • 2011-09-07
相关资源
最近更新 更多