【问题标题】:Multithreading with a large number of file IO tasks具有大量文件 IO 任务的多线程
【发布时间】:2015-09-27 05:54:38
【问题描述】:

我对 C# 并不完全陌生,但我对该语言不够熟悉,不知道如何做我需要做的事情。

我有一个文件,叫它 File1.txt。 File1.txt 有 100,000 行左右。 我将复制 File1.txt 并将其命名为 File1_untested.txt。 我还将创建一个空文件“Successes.txt” 对于文件中的每一行:

  • 从 File1_untested.txt 中删除这一行
  • 如果此行通过测试,则将其写入 Successes.txt

所以,我的问题是,我怎样才能多线程呢?

到目前为止,我的方法是创建一个对象 (LineChecker),为对象指定要检查的行,然后将对象传递到 ThreadPool 中。我了解如何通过 CountdownEvent 将 ThreadPools 用于一些任务。但是,一次将 100,000 个任务全部排队似乎是不合理的。我怎样才能逐渐喂养游泳池?也许一次 1000 行或类似的东西。

另外,我需要确保没有两个线程同时添加到 Successes.txt 或从 File1_untested.txt 中删除。我可以用 lock() 处理这个,对吧?我应该将什么传递给 lock()?我可以使用 LineChecker 的静态成员吗​​?

我只是想大致了解如何设计这样的东西。

【问题讨论】:

  • 给我们一些关于你的“测试”性质的想法。它有什么作用?它是如何工作的?计算成本低还是计算成本高?
  • 测试与网络相关。我希望大部分时间都花在等待网络服务器的响应上。
  • 看看 BlockingCollection 类。 msdn.microsoft.com/en-us/library/dd997371(v=vs.110).aspx 它实现了你想要使用的消费者/生产者模式。
  • 好的,保罗,谢谢。这正是我正在寻找的响应类型!

标签: c# multithreading


【解决方案1】:

由于测试需要相当长的时间,因此使用多个 CPU 内核是有意义的。然而,这种使用应该只用于相对昂贵的测试,而不是用于读取/更新文件。这是因为读取/更新文件相对便宜。

以下是一些您可以使用的示例代码:

假设你有一个相对昂贵的测试方法:

private bool Test(string line)
{
    //This test is expensive
}

这是一个可以利用多个 CPU 进行测试的代码示例:

这里我们将集合中的项目数限制为 10,以便从文件中读取的线程将等待其他线程赶上,然后再从文件中读取更多行。

这个输入线程的读取速度将比其他线程测试的快得多,因此在最坏的情况下,我们将读取比测试线程完成的测试多 10 行。这确保我们有良好的内存消耗。

CancellationTokenSource cancellation_token_source = new CancellationTokenSource();

CancellationToken cancellation_token = cancellation_token_source.Token;

BlockingCollection<string> blocking_collection = new BlockingCollection<string>(10);

using (StreamReader reader = new StreamReader(new FileStream(filename, FileMode.Open, FileAccess.Read)))
{
    using (
        StreamWriter writer =
            new StreamWriter(new FileStream(success_filename, FileMode.OpenOrCreate, FileAccess.Write)))
    {

        var input_task = Task.Factory.StartNew(() =>
        {
            try
            {
                while (!reader.EndOfStream)
                {
                    if (cancellation_token.IsCancellationRequested)
                        return;

                    blocking_collection.Add(reader.ReadLine());
                }
            }
            finally //In all cases, even in the case of an exception, we need to make sure that we mark that we have done adding to the collection so that the Parallel.ForEach loop will exit. Note that Parallel.ForEach will not exit until we call CompleteAdding
            {
                blocking_collection.CompleteAdding();
            }
        });


        try
        {
            Parallel.ForEach(blocking_collection.GetConsumingEnumerable(), (line) =>
            {
                bool test_reault = Test(line);


                if (test_reault)
                {
                    lock (writer)
                    {
                        writer.WriteLine(line);
                    }
                }
            });
        }
        catch
        {
            cancellation_token_source.Cancel(); //If Paralle.ForEach throws an exception, we inform the input thread to stop
            throw;
        }

        input_task.Wait(); //This will make sure that exceptions thrown in the input thread will be propagated here
    }
}

