【问题标题】:How to write to an object after TransformBlock?如何在 TransformBlock 之后写入对象?
【发布时间】:2018-05-24 00:03:25
【问题描述】:

我有一个需要并行迭代的对象列表。这是我需要做的:

foreach (var r in results)
{
    r.SomeList = await apiHelper.Get(r.Id);
}

因为我想并行化它,所以我尝试使用 Parallel.ForEach() 但它不会等到一切都真正完成后 apiHelper.Get() 正在自己内部进行等待。

Parallel.ForEach(
                results,
                async (r) =>
                {
                    r.SomeList = await apiHelper.Get(r.Id);
                });

于是我在网上搜索,发现了这个: Nesting await in Parallel.ForEach

现在我对 TPL 很陌生(20 分钟前),我可能会遗漏一些明显的东西。我该如何继续?

        var getBlock = new TransformBlock<string, List<Something>>(
            async i =>
            {
                var c = await apiHelper.Get(i);
                return c;
            }, new ExecutionDataflowBlockOptions
            {
                MaxDegreeOfParallelism = DataflowBlockOptions.Unbounded
            });

        foreach (var r in results)
        {
            r.SomeList = getBlock.Post(r.Id);  // ERROR: Can't convert boolean to list.
        }

        getBlock.Complete();

【问题讨论】:

    标签: c# async-await task-parallel-library parallel.foreach tpl-dataflow


    【解决方案1】:

    或许可以考虑改用微软的响应式框架。

    代码如下:

    var query =
        from r in results.ToObservable()
        from l in Observable.FromAsync(() => apiHelper.Get(r.Id))
        select new { r, l };
    
    query
        .Subscribe(x => x.r.SomeList = x.l);
    

    完成。并行和异步。

    只需 NuGet "System.Reactive" 并添加 using System.Reactiive.Linq;

    【讨论】:

    • 我一直想向我学习一些反应性的东西,但每次我访问该网站时,页面都会让我失望,+1 让我学到一些东西
    【解决方案2】:

    这是一个在TPL dataflow 库中使用ActionBlock Class 的示例。

    它基本上给你并行,async 并且它相当容易理解

    数据流示例

    public static async Task DoWorkLoads(List<Something> results)
    {
       var options = new ExecutionDataflowBlockOptions
                         {
                            MaxDegreeOfParallelism = 50
                         };
    
       var block = new ActionBlock<Something>(MyMethodAsync, options);
    
       foreach (var result in results)
          block.Post(result );
    
       block.Complete();
       await block.Completion;
    
    }
    
    ...
    
    public async Task MyMethodAsync(Something result)
    {       
       result.SomeList = await apiHelper.Get(result.Id);
    }
    

    显然,你需要检查错误,并添加胡椒和盐调味

    另外,它假设apiHelper 是线程安全的

    【讨论】:

      【解决方案3】:

      要异步并行调用 api,您不需要 Reactive 或 Dataflow。您所拥有的唯一复杂情况是您正在通过使用 api 调用的结果设置它的属性来改变对象r。尽管如此,你想要的还是很简单的:

      这个:

      foreach (var r in results)
      {
          r.SomeList = await apiHelper.Get(r.Id);
      }
      

      变成:

      var tasks = results.Select(async r => { r.SomeList = await apiHelper.Get(r.Id); });
      await Task.WhenAll(tasks);
      

      假设apiHelper.Get 实际上是非阻塞和异步的,那么results 中的每个项目都会对api 进行异步和并行调用。

      【讨论】:

        猜你喜欢
        • 2019-01-16
        • 1970-01-01
        • 1970-01-01
        • 2020-06-23
        • 2013-08-25
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2018-10-01
        相关资源
        最近更新 更多