【问题标题】:How to parallelize file writing using TPL?如何使用 TPL 并行化文件写入?
【发布时间】:2015-06-24 07:49:15
【问题描述】:

我正在尝试将字符串列表保存到多个文件中,每个字符串在不同的文件中,并同时进行。我是这样做的:

public async Task SaveToFilesAsync(string path, List<string> list, CancellationToken ct)
{
    int count = 0;
    foreach (var str in list)
    {
        string fullPath = path + @"\" + count.ToString() + "_element.txt";
        using (var sw = File.CreateText(fullPath))
        {
            await sw.WriteLineAsync(str);
        }
        count++;

        NLog.Trace("Saved in thread: {0} to {1}", 
           Environment.CurrentManagedThreadId,
           fullPath);

        if (ct.IsCancellationRequested)
            ct.ThrowIfCancellationRequested();
    }
}

然后这样称呼它:

try
{
   var savingToFilesTask = SaveToFilesAsync(@"D:\Test", myListOfString, ct);
}
catch(OperationCanceledException)
{
   NLog.Info("Operation has been cancelled by user.");
}

但是在日志文件中我可以清楚地看到保存总是发生在同一个线程 id 中,所以没有并行性?我究竟做错了什么?如何解决?我的目标是使用所有计算机内核尽可能快地节省所有费用。

【问题讨论】:

  • 你的操作是异步运行的,但不是并行的,是有区别的。您可以并行化循环,但我认为它不会产生任何改进 - 您的操作不受 CPU 限制,而是 IO 限制...
  • 您的 CPU 不太可能成为瓶颈。 I/O 慢了几个数量级。无需为此并行化,只有在使用机械硬盘驱动器时才会减慢处理速度。
  • 您在运行 GUI 应用程序吗?
  • @YuvalItzchakov 不,这是一个控制台应用程序

标签: c# .net multithreading task-parallel-library


【解决方案1】:

本质上,您的问题是foreach 是同步的。它使用同步的IEnumerable

要解决这个问题,首先将循环体封装到一个异步函数中。

public async Task WriteToFile(
        string path,
        string str,
        int count)
{
    var fullPath = string.Format("{0}\\{1}_element.txt", path, count);
    using (var sw = File.CreateText(fullPath))
    {
        await sw.WriteLineAsync(str);
    }

    NLog.Trace("Saved in TaskID: {0} to \"{1}\"", 
       Task.CurrentId,
       fullPath);
}

然后,不是同步循环,而是将字符串序列投射到执行封装循环体的任务序列中。这本身不是异步操作,但投影不会阻塞,即没有await

然后等待他们所有的任务按照任务调度器定义的顺序完成。

