【发布时间】:2011-09-20 05:44:18
【问题描述】:
我想:
- 立即订阅
IObservable<T>,但立即开始缓冲收到的任何T(即我的IObserver<T>尚未看到)。 - 做一些工作。
- 工作完成后,将缓冲区刷新到我的
IObserver<T>并继续
订阅是发生的第一件事,这一点非常重要。
以“大理石图”的形式,我追求这样的东西......
Time T+1 2 3 4 5 6 7 8
s1:IObservable<int> 1 2 3 4 5 6 7 8
s2:IObservable<bool> t
r: IObservable<int> 1 3 4 5 6 7 8
2
...在 T+1 时,我订阅了一个 IObservable<bool> r,它本身依赖于 IObservable<int> s1 和 IObservable<bool> s2。 s1 是我无法控制的流,s2 是我可以控制的流(主题),publish 在工作完成时开启。
我认为 SkipUntil 会帮助我,但这不会缓冲在依赖 IObservable 完成之前收到的事件。
这里有一些我认为可以工作的代码,但不是因为SkipUntil 不是缓冲区。
var are = new AutoResetEvent(false);
var events = Observable.Generate(1, i => i < 12, i => i + 1, i => i, i => TimeSpan.FromSeconds(1));
events.Subscribe(x => Console.WriteLine("events:" + x), () => are.Set());
var subject = new Subject<int>();
var completed = subject.AsObservable().Delay(TimeSpan.FromSeconds(5));
Console.WriteLine("Subscribing to events...");
events.SkipUntil(completed).Subscribe(x=> Console.WriteLine("events.SkipUntil(completed):"+ x));
Console.WriteLine("Subscribed.");
completed.Subscribe(x => Console.WriteLine("Completed"));
subject.OnNext(10);
are.WaitOne();
Console.WriteLine("Done");
我知道各种 Buffer 方法,但在这种情况下它们似乎不合适,因为我并没有真正在这里缓冲,只是在订阅开始时协调活动。
更新
我已将 Enigmativity 的响应概括为以下可能有用的扩展方法:
public static class ObservableEx
{
public static IObservable<TSource> BufferUntil<TSource, TCompleted>(this IObservable<TSource> source, IObservable<TCompleted> completed)
{
var observable = Observable.Create<TSource>(o =>
{
var replaySubject = new ReplaySubject<TSource>();
var sub1 = source.Subscribe(replaySubject);
var query =
completed.Take(1).Select(
x => replaySubject.AsObservable());
var sub2 = query.Switch().Subscribe(o);
return new CompositeDisposable(sub1, sub2);
});
return observable;
}
}
【问题讨论】: