【发布时间】:2016-01-22 00:36:47
【问题描述】:
信息:
该方法的目的是从流中获取秘密,从流中获取uri。然后用秘密和一些附加参数type和name点击uri来获取一些数据。
Secrets.GetSecret() 是 IObservable<string> 和 Urls.GetHostUrl() 是 IObservable<Uri>
问题
我遇到的问题是,如果 dataFetcher.GetFromUrl 调用引发异常,那么 observable 将终止。 RetryAfterDelay 扩展方法似乎根本不起作用。
我想要做的是能够捕获并记录异常,但理想情况下没有可观察的终止。或者如果它必须终止然后重新订阅原始流,那么它的行为就好像它已经记录/吞下了异常。
原因是网络/主机不是最稳定的,所以有时会出错。
方法
private IObservable<string> GetFromHostUrl(string type, string name, int retryTime)
{
return Observable.Interval(TimeSpan.FromSeconds(retryTime), schedulerProvider.Default).StartWith(-1L)
.CombineLatest(Observable.Defer(() => Secrets.GetSecret()),
Urls.GetHostUrl(),
(_, secret, hostUrl) => dataFetcher.GetFromUrl(hostUrl, type, name, secret).Result)
.Do(s =>
{
// log success
},
e =>
{
// log failure
})
.DistinctUntilChanged()
.RetryAfterDelay(TimeSpan.FromSeconds(retryTime), schedulerProvider.Default)
.Replay(1)
.RefCount();
}
扩展方法
// c/o: http://stackoverflow.com/questions/18978523/write-an-rx-retryafter-extension-method
public static class ExtensionMethods
{
private static IEnumerable<IObservable<TSource>> RepeateInfinite<TSource>(IObservable<TSource> source, TimeSpan dueTime, IScheduler scheduler)
{
// Don't delay the first time
yield return source;
while (true)
{
yield return source.DelaySubscription(dueTime, scheduler);
}
}
public static IObservable<TSource> RetryAfterDelay<TSource>(this IObservable<TSource> source, TimeSpan dueTime, IScheduler scheduler)
{
return RepeateInfinite(source, dueTime, scheduler).Catch();
}
}
【问题讨论】:
标签: c# .net exception-handling system.reactive