【问题标题】:How to handle exceptions in OnNext when using ObserveOn?使用 ObserveOn 时如何处理 OnNext 中的异常?
【发布时间】:2012-06-25 00:48:47
【问题描述】:

当我使用ObserveOn(Scheduler.ThreadPool) 时观察者在OnNext 中引发错误时,我的应用程序将终止。我发现处理这个问题的唯一方法是使用下面的自定义扩展方法(除了确保 OnNext 永远不会引发异常)。然后确保每个ObserveOn 后跟一个ExceptionToError

    public static IObservable<T> ExceptionToError<T>(this IObservable<T> source) {
        var sub = new Subject<T>();
        source.Subscribe(i => {
            try {
                sub.OnNext(i);
            } catch (Exception err) {
                sub.OnError(err);
            }
        }
            , e => sub.OnError(e), () => sub.OnCompleted());
        return sub;
    }

但是,这感觉不对。有没有更好的方法来解决这个问题?

示例

此程序因未捕获的异常而崩溃。

class Program {
    static void Main(string[] args) {
        try {
            var xs = new Subject<int>();

            xs.ObserveOn(Scheduler.ThreadPool).Subscribe(x => {
                Console.WriteLine(x);
                if (x % 5 == 0) {
                    throw new System.Exception("Bang!");
                }
            }, ex => Console.WriteLine("Caught:" + ex.Message)); // <- not reached

            xs.OnNext(1);
            xs.OnNext(2);
            xs.OnNext(3);
            xs.OnNext(4);
            xs.OnNext(5);
        } catch (Exception e) {
            Console.WriteLine("Caught : " + e.Message); // <- also not reached
        } finally {

            Console.ReadKey();
        }
    }
}

