【发布时间】:2016-02-09 07:41:11
【问题描述】:
在 C# 中使用 Rx 我正在尝试创建对 REST API 的轮询请求。我面临的问题是,Observable 需要按顺序发送响应。表示如果请求 A 在 X 时间到达,请求 B 在 X + dx 时间到达,并且 B 的响应在 A 之前到达,则 Observable 表达式应忽略或取消请求 A。
我编写了一个示例代码来尝试描述该场景。如何修复它以仅获取最新响应并取消或忽略以前的响应。
class Program
{
static int i = 0;
static void Main(string[] args)
{
GenerateObservableSequence();
Console.ReadLine();
}
private static void GenerateObservableSequence()
{
var timerData = Observable.Timer(TimeSpan.Zero,
TimeSpan.FromSeconds(1));
var asyncCall = Observable.FromAsync<int>(() =>
{
TaskCompletionSource<int> t = new TaskCompletionSource<int>();
i++;
int k = i;
var rndNo = new Random().Next(3, 10);
Task.Delay(TimeSpan.FromSeconds(rndNo)).ContinueWith(r => { t.SetResult(k); });
return t.Task;
});
var obs = from t in timerData
from data in asyncCall
select data;
var hot = obs.Publish();
hot.Connect();
hot.Subscribe(j =>
{
Console.WriteLine("{0}", j);
});
}
}
@Enigmativity 回答后:添加 Polling Aync 功能以始终获取最新响应:
public static IObservable<T> PollingAync<T> (Func<Task<T>> AsyncCall, double TimerDuration)
{
return Observable
.Create<T>(o =>
{
var z = 0L;
return
Observable
.Timer(TimeSpan.Zero, TimeSpan.FromSeconds(TimerDuration))
.SelectMany(nr =>
Observable.FromAsync<T>(AsyncCall),
(nr, obj) => new { nr, obj})
.Do(res => z = Math.Max(z, res.nr))
.Where(res => res.nr >= z)
.Select(res => res.obj)
.Subscribe(o);
});
}
【问题讨论】:
标签: c# recursion system.reactive reactive-programming polling