【问题标题】:How to subscribe to, but buffer data from, an IObservable until another IObservable has published?在另一个 Observable 发布之前,如何订阅但缓冲来自 Observable 的数据?
【发布时间】:2011-09-20 05:44:18
【问题描述】:

我想:

  1. 立即订阅IObservable<T>,但立即开始缓冲收到的任何T(即我的IObserver<T> 尚未看到)。
  2. 做一些工作。
  3. 工作完成后,将缓冲区刷新到我的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&lt;bool&gt; r,它本身依赖于 IObservable&lt;int&gt; s1IObservable&lt;bool&gt; s2s1 是我无法控制的流,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;
    }        
}

【问题讨论】:

    标签: c#-4.0 system.reactive


    【解决方案1】:

    这对我有用:

    var r = Observable.Create<int>(o =>
    {
        var rs = new ReplaySubject<int>();
        var subscription1 = s1.Subscribe(rs);
        var query = from f in s2.Take(1) select rs.AsObservable();
        var subscription2 = query.Switch().Subscribe(o);
        return new CompositeDisposable(subscription1, subscription2);
    });
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-12-03
      • 2018-03-13
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多