为何BehaviorSubject会报告无关Observable的未处理异常?
上下文
我正在排查一个复杂问题:某个BehaviorSubject会向观察者发送错误,但实际上没有任何代码调用它的OnError方法。唯一的关联是它与另一个存在未处理异常的Observable都在WPF应用的UI线程上执行。
以下代码大致还原了我的实现,但可能缺失了实际触发问题的关键逻辑:
private readonly BehaviorSubject<int> _exampleSubject; public void NotifyChanged(int newValue) { _exampleSubject.OnNext(newValue); } public IObservable<int> _ConstructStream() { return _exampleSubject .Do( onNext: count => Console.WriteLine($"_exampleSubject.OnNext: {count} (thread: {Thread.CurrentThread.ManagedThreadId})"), onError: exception => Console.WriteLine($"_exampleSubject.OnError: {exception} (thread: {Thread.CurrentThread.ManagedThreadId})"), onCompleted: () => Console.WriteLine($"_exampleSubject.OnCompleted (thread: {Thread.CurrentThread.ManagedThreadId})")); }
不知为何,_exampleSubject会抛出异常,这显然不合理,因为从未调用过它的OnError。而另一处Observable的未处理异常似乎对它产生了干扰。
遗憾的是我无法创建最小复现示例,这似乎是一个仅在完整代码中才会出现的竞态条件。不过我已复现了部分问题,详情如下:
环境搭建
要复现该问题,请使用默认WPF模板,并对App.xaml.cs做如下修改:
// Target framework: // net7.0-windows // Packages: // System.Reactive 6.0.0 namespace WpfApp1 { public partial class App : Application { protected override void OnStartup(StartupEventArgs e) { base.OnStartup(e); // ... Call the example methods here. } // ... Define the example methods here. } }
~~## 示例2
private async Task _Example_2() { var dispatcherSynchronizationContext = new DispatcherSynchronizationContext(Application.Current.Dispatcher); // Adding this code would no longer reproduce the issue: // using var disposable = Observable var disposable = Observable .Timer(TimeSpan.FromMilliseconds(100)) .ObserveOn(dispatcherSynchronizationContext) .SubscribeOn(dispatcherSynchronizationContext) .Subscribe(_ => throw new InvalidOperationException()); // Wait a bit for things to play out. await Task.Delay(TimeSpan.FromMilliseconds(300)); }
抛出异常时,主线程会冻结,整个窗口无响应,同一调度器上的其他操作也停止执行。未绑定到该调度器的Observable会短暂运行后,应用自动关闭。
这符合未处理异常会终止应用的预期,但奇怪的是没有任何日志信息记录该异常。
问题1: 若为disposable添加using语句,问题就不再复现。我本以为disposal会在函数末尾延迟后执行,但实际似乎提前执行了。~~
重新审视后,我发现自己混淆了示例逻辑:300ms后方法返回才会触发dispose,我之前误以为是Timeout.InfiniteTimeSpan的情况。
示例1
有趣的是,使用Observable.FromAsync结合Concat也能复现相同问题,但仅在切换线程前无延迟时才会出现。
private async Task _Example_1() { var dispatcherSynchronizationContext = new DispatcherSynchronizationContext(Application.Current.Dispatcher); Observable .Timer(TimeSpan.FromMilliseconds(100)) .ObserveOn(dispatcherSynchronizationContext) .SubscribeOn(dispatcherSynchronizationContext) .Select(index => Observable.FromAsync(async () => { // Adding this code would no longer reproduce the issue: // await Task.Delay(TimeSpan.FromMilliseconds(100)); throw new InvalidOperationException(); })) .Concat() .Subscribe(); // Wait a bit for things to play out. await Task.Delay(TimeSpan.FromMilliseconds(300)); }
问题2: 线程切换和暂停执行时似乎存在异常行为,添加延迟后Dispatcher本应崩溃,但实际却未复现问题。
内容的提问来源于stack exchange,提问作者user272507

