【问题标题】:Run async method 8 times in parallel并行运行异步方法 8 次
【发布时间】:2013-02-03 14:51:54
【问题描述】:

如何将以下内容变成 Parallel.ForEach?

public async void getThreadContents(String[] threads)
{
    HttpClient client = new HttpClient();
    List<String> usernames = new List<String>();
    int i = 0;

    foreach (String url in threads)
    {
        i++;
        progressLabel.Text = "Scanning thread " + i.ToString() + "/" + threads.Count<String>();
        HttpResponseMessage response = await client.GetAsync(url);
        String content = await response.Content.ReadAsStringAsync();
        String user;
        Predicate<String> userPredicate;
        foreach (Match match in regex.Matches(content))
        {
            user = match.Groups[1].ToString();
            userPredicate = (String x) => x == user;
            if (usernames.Find(userPredicate) != user)
            {
                usernames.Add(match.Groups[1].ToString());
            }
        }
        progressBar1.PerformStep();
    }
}

我在假设异步和并行处理相同的情况下对其进行编码,但我刚刚意识到事实并非如此。我查看了我能找到的所有问题,但我似乎真的找不到适合我的例子。他们中的大多数缺乏可读的变量名。使用不解释它们包含的内容的单字母变量名称是陈述示例的可怕方式。

我通常在名为线程的数组中有 300 到 2000 个条目(包含论坛线程的 URL),并且似乎并行处理(由于许多 HTTP 请求)会加快执行速度。

在我可以使用 Parallel.ForEach 之前,我是否必须删除所有异步(我在 foreach 之外没有任何异步,只有变量定义)?我该怎么做呢?我可以在不阻塞主线程的情况下这样做吗?

顺便说一下,我使用的是 .NET 4.5。

【问题讨论】:

  • 嗯。你的“线程”数组中有什么?为什么叫线程?它不包含线程。
  • @spender 线程与论坛线程一样,它包含论坛线程的 URL。
  • 我不明白抱怨您找到的一些答案对您的问题有何影响。
  • @svick 我写这个是因为我想要一个有意义的变量名示例。
  • 这能回答你的问题吗? Parallel foreach with asynchronous lambda

标签: c# .net parallel-processing .net-4.5


【解决方案1】:

我假设异步和并行处理是相同的

异步处理和并行处理是完全不同的。如果您不了解其中的区别,我认为您应该先阅读更多相关信息(例如what is the relation between Asynchronous and parallel programming in c#?)。

现在,您想要做的实际上并不那么简单,因为您想要异步处理一个大集合,并具有特定程度的并行性 (8)。对于同步处理,您可以使用Parallel.ForEach()(与ParallelOptions 一起配置并行度),但没有简单的替代方法可以与async 一起使用。

在您的代码中,这很复杂,因为您希望所有内容都在 UI 线程上执行。 (尽管理想情况下,您不应该直接从计算中访问 UI。相反,您应该使用 IProgress,这意味着代码不再需要在 UI 线程上执行。)

在 .Net 4.5 中执行此操作的最佳方法可能是使用 TPL 数据流。它的ActionBlock 完全符合您的要求,但它可能非常冗长(因为它比您需要的更灵活)。所以创建一个辅助方法是有意义的:

public static Task AsyncParallelForEach<T>(
    IEnumerable<T> source, Func<T, Task> body,
    int maxDegreeOfParallelism = DataflowBlockOptions.Unbounded,
    TaskScheduler scheduler = null)
{
    var options = new ExecutionDataflowBlockOptions
    {
        MaxDegreeOfParallelism = maxDegreeOfParallelism
    };
    if (scheduler != null)
        options.TaskScheduler = scheduler;

    var block = new ActionBlock<T>(body, options);

    foreach (var item in source)
        block.Post(item);

    block.Complete();
    return block.Completion;
}

在你的情况下,你会这样使用它:

await AsyncParallelForEach(
    threads, async url => await DownloadUrl(url), 8,
    TaskScheduler.FromCurrentSynchronizationContext());

这里,DownloadUrl() 是处理单个 URL(循环体)的 async Task 方法,8 是并行度(在实际代码中可能不应该是字面常量)和 @ 987654335@ 确保代码在 UI 线程上执行。

【讨论】:

  • 我在哪个命名空间中找到DataflowBlockOptionsExecutionDataflowBlockOptionsActionBlock&lt;T&gt;?我查了一下,MSDN 说 System.Threading.Tasks.Dataflow,但我不能使用那个,因为它说命名空间不存在
  • 我确实从网站上下载了它,但是当我从 NuGet 管理器那里得到它时它就起作用了。无论如何,感谢您的帮助:)
  • 我试图改编你在here 的出色工作中所做的一些事情,但是在两个嵌套的等待 foreach 和IAsyncEnumerable 中,但我得到了奇怪的结果。我错过了什么吗?也许我不应该在两个等待的AsyncParallelForEach 中使用相同的TaskScheduler.FromCurrentSynchronizationContext()?!你怎么看?