public async Task SaveToFilesAsync(
        string path,
        IEnumerable<string> list,
        CancellationToken ct)
{
    await Task.WhenAll(list.Select((str, count) => WriteToFile(path, str, count));
}

没有什么要取消的,所以没有必要向下传递取消令牌。

我使用了Select 的索引重载来提供count 值。

我已将您的日志记录代码更改为使用当前的任务 ID,这避免了调度方面的任何混乱。

【讨论】:

  • 实际上,它可以证明是有益的,因为 IO 操作涉及的不仅仅是实际的 IO。它们包括框架和操作系统级别的缓冲、序列化等 CPU 绑定操作、安全检查和其他阻止驱动器全速工作的操作。即使是像复制文件这样的“纯”IO 绑定操作,多线程也可以提高速度。尝试使用 /MT 选项运行 robocopy
  • @PanagiotisKanavos,我认为该注释可能会分散答案的注意力(我已将其删除)。我不认为这个问题真的是关于这种方法在现实世界中的优点。如果我不得不提出一些不同的建议。
  • 您可能希望将 PLINQ 的 Selectlimited max degree of parallelism 一起使用,而不是普通的 LINQ 的 Select,同时处理多个文件可能会导致一次启动处理太多文件会导致性能下降。
  • @ScottChamberlain,你说得很好,如果我的场景更复杂,我可能会使用 TPL.Dataflow,我发现它非常有用。 msdn.microsoft.com/en-us/library/hh228603(v=vs.110).aspx
【解决方案2】:

如果您想进行并行处理,您必须告诉 .NET 这样做。 我认为,如果您将代码拆分为一个附加函数,那么最简单的方法之一就会变得清晰。

这个想法是将实际的单个 IO 操作拆分为一个额外的异步函数并调用这些函数而不等待它们,而是将它们存储在一个列表中并在最后等待所有它们。

我通常不写 C# 代码,所以请原谅我可能犯的任何语法错误:

public async Task SaveToFilesAsync(string path, List<string> list, CancellationToken ct)
{
    int count = 0;
    var writeOperations = new List<Task>(list.Count);
    foreach (var str in list)
    { 
        string fullPath = path + @"\" + count.ToString() + "_element.txt";
        writeOperations.add(SaveToFileAsync(fullPath, str, ct));
        count++;
        ct.ThrowIfCancellationRequested();
    }

    await Task.WhenAll(writeOperations);
}

private async Task SaveToFileAsync(string path, string line, CancellationToken ct)
{
    using (var sw = File.CreateText(path))
    {
        await sw.WriteLineAsync(line);
    }

    NLog.Trace("Saved in thread: {0} to {1}", 
        Environment.CurrentManagedThreadId, 
        fullPath);

    ct.ThrowIfCancellationRequested();
}

这样 IO 操作由同一个线程一次又一次地触发。这应该工作得非常快。一旦使用 .NET ThreadPool 完成 IO 操作,就会触发延续。

我还删除了if (ct.IsCancellationRequested) 检查,因为无论如何这是由ct.ThrowIfCancellationRequested(); 完成的。

希望这能让您了解如何处理这些事情。

【讨论】:

  • 这是正确的,但并不重要,因为每个写入操作都会写入不同的文件。所以这个顺序是不感兴趣的。编号也由外部函数维护,因此文件的编号将是正确的。
【解决方案3】:

如果这是在并行存储 (SSD) 上,您可以通过并行化来加快速度。由于没有内置方法可以以一定的并行度并行化异步循环,因此我建议使用具有固定并行度和同步 IO 的 PLINQ。 Parallel.ForEach 不能有一个固定的 DOP(只有一个最大 DOP)。

【讨论】:

    【解决方案4】:

    我已经在原始问题中添加了我的答案,我应该在此处添加吗? C# TPL calling tasks in a parallel manner and asynchronously creating new files

    编辑:这是现在并行运行多个保存的建议解决方案。

    您需要将 foreach 循环(从第一项到最后一项按顺序运行)替换为可配置为并行性的 Parallel.ForEach() 循环。

    var cts = new CancellationTokenSource();
    Task.WaitAll(SaveFilesAsync(@"C:\Some\Path", files, cts.Token));
    cts.Dispose();
    

    然后在该方法中进行并行处理。

    public async Task SaveFilesAsync(string path, List<string> list, CancellationToken token)
    {
        int counter = 0;
    
        var options = new ParallelOptions
                          {
                              CancellationToken = token,
                              MaxDegreeOfParallelism = Environment.ProcessorCount,
                              TaskScheduler = TaskScheduler.Default
                          };
    
        await Task.Run(
            () =>
                {
                    try
                    {
                        Parallel.ForEach(
                            list,
                            options,
                            (item, state) =>
                                {
                                    // if cancellation is requested, this will throw an OperationCanceledException caught outside the Parallel loop
                                    options.CancellationToken.ThrowIfCancellationRequested();
    
                                    // safely increment and get your next file number
                                    int index = Interlocked.Increment(ref counter);
                                    string fullPath = string.Format(@"{0}\{1}_element.txt", path, index);
    
                                    using (var sw = File.CreateText(fullPath))
                                    {
                                        sw.WriteLine(item);
                                    }
    
                                    Debug.Print(
                                        "Saved in thread: {0} to {1}",
                                        Thread.CurrentThread.ManagedThreadId,
                                        fullPath);
                                });
                    }
                    catch (OperationCanceledException)
                    {
                        Debug.Print("Operation Canceled");
                    }
                });
    }
    

    【讨论】:

    • 去吧,我需要使用任务,但有趣的解决方案虽然
    • 我迷路了...需要使用任务吗?我建议的解决方案与您上面的代码相同,但将您的 foreach 循环更改为同时运行多个保存的并行循环。
    • 抱歉,我查看了不同的解决方案,不是你的
    • 没问题,如果这不是您想要的。我只是不明白“我需要使用任务”是什么意思。
    猜你喜欢
    • 2015-06-15
    • 1970-01-01
    • 2019-09-17
    • 1970-01-01
    • 1970-01-01
    • 2020-02-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多