【问题标题】:Is there a simpler way to have an IObservable be asynchronously dependent upon another IObservable?有没有更简单的方法让 IObservable 异步依赖于另一个 IObservable?
【发布时间】:2011-02-07 22:07:51
【问题描述】:

我是 RX 新手,我希望的场景运行良好,但在我看来,必须有一种更简单或更优雅的方式来实现这一点。我有一个IObservable<T>,我想订阅它,通过触发一个异步操作,为它看到的每个 T 生成一个 U,我最终得到一个 IObservable<U>,

到目前为止,我所拥有的(效果很好,但看起来很麻烦)使用中间事件流,如下所示:

public class Converter {
  public event EventHandler<UArgs> UDone;
  public IConnectableObservable<U> ToUs(IObservable<T> ts) {
    var us = Observable.FromEvent<UArgs>(this, "UDone").Select(e => e.EventArgs.U).Replay();
    ts.Subscribe(t => Observable.Start(() => OnUDone(new U(t))));
    return us;
  }
  private void OnUDone(U u) {
    var uDone = UDone;
    if (uDone != null) {
      uDone(this, u);
    }
  }
}

...

var c = new Converter();
IConnectableObservable<T> ts = ...;
var us = c.ToUs(ts);
us.Connect();

...

我确定我错过了一种更简单的方法来做到这一点......

【问题讨论】:

  • 你确定要Replay吗?

标签: asynchronous system.reactive


【解决方案1】:

SelectMany 应该做你需要做的,把IO&lt;IO&lt;T&gt;&gt; 弄平

Observable.Range(1, 10)
        .Select(ii => Observable.Start(() => 
             string.Format("{0} {1}", ii, Thread.CurrentThread.ManagedThreadId)))
        .SelectMany(id=>id)
        .Subscribe(Console.WriteLine);

【讨论】:

  • 谢谢,斯科特。我知道 SelectMany,但出于某种原因,我没有考虑过扁平化(单项)可观察对象的可观察对象的问题。
  • 然而,我注意到了一个问题 - 它在事件流的中途用 OutOfMemoryException 炸毁了我的进程。我的方法没有做到这一点......
【解决方案2】:

这正是SelectMany 的用途:

IObservable<int> ts

IObservable<string> us = ts.SelectMany(t => StartAsync(t));

us.Subscribe(u => 
    Console.WriteLine("StartAsync completed with {0}", u));

...

private IObservable<string> StartAsync(int t)
{
    return Observable.Return(t.ToString())
        .Delay(TimeSpan.FromSeconds(1));
}

请记住,如果StartAsync 的完成时间可变,您可能会以与输入值不同的顺序接收输出值。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-04-23
    相关资源
    最近更新 更多