【问题标题】:How to pass an async method inside the Observable.Do extension method?如何在 Observable.Do 扩展方法中传递异步方法?
【发布时间】:2014-12-19 05:58:51
【问题描述】:

给定:

  1. IObservable<T> src
  2. async Task F(T){...}
  3. F 只能按顺序调用。所以await F(x);await F(y); 可以,但Task.Factory.ContinueWhenAll(new[]{F(x),F(y)}, _ => {...}); 是错误的,因为F(x)F(y) 不能同时运行。

我很清楚await src.Do(F) 是错误的,因为它会同时运行F

我的问题是如何正确地做到这一点?

【问题讨论】:

    标签: system.reactive


    【解决方案1】:

    Observable.Do 仅用于副作用,不适用于顺序组合。 SelectMany 可能是您想要的。从 Rx 2.0 开始,SelectMany 的重载使得使用 Task<T> 组合 observable 变得非常容易。 (请注意additional concurrency 可能由这些和类似的 Task/Observable 合作运营商引入。)

    var q = from value in src
            from _ in F(value.X).AsVoidAsync()  // See helper definition below
            from __ in F(value.Y).AsVoidAsync()
            select value;
    

    但是,基于您特别询问Do 运算符这一事实,我怀疑 src 可能包含多个值,并且您不希望重叠调用 F 也适用于每个值。在这种情况下,考虑SelectMany 实际上就像Select->Merge;因此,您可能想要的是Select->Concat

    // using System.Reactive.Threading.Tasks
    
    var q = src.Select(x => Observable.Defer(() => F(x).ToObservable())).Concat();
    

    不要忘记使用Defer,因为F(x)热门

    AsVoidAsync 扩展

    IObservable<T> 需要 T,而 Task 代表 void,因此 Rx 的转换运算符要求我们从 Task 中获取 Task<T>。我倾向于将 Rx 的 System.Reactive.Unit 结构用于 T

    public static class TaskExtensions
    {
      public static Task<Unit> AsVoidAsync(this Task task)
      {
        return task.ContinueWith(t =>
        {
          var tcs = new TaskCompletionSource<Unit>();
    
          if (t.IsCanceled)
          {
            tcs.SetCanceled();
          }
          else if (t.IsFaulted)
          {
            tcs.SetException(t.Exception);
          }
          else
          {
            tcs.SetResult(Unit.Default);
          }
    
          return tcs.Task;
        },
        TaskContinuationOptions.ExecuteSynchronously)
        .Unwrap();
      }
    }
    

    或者,您总是可以只调用专用的ToObservable 方法。

    【讨论】:

    • 还有一个带有 max_concurrent 参数的 Merge 重载。通过将此参数设置为1,可以达到类似的效果。 msdn.microsoft.com/en-us/library/hh211914(v=vs.103).aspx查看本题答案:stackoverflow.com/questions/25436542/…
    • 这不是我推荐的方法。 1. Merge(1) 的行为类似于 Concat 并不是很明显。 Concat 具有正确的语义。 2. Concat 的实现已经调用了 Merge(1)(证明我的第一点)。 3. Merge(1) 的成本理论上可能会略高于 Concat,因为前者旨在实现有界并发。下限不一定作为热路径实现。 (由于#2,我知道这是一个有争议的问题,但这是一个可能会改变的内部实现细节。)
    【解决方案2】:

    最简单的方法是使用 TPL Dataflow,默认情况下会限制为单个并发调用,即使对于异步方法也是如此:

    var block = new ActionBlock<T>(F);
    src.Subscribe(block.AsObserver());
    await block.Completion;
    

    或者,您可以订阅一个您自己限制的异步方法:

    var semaphore = new SemaphoreSlim(1);
    src.Do(async x =>
    {
      await semaphore.WaitAsync();
      try
      {
        await F(x);
      }
      finally
      {
        semaphore.Release();
      }
    });
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2013-11-12
      • 2020-07-09
      • 1970-01-01
      • 2019-03-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多