【解决方案2】:

Stephen Toub 有一个good blog post on implementing a ForEachAsync。 Svick 的回答对于 Dataflow 可用的平台来说相当不错。

这是一个替代方案,使用来自 TPL 的分区器:

public static Task ForEachAsync<T>(this IEnumerable<T> source,
    int degreeOfParallelism, Func<T, Task> body)
{
  var partitions = Partitioner.Create(source).GetPartitions(degreeOfParallelism);
  var tasks = partitions.Select(async partition =>
  {
    using (partition) 
      while (partition.MoveNext()) 
        await body(partition.Current); 
  });
  return Task.WhenAll(tasks);
}

然后你可以这样使用它:

public async Task getThreadContentsAsync(String[] threads)
{
  HttpClient client = new HttpClient();
  ConcurrentDictionary<String, object> usernames = new ConcurrentDictionary<String, object>();

  await threads.ForEachAsync(8, async url =>
  {
    HttpResponseMessage response = await client.GetAsync(url);
    String content = await response.Content.ReadAsStringAsync();
    String user;
    foreach (Match match in regex.Matches(content))
    {
      user = match.Groups[1].ToString();
      usernames.TryAdd(user, null);
    }
    progressBar1.PerformStep();
  });
}

【讨论】:

  • 在 .Net VNext 中拥有 Parallel.ForeachAsync 功能会很好。机会有多大。
  • @rudimenter:我怀疑它在名单上并不高。大多数时候,您要么需要进行并行异步工作。对于那些极其罕见的情况,当您需要并行完成异步工作时,Task.RunTask.WhenAll 为您提供了基本的并行性,除非您需要节流,否则该并行性可以正常工作。在这种情况下,TPL Dataflow 内置了节流支持,如果您的用例真的那么复杂,那么无论如何您都应该使用 Dataflow。所以我看不到 ForEachAsync 的好用例(即,通常当人们要求这样做时,已经有更好的选择了)。
  • 我同意你的观点,大部分工作要么是并行的,要么是异步的,但新的异步功能的特点是它具有病毒性。整个应用程序现在都在使用从根到分支的异步。迭代异步任务(并行或顺序)将变得非常普遍。在 BCL 中拥有一些东西而不需要像 Dataflow 这样的额外库并且没有自定义油门实现会很好。无论如何,谢谢。
  • @rudimenter:我的意思是你已经有并行和顺序迭代:Task.WhenAllawait。极少数情况是当您有大量异步 CPU 密集型任务时,您需要限制它们。无需使用 TPL 数据流即可干净地处理其他所有情况。
  • @AndrewHanlon:调度是关键,内置的ConcurrentExclusiveSchedulerPair 就足够了。只需指定 8 作为最大并发级别,并使用 ConcurrentScheduler 属性(忽略另一个)。请记住,您需要 Unwrap 由任务调度程序执行的异步代码。
【解决方案3】:

另一种选择是使用SemaphoreSlimAsyncSemaphore(其中is included in my AsyncEx library 并支持比SemaphoreSlim 更多的平台):

public async Task getThreadContentsAsync(String[] threads)
{
  SemaphoreSlim semaphore = new SemaphoreSlim(8);
  HttpClient client = new HttpClient();
  ConcurrentDictionary<String, object> usernames = new ConcurrentDictionary<String, object>();

  await Task.WhenAll(threads.Select(async url =>
  {
    await semaphore.WaitAsync();
    try
    {
      HttpResponseMessage response = await client.GetAsync(url);
      String content = await response.Content.ReadAsStringAsync();
      String user;
      foreach (Match match in regex.Matches(content))
      {
        user = match.Groups[1].ToString();
        usernames.TryAdd(user, null);
      }
      progressBar1.PerformStep();
    }
    finally
    {
      semaphore.Release();
    }
  }));
}

【讨论】:

    【解决方案4】:

    你可以试试AsyncEnumerator NuGet PackageParallelForEachAsync扩展方法:

    using System.Collections.Async;
    
    public async void getThreadContents(String[] threads)
    {
        HttpClient client = new HttpClient();
        List<String> usernames = new List<String>();
        int i = 0;
    
        await threads.ParallelForEachAsync(async url =>
        {
            i++;
            progressLabel.Text = "Scanning thread " + i.ToString() + "/" + threads.Count<String>();
            HttpResponseMessage response = await client.GetAsync(url);
            String content = await response.Content.ReadAsStringAsync();
            String user;
            Predicate<String> userPredicate;
            foreach (Match match in regex.Matches(content))
            {
                user = match.Groups[1].ToString();
                userPredicate = (String x) => x == user;
                if (usernames.Find(userPredicate) != user)
                {
                    usernames.Add(match.Groups[1].ToString());
                }
            }
    
            // THIS CALL MUST BE THREAD-SAFE!
            progressBar1.PerformStep();
        },
        maxDegreeOfParallelism: 8);
    }
    

    【讨论】:

      猜你喜欢
      • 2016-12-02
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-04-15
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多