Rx .Net捕获错误后序列终止问题求助
Rx.NET定时轮询任务异常后持续运行的解决方案
问题核心原因
你当前的Catch操作符是绑定在整个轮询序列末尾,一旦某个Job抛出异常,整个Interval->Select->Concat的主序列会直接终止,转而订阅OnErrorFunc返回的Observable。如果返回Observable.Empty,订阅会直接结束;返回单个Job的Observable,执行完后也不会继续触发后续的轮询信号,自然无法持续运行。
正确实现方式:给单个任务添加错误隔离
要让单个任务失败不影响整体轮询,需要把错误处理逻辑放到每个子任务的Observable上,而不是整个主序列的末尾。这样单个任务抛出异常时,只会处理该任务的错误,主序列的轮询信号会继续触发。
修改后的核心代码示例
以StartAfterTimeSpan方法为例:
public void StartAfterTimeSpan() { StartTask(); Observable.Interval(Timespan) // 给每个子任务单独添加错误处理 .Select(l => Observable.FromAsync(() => Job(m_CTS.Token)) .Catch((Exception ex) => { OnErrorFunc(ex); return Observable.Empty<System.Reactive.Unit>(); })) .Concat() .Subscribe(m_CTS.Token); }
StartInstantly方法同理修改:
public void StartInstantly() { StartTask(); Observable.Interval(Timespan) .StartWith(-1L) .Select(l => Observable.FromAsync(() => Job(m_CTS.Token)) .Catch((Exception ex) => { OnErrorFunc(ex); return Observable.Empty<System.Reactive.Unit>(); })) .Concat() .Subscribe(m_CTS.Token); }
关于Retry操作符的使用场景
如果需要单个任务失败后立即重试几次,再继续等待下一轮轮询,可以在子任务中加入Retry操作符:
.Select(l => Observable.FromAsync(() => Job(m_CTS.Token)) .Retry(2) // 失败后立即重试2次 .Catch((Exception ex) => { OnErrorFunc(ex); return Observable.Empty<System.Reactive.Unit>(); }))
如果需要带延迟的重试,可以使用RetryWhen(自定义重试间隔逻辑),但如果只是要保证轮询不中断,仅给单个任务加Catch就足够。
额外优化点
- 你的
OnErrorFunc原本返回IObservable<Unit>,现在在子任务的Catch中直接调用它即可,不需要依赖它的返回值替换序列。 - 确保
Stop()方法的CancellationToken能正确传递到所有子任务,保证停止操作生效。
内容的提问来源于stack exchange,提问作者Bluescreenterror
相关产品推荐
相关产品推荐

