【问题标题】:ConcurrentQueue<string> with multithreading/tasks?ConcurrentQueue<string> 与多线程/任务?
【发布时间】:2019-06-19 15:03:19
【问题描述】:

如何将多线程添加到并发队列中,目前我正在使用 1 个线程上的并发队列处理文本文件,但是如果我想在多个线程上运行它以减少整体处理时间怎么办?

当前方法示例 -

    private static ConcurrentQueue<string> queue;

    static void Main(string[] args)
    {

        queue = new ConcurrentQueue<string>(System.IO.File.ReadAllLines("input.txt"));
        Process();

    }

    static void Process()
    {

        while (queue.Count > 0)
        {
            string entry;
            if (queue.TryDequeue(out entry))
            {
                Console.WriteLine(entry);
                log("out.txt", entry);
            }
        }
    }

    private static void log(string file, string data)
    {
        using (StreamWriter writer = System.IO.File.AppendText(file))
        {
            writer.WriteLine(data);
            writer.Flush();
            writer.Close();
        }
    }

代码分解-

queue = new ConcurrentQueue<string>(System..) // assigns queue to a text file

Process(); // Executes the Process method


static void Process() {

    while ... // runs a loop whilst queue.count is not equal to 0

    if (queueTryDequeue... // takes one line from queue and assigns it to 'string entry'

    Console.. // Writes 'entry' to console

    log.. // adds 'string entry' to a new line inside 'out.txt'

input.txt 例如包含 1000 个条目,我想创建 10 个线程,它们从 input.txt 中获取一个条目并对其进行处理,同时避免使用与另一个线程相同的条目/复制相同的进程。我将如何实现这个?

【问题讨论】:

  • 很大程度上取决于这里的和处理的含义。 CPU 或 I/O 绑定?

标签: c# multithreading concurrency


【解决方案1】:

您应该使用Parallel 循环:

注意:它不会按原始顺序循环项目!

private static StreamWriter logger;

static void Main(string[] args)
{
    // Store your entries from a file in a queue.
    ConcurrentQueue<string> queue = new ConcurrentQueue<string>(System.IO.File.ReadAllLines("input.txt"));

    // Open StreamWriter here.
    logger = File.AppendText("log.txt");

    // Call process method.
    ProcessParallel(queue);

    // Close the StreamWriter after processing is done.
    logger.Close();
}

static void ProcessParallel(ConcurrentQueue<string> collection)
{
    ParallelOptions options = new ParallelOptions()
    {
        // A max of 10 threads can access the file at one time.
        MaxDegreeOfParallelism = 10
    };

    // Start the loop and store the result, so we can check if all the threads are done.
    // The Parallel.For will do all the mutlithreading for you!
    ParallelLoopResult result = Parallel.For(0, collection.Count, options, (i) =>
    {
        string entry;
        if (collection.TryDequeue(out entry))
        {
            Console.WriteLine(entry);
            log(entry);
        }
    });
    // Parallel.ForEach can also be used.

    // Block the main thread while it is still processing the entries...
    while (!result.IsCompleted) ;

    // Every thread is done
    Console.WriteLine("Multithreaded loop is done!");
}

private static void log(string data)
{
    if (logger.BaseStream == null)
    {
        // Cannot log, because logger.Close(); was called.
        return;
    }

    logger.WriteLine(data);
}

【讨论】:

  • 这似乎有效,但是如果我尝试使用我的日志方法而不是您提供的方法,它会引发异常'进程无法访问此文件,因为它正在被另一个进程使用' ..除非我采用您的方法并将主要内部的记录器设置为“log.txt”
  • 我已标记为正确答案,因为这个问题是关于线程的,您的答案是正常的。
  • @R2-D2 您的日志记录方法不起作用,因为循环同时记录多个线程。日志函数打开“log.txt”文件,写入一些内容并关闭它。如果另一个线程在另一个线程仍在写入日志文件时访问日志方法,您将收到该异常,必须首先通过调用 writer.Close(); 删除文件上的“锁定”。
猜你喜欢
  • 1970-01-01
  • 2016-08-29
  • 1970-01-01
  • 2010-12-18
  • 2014-04-25
  • 2015-09-14
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多