【讨论】:

  • 哇,这真的很有帮助,谢谢!为什么给 BlockingCollection 的上限是 10?此外,Parallel.ForEach 会阻塞,直到它测试每一行?
  • 我更新了答案以包含对 10 值的解释。是的,Parallel.ForEach 将阻塞,直到它测试所有行。
  • 我看到你已经添加了取消令牌,这是有道理的。不错的补充
  • 是的。确保我们正确处理异常。
  • BlockingCollection 的上限不影响并行度,它只是确保例如输入线程没有读取所有 100000 行,而测试线程只完成了大约 200 行。例如,指定 10 将确保如果测试线程已完成 200 行,则读取/输入线程最多读取 210 行。所以 10 是已读取并准备测试的行数。
【解决方案2】:

如果您的“测试”速度很快,那么多线程不会给您带来任何优势,因为您的代码将 100% 与磁盘绑定,并且您可能会将所有文件都放在同一个磁盘上:您无法改进单个磁盘的多线程吞吐量。

但是由于您的“测试”将等待来自网络服务器的响应,这意味着测试将会很慢,因此多线程还有很大的改进空间。基本上,您需要的线程数取决于网络服务器可以同时处理多少请求,而不会降低网络服务器的性能。这个数字可能仍然很低,因此您最终可能一无所获,但至少您可以尝试。

如果您的文件不是很大,那么您可以一次读取所有文件,一次写入所有文件。如果每行只有 80 个字符长,那么这意味着您的文件只有 8 兆字节,这是花生,所以您可以将所有行读入一个列表,在列表上工作,生成另一个列表,最后写出整个列表。

这将允许您创建一个结构,例如MyLine,其中包含每行的索引和每行的文本,这样您就可以在编写之前对所有行进行排序,这样您就不必担心关于来自服务器的无序响应。

然后,您需要按照@Paul 的建议使用像BlockingCollection 这样的边界阻塞队列。

BlockingCollection 接受其最大容量作为构造函数参数。这意味着一旦达到其最大容量,任何进一步的添加尝试都会被阻止(调用者坐在那里等待),直到删除一些项目。因此,如果您想同时处理多达 10 个未决请求,您可以按如下方式构建它:

var sourceCollection = new BlockingCollection<MyLine>(10);

您的主线程将用MyLine 对象填充sourceCollection,并且您将有10 个线程阻塞等待从集合中读取MyLines。每个线程向服务器发送请求,等待响应,将结果保存到线程安全的resultCollection,并尝试从sourceCollection 获取下一项。

您可以使用 C# 的 async 功能代替使用多线程,但我对它们不是很熟悉,因此我无法准确地建议您如何做到这一点。

最后将resultCollection的内容复制到一个List中,对列表进行排序,写入输出文件。 (复制到单独的List 可能是个好主意,因为对线程安全的resultCollection 进行排序可能比对非线程安全的List 进行排序要慢得多。我说可能。)

【讨论】:

  • 优秀的描述迈克,谢谢。你对文件“File1_untested.txt”有什么想法吗?我的意思是,一种从其中删除线条作为测试的干净方法。 (万一程序意外中断)
  • 这很难。从文本文件中删除行是昂贵的。本质上,您必须读取整个文件并写入整个文件,跳过要删除的行。将已处理的最大行号继续存储在一个单独的文件中可能会容易得多,该文件仅包含一行,仅包含一个数字。然后,您将在每次开始处理之前查阅此文件。但是如果在一个紧密的循环中保存该文件也会很昂贵,因此请尝试每隔一百行左右更新一次。
猜你喜欢
  • 2013-10-06
  • 2018-12-12
  • 1970-01-01
  • 2020-04-09
  • 2021-07-28
  • 2014-05-06
  • 2021-11-09
  • 2020-06-16
  • 1970-01-01
相关资源
最近更新 更多