【问题标题】:C# Lock Threading IssueC#锁线程问题
【发布时间】:2015-01-18 07:07:44
【问题描述】:

这里的任务很简单(或者我是这么认为的......)。我想用要执行的方法填充一个队列(所有这些都将返回一个对象结果),然后我想让一些任意数量的线程从这个队列中拉出,执行这些方法,并将结果添加到其他集合(在这种情况下为字典)将在所有工作完成后返回。将在主线程中调用一个 main 方法,该方法将开始处理并应该阻塞,直到所有线程完成它们正在做的任何事情并返回带有结果的集合。所以我把这个类放在一起:

public class BackgroundWorkManager
{
    public delegate object ThreadTask();

    private Thread[] workers;
    private ManualResetEvent workerThreadMre;
    private ManualResetEvent mainThreadMre;
    private Queue<WorkItem> workQueue;
    private Dictionary<string, object> results;
    private object writeLock;
    private int activeTasks;

    private struct WorkItem
    {
        public string name;
        public ThreadTask task;

        public WorkItem(string name, ThreadTask task)
        {
            this.name = name;
            this.task = task;
        }
    }

    private void workMethod()
    {
        while (true)
        {
            workerThreadMre.WaitOne();

            WorkItem task;

            lock (workQueue)
            {
                if (workQueue.Count == 0)
                {
                    workerThreadMre.Reset();
                    continue;
                }

                task = workQueue.Dequeue();
            }

            object result = task.task();

            lock (writeLock)
            {
                results.Add(task.name, result);
                activeTasks--;

                if (activeTasks == 0)
                    mainThreadMre.Set();
            }
        }
    }

    public BackgroundWorkManager()
    {
        workers = new Thread[Environment.ProcessorCount];
        workerThreadMre = new ManualResetEvent(false);
        mainThreadMre = new ManualResetEvent(false);
        workQueue = new Queue<WorkItem>();
        writeLock = new object();
        activeTasks = 0;

        for (int i = 0; i < Environment.ProcessorCount; i++)
        {
            workers[i] = new Thread(workMethod);
            workers[i].Priority = ThreadPriority.Highest;
            workers[i].Start();
        }
    }

    public void addTask(string name, ThreadTask task)
    {
        workQueue.Enqueue(new WorkItem(name, task));
    }

    public Dictionary<string, object> process()
    {
        results = new Dictionary<string, object>();

        activeTasks = workQueue.Count;

        mainThreadMre.Reset();
        workerThreadMre.Set();
        mainThreadMre.WaitOne();
        workerThreadMre.Reset();

        return results;
    }
}

如果我使用对象一次来处理方法队列,这很好,但如果我尝试这样的事情

BackgroundWorkManager manager = new BackgroundWorkManager();

for (int i = 0; i < 20; i++)
{
    manager.addTask("result1", (BackgroundWorkManager.ThreadTask)delegate
    {
        return (object)(1);
    });

    manager.process();
}

事情破裂了。我要么死锁,要么得到一个异常,说我正在编写结果的字典已经包含密钥(但 Visual Studio 调试器说它是空的)。在工作方法中添加“Thread.Sleep(1)”似乎可以解决它,这很奇怪。这是我第一次使用线程,所以我不确定我是否严重滥用了锁,或者什么。如果有人能提供一些关于我做错了什么的见解,将不胜感激。

【问题讨论】:

  • 您正在添加多个同名任务。为什么收到Dictionary already contains the key 会感到惊讶?此外,使用 PLINQ 应该很容易满足您的要求。简单的 AsParallel + ToDictionary 应该可以解决问题,无需锁定或任何东西。
  • 因为 process 方法创建了一个新的字典实例,然后将结果添加到该实例中。
  • 它创建新字典的事实没有任何区别。您将其分配给多个不同线程正在使用的类中的字段,因此很可能多个线程将写入Dictionary&lt;TKey, TValue&gt; 的同一实例。
  • 对,但锁定的目的是它们不会。只有一个线程应该能够获取委托并执行它,结果应该添加到字典中,此时 process 方法返回字典。进程方法是阻塞的。
  • 这是经典的生产者-消费者模式。我建议你看看BlockingCollection 和TPL Dataflow 而不是重新发明轮子。

标签: c# multithreading dictionary locks


【解决方案1】:

Parallel 类的版本:

List<Func<object>> actions = new List<Func<object>>();

actions.Add(delegate { return (object)(1); });
actions.Add(delegate { return (object)(1); });
actions.Add(delegate { return (object)(1); });

Dictionary<string, object> results = new Dictionary<string,object>();

Parallel.ForEach(actions,(f)=> {
    lock (results)
    {
        results.Add(Guid.NewGuid().ToString(), f());
    }
});

【讨论】:

    【解决方案2】:

    关于如何使用生产者-消费者模式,有很多选择。例如,您可以使用ActionBlock&lt;T&gt;(它是TPL Dataflow 的一部分)大大简化您的代码:

    var concurrentDictionary = new ConcurrentDictionary<string, object>();
    
    ActionBlock<Func<object>> actionBlock = new ActionBlock<Func<object>>((func) => 
    {
        var obj = func();
        concurrentDictionary.AddOrUpdate("someKey", obj, (s,o) => o);
    }, new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism =
                                           Environment.ProcessorCount });
    

    然后简单地发布您的代表:

    foreach (var task in tasks)
    {
        actionBlock.Post(() => (object) 1);
    }
    

    【讨论】:

    • 我没有意识到这个 .NET 可以为您做多少。谢谢。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-08-28
    • 1970-01-01
    • 2018-07-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-01-28
    相关资源
    最近更新 更多