具有多线程的ConcurrentQueue

Nav*_*mar 1 c# multithreading concurrent-queue

我是多线程概念的新手.我需要将一定数量的字符串添加到队列中并使用多个线程处理它们.使用ConcurrentQueue哪个是线程安全的.

这就是我尝试过的.但是不会处理添加到并发队列中的所有项目.只处理前4个项目.

class Program
{
    ConcurrentQueue<string> iQ = new ConcurrentQueue<string>();
    static void Main(string[] args)
    {
        new Program().run();
    }

    void run()
    {
        int threadCount = 4;
        Task[] workers = new Task[threadCount];

        for (int i = 0; i < threadCount; ++i)
        {
            int workerId = i;
            Task task = new Task(() => worker(workerId));
            workers[i] = task;
            task.Start();
        }

        for (int i = 0; i < 100; i++)
        {
            iQ.Enqueue("Item" + i);
        }

        Task.WaitAll(workers);
        Console.WriteLine("Done.");

        Console.ReadLine();
    }

    void worker(int workerId)
    {
        Console.WriteLine("Worker {0} is starting.", workerId);
        string op;
        if(iQ.TryDequeue(out op))
        {
            Console.WriteLine("Worker {0} is processing item {1}", workerId, op);
        }

        Console.WriteLine("Worker {0} is stopping.", workerId);
    }


}
Run Code Online (Sandbox Code Playgroud)

Sef*_*efe 5

您的实施存在一些问题.第一个也是显而易见的一个是该worker方法只将零个或一个项目出列然后停止:

    if(iQ.TryDequeue(out op))
    {
        Console.WriteLine("Worker {0} is processing item {1}", workerId, op);
    }
Run Code Online (Sandbox Code Playgroud)

它应该是:

    while(iQ.TryDequeue(out op))
    {
        Console.WriteLine("Worker {0} is processing item {1}", workerId, op);
    }
Run Code Online (Sandbox Code Playgroud)

但这并不足以使您的程序正常工作.如果你的工人出列的速度比主线程入队的速度快,那么当主要任务仍然排队时它们会停止.你需要告诉工人他们可以停下来.您可以定义一个布尔变量,该变量将被设置为true一旦排队完成:

for (int i = 0; i < 100; i++)
{
    iQ.Enqueue("Item" + i);
}
Volatile.Write(ref doneEnqueueing, true);
Run Code Online (Sandbox Code Playgroud)

工人将检查价值:

void worker(int workerId)
{
    Console.WriteLine("Worker {0} is starting.", workerId);
    do {
        string op;
        while(iQ.TryDequeue(out op))
        {
            Console.WriteLine("Worker {0} is processing item {1}", workerId, op);
        }
        SpinWait.SpinUntil(() => Volatile.Read(ref doneEnqueueing) || (iQ.Count > 0));
    }
    while (!Volatile.Read(ref doneEnqueueing) || (iQ.Count > 0))
    Console.WriteLine("Worker {0} is stopping.", workerId);
}  
Run Code Online (Sandbox Code Playgroud)