【问题标题】:Parallel.ForEach not adding items as expected in ConcurrentBag in C#Parallel.ForEach 未按预期在 C# 的 ConcurrentBag 中添加项目
【发布时间】:2020-02-21 10:00:05
【问题描述】:

在我的Asp.Net Core WebApi Controller 中,我收到了IFormFile[] files。我需要将其转换为List<DocumentData>。我第一次使用foreach。它工作正常。但后来决定更改为Parallel.ForEach,因为我收到了很多(> 5)个文件。

这是我的DocumentData 班级:

public class DocumentData
{
    public byte[] BinaryData { get; set; }
    public string FileName { get; set; }
}

这是我的Parallel.ForEach逻辑:

var documents = new ConcurrentBag<DocumentData>();
Parallel.ForEach(files, async (currentFile) =>
{
    if (currentFile.Length > 0)
    {
        using (var ms = new MemoryStream())
        {
            await currentFile.CopyToAsync(ms);
            documents.Add(new DocumentData
            {
                BinaryData = ms.ToArray(),
                FileName = currentFile.FileName
            });
        }
    }
});

例如,即使是两个文件作为输入,documents 也总是给出一个文件作为输出。我错过了什么吗?

我最初拥有List&lt;DocumentData&gt;。我发现它不是线程安全的并更改为ConcurrentBag&lt;DocumentData&gt;。但我仍然得到了意想不到的结果。请帮忙看看我哪里错了?

【问题讨论】:

    标签: c# multithreading concurrency asp.net-core-webapi parallel.foreach


    【解决方案1】:

    我猜是因为Parallel.Foreach 不支持async/await。它只需要Action 作为输入并为每个项目执行它。在异步委托的情况下,它将以一种即发即弃的方式执行它们。 在这种情况下,传递的 lambda 将被视为 async void 函数并且无法等待 async void

    如果有需要 Func&lt;Task&gt; 的过载,那么它会起作用。

    我建议你在Select的帮助下创建Tasks,同时使用Task.WhenAll来执行它们。

    例如:

    var tasks = files.Select(async currentFile =>
    {
        if (currentFile.Length > 0)
        {
            using (var ms = new MemoryStream())
            {
                await currentFile.CopyToAsync(ms);
                documents.Add(new DocumentData
                {
                    BinaryData = ms.ToArray(),
                    FileName = currentFile.FileName
                });
            }
        }
    });
    
    await Task.WhenAll(tasks);
    

    此外,您可以通过从该方法返回 DocumentData 实例来改进该代码,在这种情况下,无需修改 documents 集合。 Task.WhenAll 具有将 IEnumerable&lt;Task&lt;TResult&gt; 作为输入并生成 TResult 数组的 Task 的重载。所以,结果会是这样:

    var tasks = files.Select(async currentFile =>
        {
            if (currentFile.Length > 0)
            {
                using (var ms = new MemoryStream())
                {
                    await currentFile.CopyToAsync(ms);
                    return new DocumentData
                    {
                        BinaryData = ms.ToArray(),
                        FileName = currentFile.FileName
                    };
                }
            }
    
            return null;
        });
    
    var documents =  (await Task.WhenAll(tasks)).Where(d => d != null).ToArray();
    

    【讨论】:

    • 如果你使用Task.WhenAll,那么你需要await它。
    • @FarhadJabiyev,现在不需要ConcurrentBag&lt;DocumentData&gt; 对吧?相反,我可以使用List&lt;DocumentData&gt;
    • @fingers10 当然可以,因为Parallel.ForEach 不支持异步 lambda。您可以使用简单的foreach 来顺序执行它们,或者使用Task.WhenAll 来同时执行它们
    • 作为奖励,您可以从Select lambda 返回 DocumentData,然后await Task.WhenAll(..) 将产生DocumentData[]
    • @StephenCleary 好主意。值得用这个想法更新答案。会的。
    【解决方案2】:

    您对 并发集合 的想法是正确的,但误用了 TPL 方法

    简而言之,您需要非常小心异步 lambdas,如果您将它们传递给 ActionFunc&lt;Task&gt;

    您的问题是因为Parallel.For / ForEach 不适合异步和等待模式IO 绑定任务。它们适用于 cpu 绑定的工作负载。这意味着它们本质上具有Action 参数,让任务调度程序为您创建任务

    如果您想同时运行多个任务,请使用 Task.WhenAllTPL 数据流 ActionBlock,它可以有效地处理 CPU 限制 和 IO bound 工作负载,或者更直接地说,它们可以处理 tasks,这就是异步方法。

    根本问题是当您在 Action 上调用 async lambda 时,实际上是在创建一个 async void 方法,该方法将作为未观察到的 task 运行.也就是说,您的TPL 方法 只是并行创建一堆任务 来运行一堆未观察到的任务,而不是等待它们。

    这样想,你请一群朋友去给你买些杂货,然后他们又告诉别人去买你的杂货,但你的朋友向你报告说他们的工作已经完成。显然不是,而且你没有杂货。

    【讨论】:

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