You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.12 22:05:12