【问题标题】:Write an Rx "RetryAfter" extension method编写一个 Rx "RetryAfter" 扩展方法
【发布时间】:2013-09-24 10:03:19
【问题描述】:

IntroToRx一书中,作者建议为I/O编写一个“智能”重试,在一段时间后重试I/O请求,如网络请求。

这是确切的段落:

添加到您自己的库中的一个有用的扩展方法可能是“返回 Off and Retry”的方法。我合作过的团队发现了这样一个 在执行 I/O(尤其是网络请求)时很有用的功能。这 概念是尝试,失败时等待给定的时间段,然后 然后再试一次。您的此方法版本可能会考虑到 您要重试的异常类型,以及最大数量 重试的次数。您甚至可能希望将等待时间延长至 以后每次重试时不要那么激进。

不幸的是,我不知道如何编写这个方法。 :(

【问题讨论】:

    标签: c# system.reactive


    【解决方案1】:

    实现回退重试的关键是deferred observables。在有人订阅它之前,延迟的 observable 不会执行它的工厂。它会为每个订阅调用工厂,使其成为我们重试场景的理想选择。

    假设我们有一个触发网络请求的方法。

    public IObservable<WebResponse> SomeApiMethod() { ... }
    

    为了这个小sn-p的目的,让我们将deferred定义为source

    var source = Observable.Defer(() => SomeApiMethod());
    

    每当有人订阅源时,它都会调用 SomeApiMethod 并启动一个新的 Web 请求。失败时重试的简单方法是使用内置的 Retry 运算符。

    source.Retry(4)
    

    不过,这对 API 来说不是很好,也不是您想要的。我们需要在每次尝试之间延迟发起请求。一种方法是使用delayed subscription

    Observable.Defer(() => source.DelaySubscription(TimeSpan.FromSeconds(1))).Retry(4)
    

    这并不理想,因为即使在第一次请求时也会增加延迟,让我们解决这个问题。

    int attempt = 0;
    Observable.Defer(() => { 
       return ((++attempt == 1)  ? source : source.DelaySubscription(TimeSpan.FromSeconds(1)))
    })
    .Retry(4)
    .Select(response => ...)
    

    只是暂停一秒钟并不是一个很好的重试方法,所以让我们将该常量更改为一个接收重试计数并返回适当延迟的函数。指数回退很容易实现。

    Func<int, TimeSpan> strategy = n => TimeSpan.FromSeconds(Math.Pow(n, 2));
    
    ((++attempt == 1)  ? source : source.DelaySubscription(strategy(attempt - 1)))
    

    我们现在差不多完成了,我们只需要添加一种方法来指定我们应该重试哪些异常。让我们添加一个函数,给定一个异常返回,无论重试是否有意义,我们将其称为 retryOnError。

    现在我们需要编写一些看起来很吓人的代码,但请耐心等待。

    Observable.Defer(() => {
        return ((++attempt == 1)  ? source : source.DelaySubscription(strategy(attempt - 1)))
            .Select(item => new Tuple<bool, WebResponse, Exception>(true, item, null))
            .Catch<Tuple<bool, WebResponse, Exception>, Exception>(e => retryOnError(e)
                ? Observable.Throw<Tuple<bool, WebResponse, Exception>>(e)
                : Observable.Return(new Tuple<bool, WebResponse, Exception>(false, null, e)));
    })
    .Retry(retryCount)
    .SelectMany(t => t.Item1
        ? Observable.Return(t.Item2)
        : Observable.Throw<T>(t.Item3))
    

    所有这些尖括号都用于编组异常,我们不应该重试超过.Retry()。我们将内部 observable 设为 IObservable&lt;Tuple&lt;bool, WebResponse, Exception&gt;&gt;,其中第一个 bool 表示我们是否有响应或异常。如果 retryOnError 指示我们应该为特定异常重试,则内部 observable 将抛出并且将被重试拾取。 SelectMany 只是解开我们的 Tuple 并使生成的 observable 再次成为 IObservable&lt;WebRequest&gt;

    查看我的gist with full source and tests 了解最终版本。有了这个运算符,我们可以非常简洁地编写重试代码

    Observable.Defer(() => SomApiMethod())
      .RetryWithBackoffStrategy(
         retryCount: 4, 
         retryOnError: e => e is ApiRetryWebException
      )
    

    【讨论】:

    • 在我看来,对于这个实现,源 observable 永远不会取消订阅。在这里粘贴有点困难,但试试这个,你会看到间隔一直在滴答作响: Observable.Interval(TimeSpan.FromSeconds(1)).Do(Console.WriteLine).RetryWithBackoffStrategy().Take(1).Subscribe ();
    • @NiallConnaughton 不错的收获!它没有从源中取消订阅的原因与我最初在我们拥有的另一个内部操作符之后对该方法建模有关,该操作符产生热可观察量。该操作员不应该这样做。我更改了代码以生成冷可观察对象,并添加了一个测试来验证它是否取消订阅。谢谢!
    • Marcus 我曾尝试使用您的代码,但有一个小问题,我试图在 Github 上提问,但在我工作的地方被阻止了。所以我不得不打开一个新的 SO 问题stackoverflow.com/questions/20189166/rx-back-off-and-retry,如果你有时间可以快速看一下。非常感谢
    • 当使用指数回退(或基本上任何基于attempt 的回退策略)时,回退时间永远不会重置可能是出乎意料的。我的意思是 attempt 会随着 observable 生成的每个 错误而增加,无论其间有多少 good 值。
    【解决方案2】:

    也许我过度简化了这种情况,但如果我们看一下 Retry 的实现,它只是一个 Observable.Catch 对一个无限枚举的 observables:

    private static IEnumerable<T> RepeatInfinite<T>(T value)
    {
        while (true)
            yield return value;
    }
    
    public virtual IObservable<TSource> Retry<TSource>(IObservable<TSource> source)
    {
        return Observable.Catch<TSource>(QueryLanguage.RepeatInfinite<IObservable<TSource>(source));
    }
    

    所以如果我们采用这种方法,我们可以在第一次收益之后添加延迟。

    private static IEnumerable<IObservable<TSource>> RepeateInfinite<TSource> (IObservable<TSource> source, TimeSpan dueTime)
    {
        // Don't delay the first time        
        yield return source;
    
        while (true)
            yield return source.DelaySubscription(dueTime);
        }
    
    public static IObservable<TSource> RetryAfterDelay<TSource>(this IObservable<TSource> source, TimeSpan dueTime)
    {
        return RepeateInfinite(source, dueTime).Catch();
    }
    

    通过重试计数捕获特定异常的重载可以更加简洁:

    public static IObservable<TSource> RetryAfterDelay<TSource, TException>(this IObservable<TSource> source, TimeSpan dueTime, int count) where TException : Exception
    {
        return source.Catch<TSource, TException>(exception =>
        {
            if (count <= 0)
            {
                return Observable.Throw<TSource>(exception);
            }
    
            return source.DelaySubscription(dueTime).RetryAfterDelay<TSource, TException>(dueTime, --count);
        });
    }
    

    注意这里的重载是使用递归。乍一看,如果 count 是 Int32.MaxValue 之类的东西,StackOverflowException 似乎是可能的。但是,DelaySubscription 使用调度程序来运行订阅操作,因此堆栈溢出是不可能的(即使用“蹦床”)。我想这通过查看代码并不是很明显。我们可以通过将 DelaySubscription 重载中的调度程序显式设置为 Scheduler.Immediate 并传入 TimeSpan.Zero 和 Int32.MaxValue 来强制堆栈溢出。我们可以传入一个非即时调度器来更明确地表达我们的意图,例如:

    return source.DelaySubscription(dueTime, TaskPoolScheduler.Default).RetryAfterDelay<TSource, TException>(dueTime, --count);
    

    更新:添加重载以接收特定调度程序。

    public static IObservable<TSource> RetryAfterDelay<TSource, TException>(
        this IObservable<TSource> source,
        TimeSpan retryDelay,
        int retryCount,
        IScheduler scheduler) where TException : Exception
    {
        return source.Catch<TSource, TException>(
            ex =>
            {
                if (retryCount <= 0)
                {
                    return Observable.Throw<TSource>(ex);
                }
    
                return
                    source.DelaySubscription(retryDelay, scheduler)
                        .RetryAfterDelay<TSource, TException>(retryDelay, --retryCount, scheduler);
            });
    } 
    

    【讨论】:

    • 如果你愿意,你也可以用工厂替换dueTime参数,就像上面的例子一样。
    • 我喜欢第一个版本,对递归版本不太确定。理想情况下,您希望让用户通过调度程序,此时您不能再保证他们将使用哪个。
    • 感谢 Benjol 的评论。我同意。我实际上使用了一个需要调度程序的重载(更新到上面添加的代码)。你是对的,尽管用户可以立即通过。如果传入 Scheduler.Immediate,一个潜在的解决方案可能是抛出异常。
    • 我在其他人的帮助下将这个最新的部分翻译成 F#(当然,糟糕的代码和错误是我的)。代码见stackoverflow.com/questions/23404185/…
    【解决方案3】:

    这是我正在使用的:

    public static IObservable<T> DelayedRetry<T>(this IObservable<T> src, TimeSpan delay)
    {
        Contract.Requires(src != null);
        Contract.Ensures(Contract.Result<IObservable<T>>() != null);
    
        if (delay == TimeSpan.Zero) return src.Retry();
        return src.Catch(Observable.Timer(delay).SelectMany(x => src).Retry());
    }
    

    【讨论】:

    • src.Catch(Observable.Timer(delay).SelectMany(x =&gt; src).Retry())src.Catch(src.DelaySubscription(delay).Retry())有什么区别吗?
    • 看起来不像。
    • 我害怕错过了什么。我认为DelaySubscription 比Observable.Timer + SelectMany 更容易理解。感谢您分享您的解决方案,帮助我找到自己的解决方案:)
    【解决方案4】:

    根据 Markus 的回答,我写了以下内容:

    public static class ObservableExtensions
    {
        private static IObservable<T> BackOffAndRetry<T>(
            this IObservable<T> source,
            Func<int, TimeSpan> strategy,
            Func<int, Exception, bool> retryOnError,
            int attempt)
        {
            return Observable
                .Defer(() =>
                {
                    var delay = attempt == 0 ? TimeSpan.Zero : strategy(attempt);
                    var s = delay == TimeSpan.Zero ? source : source.DelaySubscription(delay);
    
                    return s
                        .Catch<T, Exception>(e =>
                        {
                            if (retryOnError(attempt, e))
                            {
                                return source.BackOffAndRetry(strategy, retryOnError, attempt + 1);
                            }
                            return Observable.Throw<T>(e);
                        });
                });
        }
    
        public static IObservable<T> BackOffAndRetry<T>(
            this IObservable<T> source,
            Func<int, TimeSpan> strategy,
            Func<int, Exception, bool> retryOnError)
        {
            return source.BackOffAndRetry(strategy, retryOnError, 0);
        }
    }
    

    我更喜欢它,因为

    • 它不修改attempts,而是使用递归。
    • 它不使用retries,而是将尝试次数传递给retryOnError

    【讨论】:

      【解决方案5】:

      这是我在研究Rxx 是如何做到的时提出的另一个稍微不同的实现。所以它在很大程度上是 Rxx 方法的缩减版。

      签名与 Markus 的版本略有不同。您指定要重试的异常类型,延迟策略采用异常和重试计数,因此每次连续重试可能会有更长的延迟,等等。

      我不能保证它是错误证明或最佳方法,但它似乎有效。

      public static IObservable<TSource> RetryWithDelay<TSource, TException>(this IObservable<TSource> source, Func<TException, int, TimeSpan> delayFactory, IScheduler scheduler = null)
      where TException : Exception
      {
          return Observable.Create<TSource>(observer =>
          {
              scheduler = scheduler ?? Scheduler.CurrentThread;
              var disposable = new SerialDisposable();
              int retryCount = 0;
      
              var scheduleDisposable = scheduler.Schedule(TimeSpan.Zero,
              self =>
              {
                  var subscription = source.Subscribe(
                  observer.OnNext,
                  ex =>
                  {
                      var typedException = ex as TException;
                      if (typedException != null)
                      {
                          var retryDelay = delayFactory(typedException, ++retryCount);
                          self(retryDelay);
                      }
                      else
                      {
                          observer.OnError(ex);
                      }
                  },
                  observer.OnCompleted);
      
                  disposable.Disposable = subscription;
              });
      
              return new CompositeDisposable(scheduleDisposable, disposable);
          });
      }
      

      【讨论】:

      • 我发现了这个 impl 中断的边缘情况(混合 Immediate 和 CurrentThread 调度程序):int a = 0; Observable.Defer(() => a++ (new Exception())).RetryWithDelay((ex, i) => TimeSpan.Zero, Scheduler.Immediate).Subscribe(i => Console.WriteLine(i));支持 Scheduler.Immediate 的一个简单修复方法是在分配 serial-disposable 之前检查订阅值是否在 Subscribe() 调用期间发生了变化。
      • 根据 RetryWithDelay() 函数的预期语义,在 OnNext 处理程序中将 retryCount 重置为 0 可能是有意义的,例如onNext: x => { 观察者.OnNext(x);重试次数 = 0; }.
      【解决方案6】:

      这是我想出的。

      不想将单个重试的项目连接到一个序列中,而是在每次重试时将源序列作为一个整体发出 - 因此运算符返回IObservable&lt;IObservable&lt;TSource&gt;&gt;。如果不希望这样做,可以简单地将其Switch()ed 重新编入一个序列。

      (背景:在我的用例中,源是一个热热序列,我GroupByUntil出现了一个关闭组的项目。如果在两次重试之间丢失了该项目,则该组永远不会关闭,从而导致内存泄漏。有序列序列允许仅对内部序列进行分组(或异常处理或...)。)

      /// <summary>
      /// Repeats <paramref name="source"/> in individual windows, with <paramref name="interval"/> time in between.
      /// </summary>
      public static IObservable<IObservable<TSource>> RetryAfter<TSource>(this IObservable<TSource> source, TimeSpan interval, IScheduler scheduler = null)
      {
          if (scheduler == null) scheduler = Scheduler.Default;
          return Observable.Create<IObservable<TSource>>(observer =>
          {
              return scheduler.Schedule(self =>
              {
                  observer.OnNext(Observable.Create<TSource>(innerObserver =>
                  {
                      return source.Subscribe(
                          innerObserver.OnNext,
                          ex => { innerObserver.OnError(ex); scheduler.Schedule(interval, self); },
                          () => { innerObserver.OnCompleted(); scheduler.Schedule(interval, self); });
                  }));
              });
          });
      }
      

      【讨论】:

      • 顺便说一句,我在 stackoverflow 上的第一篇文章 ***
      猜你喜欢
      • 2018-02-23
      • 1970-01-01
      • 2015-09-01
      • 1970-01-01
      • 2012-03-10
      • 2016-07-29
      • 2014-10-15
      • 1970-01-01
      相关资源
      最近更新 更多