基于C# Observable实现带次数与间隔的重试机制技术问询
基于Rx.NET的可定制重试机制实现
我来帮你完善这个基于Rx.NET的重试逻辑,刚好能贴合你提出的带重试次数控制、固定重试间隔,执行时触发OnExecute、捕获异常时触发OnCatch的需求~
首先我们先明确核心需求点:
- 支持配置最大重试次数
- 每次重试前等待指定时间间隔
- 执行逻辑通过
OnExecute()触发,且每次重试都要重新执行该方法 - 捕获到异常时立即调用
OnCatch()回调处理
1. 先定义配套的请求类与动作接口
首先我们需要把你提到的GenericRetryExecutorRequest<T>和动作接口补全,让重试逻辑的参数和回调更清晰:
// 重试请求配置类 public class GenericRetryExecutorRequest<T> { // 最大重试次数(注意:这里指失败后重试的次数,总尝试次数=重试次数+1) public int RetryCount { get; set; } // 每次重试的间隔时间 public TimeSpan Interval { get; set; } // 重试相关的动作回调集合 public IGenericRetryActions<T> GenericRetryActions { get; set; } } // 重试动作回调接口 public interface IGenericRetryActions<T> { // 要执行的核心业务逻辑 T OnExecute(); // 捕获到异常时的回调方法 void OnCatch(Exception ex); }
2. 实现重试逻辑的核心方法
这里我们用RetryWhen来替代你原来的Timer+Retry方案,因为RetryWhen能更精准地控制仅当异常发生时才触发重试,而不是定时轮询,更符合重试机制的设计:
public static IObservable<T> Retry<T>(GenericRetryExecutorRequest<T> request) { // 用Defer包装执行逻辑,确保每次重试都会重新执行OnExecute(避免复用之前的执行序列) return Observable.Defer(() => Observable.Start(() => request.GenericRetryActions.OnExecute())) .RetryWhen(errors => errors // 给每个异常绑定重试索引,用来判断是否达到最大重试次数 .Select((ex, retryIndex) => new { Exception = ex, RetryIndex = retryIndex }) // 控制重试次数:仅当重试索引小于配置的RetryCount时继续重试 .TakeWhile(retryInfo => retryInfo.RetryIndex < request.RetryCount) // 触发异常回调:每次捕获异常时调用OnCatch .Do(retryInfo => request.GenericRetryActions.OnCatch(retryInfo.Exception)) // 延迟指定间隔后再发起下一次重试 .Delay(_ => Observable.Timer(request.Interval))); }
3. 代码逻辑说明
对上面的核心实现拆解一下:
Observable.Defer:确保每次重试都会重新创建执行序列,保证OnExecute()在每次重试时都会被重新调用,而不是复用第一次的执行结果。Observable.Start:把同步的OnExecute()方法包装成Observable,默认会在后台线程执行(如果需要指定调度器,可以传入Scheduler参数,比如Scheduler.CurrentThread)。RetryWhen:比Rx内置的Retry方法更灵活,允许我们自定义重试的触发条件、间隔和次数。Select((ex, retryIndex)):给每个异常加上重试索引,用来判断是否已经达到配置的最大重试次数。TakeWhile:当重试索引小于配置的RetryCount时,继续重试;超过次数后会停止重试,将最终的异常抛给订阅者。Do:在每次捕获异常时立即调用OnCatch()回调,你可以在这里做异常日志、告警等操作。Delay:通过Observable.Timer实现重试间隔,确保每次异常后等待指定时间再发起下一次重试。
4. 使用示例
我们可以写一个简单的测试类来验证这个重试机制:
// 实现重试动作接口的示例类 public class SampleRetryActions : IGenericRetryActions<string> { private int _attemptNumber = 0; public string OnExecute() { _attemptNumber++; Console.WriteLine($"开始执行第 {_attemptNumber} 次尝试"); // 模拟前2次执行失败,第3次成功 if (_attemptNumber < 3) { throw new InvalidOperationException("模拟业务执行失败"); } return "业务执行成功!"; } public void OnCatch(Exception ex) { Console.WriteLine($"捕获到异常:{ex.Message},等待1秒后重试..."); } } // 调用重试机制 var retryRequest = new GenericRetryExecutorRequest<string> { RetryCount = 3, // 最多重试3次 Interval = TimeSpan.FromSeconds(1), // 每次重试间隔1秒 GenericRetryActions = new SampleRetryActions() }; // 订阅重试结果 Retry(retryRequest) .Subscribe( successResult => Console.WriteLine($"最终结果:{successResult}"), finalError => Console.WriteLine($"所有重试失败:{finalError.Message}"));
运行这段代码后,你会看到如下输出:
开始执行第 1 次尝试 捕获到异常:模拟业务执行失败,等待1秒后重试... 开始执行第 2 次尝试 捕获到异常:模拟业务执行失败,等待1秒后重试... 开始执行第 3 次尝试 最终结果:业务执行成功!
内容的提问来源于stack exchange,提问作者ohadinho
相关产品推荐
相关产品推荐

