【发布时间】:2014-11-15 02:55:41
【问题描述】:
以下代码将结果延迟 2 秒。我想要的是立即返回结果,但每 2 秒启动一个新的 observable。我错过了什么?
输出:
**The current output is:**
05: 1. Run
07 Result: 1
07: 2. Run
09 Result: 2
09: 3. Run
11 Result: 3
**Desired output is:**
05: 1. Run
05 Result: 1
07: 2. Run
07 Result: 2
09: 3. Run
09 Result: 3
代码:
var sources = Enumerable.Range(1, 8).Select(i =>
{
Console.WriteLine("{0}: {1}. Run", DateTimeOffset.Now.ToString("ss"), i);
return Observable.Return(i, CurrentThreadScheduler.Instance);
});
Observable.Generate(sources.GetEnumerator(), e => e.MoveNext(), e => e, e => e.Current, e => TimeSpan.FromMilliseconds(2000), ThreadPoolScheduler.Instance)
.Merge()
.Timestamp()
.Do(r =>
{
Console.WriteLine("{0} Result: {1}{2}", r.Timestamp.ToString("ss"), r.Value, Environment.NewLine);
},
ex =>
{
Console.WriteLine(ex.ToString());
},
() =>
{
Console.WriteLine("Completed");
})
.Subscribe();
【问题讨论】:
-
真的不清楚你在这里追求什么 - 我认为你对保罗的回答的评论意味着你想要可变的数据驱动间隔,但除此之外(对我训练有素的 Rx 眼)还有很多代码中的“奇怪的东西”。也许用非 rx 术语解释您要实现的目标会很有用。例如,不清楚为什么要为
sources创建IEnumerable<IObservable>,因为看起来IEnumerable<>会这样做。也不清楚sources是否应该包含指示所需间隔的数据。
标签: c# system.reactive scheduler observable