【问题讨论】:

    标签: c# system.reactive


    【解决方案1】:

    我们在 Rx v2.0 中解决了这个问题,从 RC 版本开始。您可以在我们的博客http://blogs.msdn.com/rxteam 上阅读所有相关信息。它基本上归结为管道本身中更严格的错误处理,结合了 SubscribeSafe 扩展方法(将订阅期间的错误重定向到 OnError 通道)和 IScheduler 上的 Catch 扩展方法(将调度程序与异常处理逻辑包装在调度行动)。

    关于这里提出的 ExceptionToError 方法,它有一个缺陷。回调运行时,IDisposable 订阅对象仍然可以为空;有一个基本的竞争条件。要解决此问题,您必须使用 SingleAssignmentDisposable。

    【讨论】:

    【解决方案2】:

    订阅错误和可观察错误之间存在差异。快速测试:

    var xs = new Subject<int>();
    
    xs.Subscribe(x => { Console.WriteLine(x); if (x % 3 == 0) throw new System.Exception("Error in subscription"); }, 
                 ex => Console.WriteLine("Error in source: " + ex.Message));
    

    用这个运行,你会在源代码中得到一个很好处理的错误:

    xs.OnNext(1);
    xs.OnNext(2);
    xs.OnError(new Exception("from source"));
    

    运行这个,你会在订阅中得到一个未处理的错误:

    xs.OnNext(1);
    xs.OnNext(2);
    xs.OnNext(3);
    

    您的解决方案所做的是在订阅中出错并在源中使它们出错。而且您已经在原始流上完成了此操作,而不是基于每个订阅。您可能打算也可能不打算这样做,但这几乎肯定是错误的。

    做到这一点的“正确”方法是将您需要的错误处理直接添加到它所属的订阅操作中。如果不想直接修改订阅功能,可以使用一个小帮手:

    public static Action<T> ActionAndCatch<T>(Action<T> action, Action<Exception> catchAction)
    {
        return item =>
        {
            try { action(item); }
            catch (System.Exception e) { catchAction(e); }
        };
    }
    

    现在使用它,再次显示不同错误之间的区别:

    xs.Subscribe(ActionAndCatch<int>(x => { Console.WriteLine(x); if (x % 3 == 0) throw new System.Exception("Error in subscription"); },
                                     ex => Console.WriteLine("Caught error in subscription: " + ex.Message)),
                 ex => Console.WriteLine("Error in source: " + ex.Message));
    

    现在我们可以(分别)处理源中的错误和订阅中的错误。当然,这些动作中的任何一个都可以在一个方法中定义,使上面的代码变得简单(可能):

    xs.Subscribe(ActionAndCatch(Handler, ExceptionHandler), SourceExceptionHandler);
    

    编辑

    在 cmets 中,我们开始讨论订阅中的错误指向流本身的错误,并且您不希望该流上的其他订阅者。这是一个完全不同类型的问题。我倾向于编写一个可观察的Validate 扩展来处理这种情况:

    public static IObservable<T> Validate<T>(this IObservable<T> source, Predicate<T> valid)
    {
        return Observable.Create<T>(o => {
            return source.Subscribe(
                x => {
                    if (valid(x)) o.OnNext(x);
                    else       o.OnError(new Exception("Could not validate: " + x));
                }, e => o.OnError(e), () => o.OnCompleted()
            );
        });
    }
    

    然后简单易用,无需混合隐喻(仅在源代码中出现错误):

    xs
    .Validate(x => x != 3)
    .Subscribe(x => Console.WriteLine(x),
                 ex => Console.WriteLine("Error in source: " + ex.Message));
    

    如果您仍然希望在 Subscribe 中抑制异常,您应该使用其他讨论的方法之一。

    【讨论】:

    • 我同意在源代码中将其作为错误可能是不正确的。但是,对于您的解决方案,每个观察者都必须自己进行错误处理。我的方法可以改为ExceptionToAction&lt;T&gt;(this IObservable&lt;T&gt; source, Action&lt;Exception&gt; handler)。不管怎样,我觉得这还不太合适。
    • 我不太明白。在每种情况下,您都需要提供异常处理程序。您可以在订阅时附加异常处理程序(首选方法),也可以将异常处理程序附加到源 observable。后者违背了一般的 RX 设计,因为您正在修改 observable 的未来订阅,而不是 observable 本身。非常“有副作用”。您建议的更改 (ExceptionToAction) 至少需要,但它仍然没有明确表明这是订阅中的例外情况。
    • 还要注意你说“每个观察者必须自己做错误处理”。如果您想“扩展”某些东西,也许可以创建一个标准的观察者来处理错误,而仅仅是那个?这就是ActionAndCatch 的本质,除了它是OnNext 的最低限度的观察者。
    • 这是一个有趣的想法!类似于SubscribeAndCatch。我的目标有两个:1)我想要一个安全网,以确保当流中出现异常时我的整个应用程序不会崩溃,以及 2)防止流中的无效数据。安全网必须在源头(或调度程序?)以捕获所有内容,而另一个可能靠近观察者。在我的情况下,缺少一个值会使流无效,因此调用 OnError 对我来说似乎是正确的选择,但我可能弄错了吗?
    • 现在我们正在讨论流中的无效数据!好的,根据这个坏男孩(Rx 设计指南:go.microsoft.com/fwlink/?LinkID=205219),不应在订阅中抛出任何异常,但它们应该通过隧道返回到流中以阻止它。这是您所追求的场景 - 在某些条件下终止流。因此,换句话说,继续使用Observable.Create(或SelectManyReturnThrow 的组合)创建自定义运算符并使用它来验证流。但是在其他地方处理订阅错误(这是一个不同的野兽)。
    【解决方案3】:

    您当前的解决方案并不理想。正如一位 Rx 人 here 所说:

    Rx 运算符不会捕获在调用 OnNext、OnError 或 OnCompleted 时发生的异常。这是因为我们期望 (1) 观察者实现者最清楚如何处理这些异常,我们不能对它们做任何合理的事情; (2) 如果发生异常,那么我们希望它冒泡而不被 Rx 处理.

    您当前的解决方案让 IObservable 来处理 IObserver 抛出的错误,这没有意义,因为 IObservable 在语义上应该不知道观察它的事物。考虑以下示例:

    var errorFreeSource = new Subject<int>();
    var sourceWithExceptionToError = errorFreeSource.ExceptionToError();
    var observerThatThrows = Observer.Create<int>(x =>
      {
          if (x % 5 == 0)
              throw new Exception();
      },
      ex => Console.WriteLine("There's an argument that this should be called"),
      () => Console.WriteLine("OnCompleted"));
    var observerThatWorks = Observer.Create<int>(
        x => Console.WriteLine("All good"),
        ex => Console.WriteLine("But definitely not this"),
        () => Console.WriteLine("OnCompleted"));
    sourceWithExceptionToError.Subscribe(observerThatThrows);
    sourceWithExceptionToError.Subscribe(observerThatWorks);
    errorFreeSource.OnNext(1);
    errorFreeSource.OnNext(2);
    errorFreeSource.OnNext(3);
    errorFreeSource.OnNext(4);
    errorFreeSource.OnNext(5);
    Console.ReadLine();
    

    这里源或observerThatWorks没有问题,但由于另一个Observer的不相关错误将调用其OnError。要阻止不同线程中的异常结束进程,您必须在该线程中捕获它们,因此在您的观察者中放置一个 try/catch 块。

    【讨论】:

      【解决方案4】:

      我查看了应该解决此问题的本机 SubscribeSafe 方法,但我无法使其工作。此方法有一个接受 IObserver&lt;T&gt; 的重载:

      // Subscribes to the specified source, re-routing synchronous exceptions during
      // invocation of the IObservable<T>.Subscribe(IObserver<T>) method to the
      // observer's IObserver<T>.OnError(Exception) channel. This method is typically
      // used when writing query operators.
      public static IDisposable SubscribeSafe<T>(this IObservable<T> source,
          IObserver<T> observer);
      

      我尝试传递由Observer.Create 工厂方法创建的观察者,但onNext 处理程序中的异常继续使进程崩溃¹,就像它们对普通Subscribe 所做的那样。所以我最终编写了我自己的SubscribeSafe 版本。这个接受三个处理程序作为参数,并将onNextonCompleted 处理程序抛出的任何异常集中到onError 处理程序。

      /// <summary>Subscribes an element handler, an error handler, and a completion
      /// handler to an observable sequence. Any exceptions thrown by the element or
      /// the completion handler are propagated through the error handler.</summary>
      public static IDisposable SubscribeSafe<T>(this IObservable<T> source,
          Action<T> onNext, Action<Exception> onError, Action onCompleted)
      {
          // Arguments validation omitted
          var disposable = new SingleAssignmentDisposable();
          disposable.Disposable = source.Subscribe(
              value =>
              {
                  try { onNext(value); } catch (Exception ex) { onError(ex); disposable.Dispose(); }
              }, onError, () =>
              {
                  try { onCompleted(); } catch (Exception ex) { onError(ex); }
              }
          );
          return disposable;
      }
      

      注意,onError 处理程序中未处理的异常仍会使进程崩溃!

      ¹ 只有在 ThreadPool 上异步调用处理程序时才会引发异常。

      【讨论】:

        【解决方案5】:

        你是对的 - 它应该感觉很糟糕。像这样使用和返回主题不是一个好方法。

        至少你应该像这样实现这个方法:

        public static IObservable<T> ExceptionToError<T>(this IObservable<T> source)
        {
            return Observable.Create<T>(o =>
            {
                var subscription = (IDisposable)null;
                subscription = source.Subscribe(x =>
                {
                    try
                    {
                        o.OnNext(x);
                    }
                    catch (Exception ex)
                    {
                        o.OnError(ex);
                        subscription.Dispose();
                    }
                }, e => o.OnError(e), () => o.OnCompleted());
                return subscription;
            });
        }
        

        请注意,没有使用任何主题,如果我发现一个错误,我会处理订阅以防止序列继续超过错误。

        但是,为什么不在订阅中添加 OnError 处理程序。有点像这样:

        var xs = new Subject<int>();
        
        xs.ObserveOn(Scheduler.ThreadPool).Subscribe(x =>
        {
            Console.WriteLine(x);
            if (x % 5 == 0)
            {
                throw new System.Exception("Bang!");
            }
        }, ex => Console.WriteLine(ex.Message));
        
        xs.OnNext(1);
        xs.OnNext(2);
        xs.OnNext(3);
        xs.OnNext(4);
        xs.OnNext(5);
        

        此代码在订阅中正确捕获错误。

        另一种方法是使用Materialize 扩展方法,但这可能有点矫枉过正,除非上述解决方案不起作用。

        【讨论】:

        • 在您的 OnError 示例中不会调用 OnError,除非添加了 ExceptionToError(使用控制台应用程序和单元测试项目尝试)。你测试了吗?非常感谢 ExceptionToError 的改进,但是仍然感觉不自然。也许应该将错误处理添加到调度程序中?即一个 SafeScheduler 包装器。
        • 我查看了Materialize,但不知道如何使用它。没有ExceptionToError 就不会调用 OnError。测试它: public static IObservable Log(IObservable source) { return source.Materialize().Select(n => { Trace.WriteLine("Error hunt: " + n); return n ; }).Dematerialize(); }
        • @Herman - 你能提供一些代码来证明你的错误吗?
        • @Enigmativity Herman 的意思是,如果您运行您提供的代码(OnErrorSubscribe 方法操作),整个程序仍然会抛出异常并提前“终止”。他希望他处理这个错误。 OnError 方法重载仅处理源中的错误,而不是订阅中的错误。有关详细信息,请参阅我的答案。
        猜你喜欢
        • 2014-10-10
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2022-01-11
        • 1970-01-01
        相关资源
        最近更新 更多