【问题标题】:Read and process files in parallel C#并行 C# 读取和处理文件
【发布时间】:2014-01-05 00:53:21
【问题描述】:

我有非常大的文件需要阅读和处理。这可以使用线程并行完成吗?

这是我完成的一些代码。但是,一个接一个地读取和处理文件似乎并没有缩短执行时间。

String[] files = openFileDialog1.FileNames;

Parallel.ForEach(files, f =>
{
    readTraceFile(f);
});        

private void readTraceFile(String file)
{
    StreamReader reader = new StreamReader(file);
    String line;

    while ((line = reader.ReadLine()) != null)
    {
        String pattern = "\\s{4,}";

        foreach (String trace in Regex.Split(line, pattern))
        {
            if (trace != String.Empty)
            {
                String[] details = Regex.Split(trace, "\\s+");

                Instruction instruction = new Instruction(details[0],
                    int.Parse(details[1]),
                    int.Parse(details[2]));
                Console.WriteLine("computing...");
                instructions.Add(instruction);
            }
        }
    }
}

【问题讨论】:

  • 你是 CPU 密集型还是 IO 密集型?
  • instructions 线程安全吗? (答案:没有)
  • IO 系统不如你的 CPU 快,所以当涉及到 IO 时使用多线程没有任何好处也就不足为奇了。
  • @user2936347 它可以.. 但前提是您在它在内存中(受 CPU 限制)时对其执行计算。如果您正在等待将大文件加载到内存中(I/O-Bound).. 那么没有。
  • @user2936347,这就是要走的路。您有一个用于 IO 的专用线程(或异步Task),另一个线程或Task 在内容可用时处理它(如果您的处理成为瓶颈,可能并行处理)。经典的生产者-消费者。请参阅并行编程模式 (microsoft.com/en-au/download/details.aspx?id=19222)。第 55 页几乎是完全您的方案。

标签: c# multithreading


【解决方案1】:

您的应用程序的性能似乎主要受 IO 限制。但是,您的代码中仍有一些 CPU 密集型工作。这两项工作是相互依赖的:在 IO 完成其工作之前,您的 CPU 密集型工作无法开始,并且在您的 CPU 完成前一个工作项之前,I​​O 不会继续执行下一个工作项。他们都在互相扶持。因此,可能(在最底部解释)如果您并行执行 IO 和 CPU 密集型工作,您将看到吞吐量有所提高,如下所示:

void ReadAndProcessFiles(string[] filePaths)
{
    // Our thread-safe collection used for the handover.
    var lines = new BlockingCollection<string>();

    // Build the pipeline.
    var stage1 = Task.Run(() =>
    {
        try
        {
            foreach (var filePath in filePaths)
            {
                using (var reader = new StreamReader(filePath))
                {
                    string line;

                    while ((line = reader.ReadLine()) != null)
                    {
                        // Hand over to stage 2 and continue reading.
                        lines.Add(line);
                    }
                }
            }
        }
        finally
        {
            lines.CompleteAdding();
        }
    });

    var stage2 = Task.Run(() =>
    {
        // Process lines on a ThreadPool thread
        // as soon as they become available.
        foreach (var line in lines.GetConsumingEnumerable())
        {
            String pattern = "\\s{4,}";

            foreach (String trace in Regex.Split(line, pattern))
            {
                if (trace != String.Empty)
                {
                    String[] details = Regex.Split(trace, "\\s+");

                    Instruction instruction = new Instruction(details[0],
                        int.Parse(details[1]),
                        int.Parse(details[2]));
                    Console.WriteLine("computing...");
                    instructions.Add(instruction);
                }
            }
        }
    });

    // Block until both tasks have completed.
    // This makes this method prone to deadlocking.
    // Consider using 'await Task.WhenAll' instead.
    Task.WaitAll(stage1, stage2);
}

我非常怀疑是您的 CPU 工作造成了问题,但如果发生这种情况,您也可以像这样并行化第 2 阶段:

    var stage2 = Task.Run(() =>
    {
        var parallelOptions = new ParallelOptions { MaxDegreeOfParallelism = Environment.ProcessorCount };

        Parallel.ForEach(lines.GetConsumingEnumerable(), parallelOptions, line =>
        {
            String pattern = "\\s{4,}";

            foreach (String trace in Regex.Split(line, pattern))
            {
                if (trace != String.Empty)
                {
                    String[] details = Regex.Split(trace, "\\s+");

                    Instruction instruction = new Instruction(details[0],
                        int.Parse(details[1]),
                        int.Parse(details[2]));
                    Console.WriteLine("computing...");
                    instructions.Add(instruction);
                }
            }
        });
    });

请注意,如果您的 CPU 工作组件与 IO 组件相比可以忽略不计,那么您将不会看到太多的加速。与顺序处理相比,工作负载越均匀,管道的性能就越好。

既然我们谈论的是性能,请注意我对上述代码中阻塞调用的数量并不特别兴奋。如果我在自己的项目中这样做,我会选择 async/await 路线。在这种情况下我选择不这样做,因为我想让事情易于理解和易于集成。

【讨论】:

  • 很好的示例代码。关于“阻塞呼叫数”。更改最后一行来等待 Task.WhenAll 还不够吗?当然,将方法更改为异步任务。
【解决方案2】:

从您尝试执行的操作来看,您几乎可以肯定是受 I/O 限制的。在这种情况下尝试并行处理将无济于事,实际上可能会由于磁盘驱动器上的附加寻道操作而减慢处理速度(除非您可以将数据拆分到多个主轴上)。

【讨论】:

  • 如果我受 I/O 限制,有什么办法可以提高性能吗?
  • @user2936347 通常执行许多异步调用对 I/O 来说更好。看看新的async-await 模式
  • @user2936347:有一些策略可以帮助解决 I/O 问题。然而,大多数都需要对硬件进行投资。这是否意味着单个更快的驱动器(如 SSD)、RAID 0 或 1,甚至只是将文件拆分到多个驱动器上,每个驱动器都有自己的独立控制器或它们的某种组合。
【解决方案3】:

尝试改为并行处理这些行。例如:

var q = from file in files
        from line in File.ReadLines(file).AsParallel()    // for smaller files File.ReadAllLines(file).AsParallel() might be faster
        from trace in line.Split(new [] {"    "}, StringSplitOptions.RemoveEmptyEntries)  // split by 4 spaces and no need for trace != "" check
        let details = trace.Split(null as char[], StringSplitOptions.RemoveEmptyEntries)  // like Regex.Split(trace, "\\s+") but removes empty strings too
        select new Instruction(details[0], int.Parse(details[1]), int.Parse(details[2]));

List<Instruction> instructions = q.ToList();  // all of the file reads and work is done here with .ToList

随机访问非 SSD 硬盘驱动器(当您尝试同时读取/写入不同文件或碎片文件时)通常比顺序访问(例如读取单个碎片整理文件)慢得多,所以我希望并行处理单个文件以加快碎片整理文件的速度。

此外,跨线程共享资源(例如 Console.Write 或添加到线程安全的阻塞集合)可能会减慢或阻塞/死锁执行,因为某些线程必须等待其他线程完成访问该资源。

【讨论】:

  • 谢谢,但这是一个两年前的话题,我在学校任务中​​需要它:)
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2012-11-30
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多