消耗多个任务/消费者的阻塞集合

Dar*_*der 2 c# multithreading task-parallel-library

我有以下代码,我从一个源填充用户,为例子如下所示.我想要做的是消耗多个消费者的BlockingCollection.

低于正确的方式吗?还有什么是最好的线程数?好吧这取决于硬件,内存等.或者我怎么能以更好的方式做到这一点?

还将在下面的实现确保我将处理集合中的所有内容,直到它为空?

    class Program
    {
        public static readonly BlockingCollection<User> users = new BlockingCollection<User>();

        static void Main(string[] args)
        {
            for (int i = 0; i < 100000; i++)
            {
                var u = new User {Id = i, Name = "user " + i};
                users.Add(u);
            }

            Run(); 
        }

        static void Run()
        {
            for (int i = 0; i < 100; i++)
            {
                Task.Factory.StartNew(Process, TaskCreationOptions.LongRunning);
            }
        }

        static void Process()
        {
            foreach (var user in users.GetConsumingEnumerable())
            {
                Console.WriteLine(user.Id);
            }
        }
    }

    public class User
    {
        public int Id { get; set; }
        public string Name { get; set; }
    }
Run Code Online (Sandbox Code Playgroud)

Sco*_*ain 7

一些小事情

  1. 您从未调用过CompleteAdding,因为没有这样做,您的消费foreach循环永远不会完成并永久挂起.通过users.CompleteAdding()在初始for循环之后执行修复.
  2. 你永远不会等待工作完成,Run()将启动你的100个线程(除非你的真实过程涉及大量等待无争议的资源,否则可能过多).由于任务不是前台线程,因此当您Main退出时,它们不会使您的程序保持打开状态.您需要一个CountdownEvent来跟踪所有内容的完成情况.
  3. 在生产者完成所有工作之后,你不会启动你的消费者,你应该将生产者分离到一个单独的线程或首先启动消费者,这样他们就可以在你在主线程上填充生产者时工作了.

这是带有修复程序的代码的更新版本

class Program
{
    private const int MaxThreads = 100; //way to high for this example.
    private static readonly CountdownEvent cde = new CountdownEvent(MaxThreads);
    public static readonly BlockingCollection<User> users = new BlockingCollection<User>();

    static void Main(string[] args)
    {
        Run(); 

        for (int i = 0; i < 100000; i++)
        {
            var u = new User {Id = i, Name = "user " + i};
            users.Add(u);
        }
        users.CompleteAdding();
        cde.Wait();
    }

    static void Run()
    {
        for (int i = 0; i < MaxThreads; i++)
        {
            Task.Factory.StartNew(Process, TaskCreationOptions.LongRunning);
        }
    }

    static void Process()
    {
        foreach (var user in users.GetConsumingEnumerable())
        {
            Console.WriteLine(user.Id);
        }
        cde.Signal();
    }
}

public class User
{
    public int Id { get; set; }
    public string Name { get; set; }
}
Run Code Online (Sandbox Code Playgroud)

对于我之前说过的"最佳线程数",这实际上取决于你在等什么.

如果您正在处理的是CPU绑定,则最佳线程数可能是Enviorment.ProcessorCount.

如果您正在做的是等待外部资源,但新请求不会影响旧请求(例如,询问20个不同的服务器以获取信息,服务器上的服务器n上的负载不会影响服务器上的负载n+1)在这种情况下我会让并行.ForEach只为你选择线程数.

如果您正在等待争用的资源(例如读/写硬盘),您将不想使用很多线程(甚至可能只使用一个).我刚刚在另一个问题中发布了一个答案,当从硬盘读入时,你应该一次只使用一个线程,这样硬盘驱动器就不会一遍遍地试图完成所有的读取.