【问题标题】:What's the best way to have multiple threads doing work, and waiting for all of them to complete?让多个线程工作并等待所有线程完成的最佳方法是什么?
【发布时间】:2009-12-16 15:00:52
【问题描述】:

我正在编写一个简单的应用程序(对于我的妻子来说同样如此 :-P ),它对可能大量的图像进行一些图像处理(调整大小、时间戳等)。所以我正在编写一个可以同步和异步执行此操作的库。我决定使用Event-based Asynchronous Pattern。使用此模式时,您需要在工作完成后引发事件。这是我在知道何时完成时遇到问题的地方。所以基本上,在我的 DownsizeAsync 方法(缩小图像的异步方法)中,我正在做这样的事情:

    public void DownsizeAsync(string[] files, string destination)
    {
        foreach (var name in files)
        {
            string temp = name; //countering the closure issue
            ThreadPool.QueueUserWorkItem(f =>
            {
                string newFileName = this.DownsizeImage(temp, destination);
                this.OnImageResized(newFileName);
            });
        }
     }

现在最棘手的部分是知道它们何时全部完成。

这是我考虑过的:使用 ManualResetEvents 像这里:http://msdn.microsoft.com/en-us/library/3dasc8as%28VS.80%29.aspx 但我遇到的问题是您只能等待 64 个或更少的事件。我可能还有更多图片。

第二种选择:有一个计数器来计算已经完成的图像,并在计数达到总数时引发事件:

public void DownsizeAsync(string[] files, string destination)
{
    foreach (var name in files)
    {
        string temp = name; //countering the closure issue
        ThreadPool.QueueUserWorkItem(f =>
        {
            string newFileName = this.DownsizeImage(temp, destination);
            this.OnImageResized(newFileName);
            total++;
            if (total == files.Length)
            {
                this.OnDownsizeCompleted(new AsyncCompletedEventArgs(null, false, null));
            }
        });
    }


}

private volatile int total = 0;

现在感觉“hacky”,我不完全确定这是否是线程安全的。

那么,我的问题是,最好的方法是什么?有没有另一种方法来同步所有线程?我不应该使用线程池吗?谢谢!!

更新根据 cmets 的反馈和一些答案,我决定采用这种方法:

首先,我创建了一个扩展方法,将一个可枚举的对象批量化为“​​批次”:

    public static IEnumerable<IEnumerable<T>> GetBatches<T>(this IEnumerable<T> source, int batchCount)
    {
        for (IEnumerable<T> s = source; s.Any(); s = s.Skip(batchCount))
        {
            yield return s.Take(batchCount);
        }
    }

基本上,如果你这样做:

        foreach (IEnumerable<int> batch in Enumerable.Range(1, 95).GetBatches(10))
        {
            foreach (int i in batch)
            {
                Console.Write("{0} ", i);
            }
            Console.WriteLine();
        }

你得到这个输出:

1 2 3 4 5 6 7 8 9 10
11 12 13 14 15 16 17 18 19 20
21 22 23 24 25 26 27 28 29 30
31 32 33 34 35 36 37 38 39 40
41 42 43 44 45 46 47 48 49 50
51 52 53 54 55 56 57 58 59 60
61 62 63 64 65 66 67 68 69 70
71 72 73 74 75 76 77 78 79 80
81 82 83 84 85 86 87 88 89 90
91 92 93 94 95

这个想法是(正如 cmets 中的某个人指出的那样)没有必要为每个图像创建单独的线程。因此,我将图像分批成 [machine.cores * 2] 个批次。然后,我将使用第二种方法,即简单地保持计数器运行,当计数器达到我期望的总数时,我就知道我已经完成了。

我现在确信它实际上是线程安全的原因是因为我已根据MSDN 将总变量标记为 volatile:

通常使用 volatile 修饰符 对于被访问的字段 多个线程不使用 lock 语句来序列化访问。 使用 volatile 修饰符可确保 一个线程检索最多 另一个人写入的最新值 线程

表示我应该清楚(如果没有,请告诉我!!)

这是我要使用的代码:

    public void DownsizeAsync(string[] files, string destination)
    {
        int cores = Environment.ProcessorCount * 2;
        int batchAmount = files.Length / cores;

        foreach (var batch in files.GetBatches(batchAmount))
        {
            var temp = batch.ToList(); //counter closure issue
            ThreadPool.QueueUserWorkItem(b =>
            {
                foreach (var item in temp)
                {
                    string newFileName = this.DownsizeImage(item, destination);
                    this.OnImageResized(newFileName);
                    total++;
                    if (total == files.Length)
                    {
                        this.OnDownsizeCompleted(new AsyncCompletedEventArgs(null, false, null));
                    }
                }
            });
        }
    }

我愿意接受反馈,因为我绝不是多线程方面的专家,所以如果有人对此有任何问题,或者有更好的想法,请告诉我。 (是的,这只是一个自制的应用程序,但我对如何利用我在这里获得的知识来改进我们在工作中使用的搜索/索引服务有一些想法。)现在我会一直保持这个问题,直到我感觉我正在使用正确的方法。感谢大家的帮助。

【问题讨论】:

  • total++ 在我看来不是线程安全的!
  • 计数器方法具有很高的可扩展性。只需使用 Interlocked.Increment 和 .Decrement.. 使其线程安全。
  • 在实践中,64 是这个限制吗?即使使用 64 个物理内核,您也会遇到内存和磁盘瓶颈,这些瓶颈在并行访问下会降级得更快。但是,这些可能会在一两年内消失。
  • 我只是对为什么线程不能重新分配更多任务感到困惑;你真的需要为每个文件创建一个线程吗?这似乎效率低下。
  • @Dean J:很好。请参阅我的编辑。我想我将采用“批量”方法,每个线程将负责更多的图像。

标签: c# multithreading threadpool


【解决方案1】:

最简单的方法是创建新线程,然后在每个线程上调用Thread.Join。您可以使用信号量或类似的东西 - 但创建新线程可能更容易。

在 .NET 4.0 中,您可以使用并行扩展非常轻松地处理任务。

作为另一个使用线程池的替代方案,您可以创建一个委托并在其上调用BeginInvoke,以返回一个IAsyncResult - 然后您可以获得每个结果的WaitHandle通过AsyncWaitHandle 属性,并调用WaitHandle.WaitAll

编辑:正如 cmets 中所指出的,在某些实现中,您一次最多只能使用 64 个句柄调用 WaitAll。替代方案可以依次调用WaitOne,或者批量调用WaitAll。只要您是从不会阻塞线程池的线程中执行此操作,这并不重要。另请注意,您不能从 STA 线程调用 WaitAll

【讨论】:

  • Jon- 你忘记了 WaitAll 的 64 位限制。
  • @RichardOD:谢谢,已经添加了一个注释。
【解决方案2】:

您仍然希望使用 ThreadPool,因为它将管理它同时运行的线程数。我最近遇到了类似的问题并这样解决了:

var dispatcher = new ThreadPoolDispatcher();
dispatcher = new ChunkingDispatcher(dispatcher, 10);

foreach (var image in images)
{
    dispatcher.Add(new ResizeJob(image));
}

dispatcher.WaitForJobsToFinish();

IDispatcher 和 IJob 如下所示:

public interface IJob
{
    void Execute();
}

public class ThreadPoolDispatcher : IDispatcher
{
    private IList<ManualResetEvent> resetEvents = new List<ManualResetEvent>();

    public void Dispatch(IJob job)
    {
        var resetEvent = CreateAndTrackResetEvent();
        var worker = new ThreadPoolWorker(job, resetEvent);
        ThreadPool.QueueUserWorkItem(new WaitCallback(worker.ThreadPoolCallback));
    }

    private ManualResetEvent CreateAndTrackResetEvent()
    {
        var resetEvent = new ManualResetEvent(false);
        resetEvents.Add(resetEvent);
        return resetEvent;
    }

    public void WaitForJobsToFinish()
    {
        WaitHandle.WaitAll(resetEvents.ToArray() ?? new ManualResetEvent[] { });
        resetEvents.Clear();
    }
}

然后用一个装饰器来分块使用ThreadPool:

public class ChunkingDispatcher : IDispatcher
{
    private IDispatcher dispatcher;
    private int numberOfJobsDispatched;
    private int chunkSize;

    public ChunkingDispatcher(IDispatcher dispatcher, int chunkSize)
    {
        this.dispatcher = dispatcher;
        this.chunkSize = chunkSize;
    }

    public void Dispatch(IJob job)
    {
        dispatcher.Dispatch(job);

        if (++numberOfJobsDispatched % chunkSize == 0)
            WaitForJobsToFinish();
    }

    public void WaitForJobsToFinish()
    {
        dispatcher.WaitForJobsToFinish();
    }
}

IDispatcher 抽象非常适合替换您的线程技术。我有另一个实现,它是 SingleThreadedDispatcher,你可以像 Jon Skeet 建议的那样制作一个 ThreadStart 版本。然后很容易运行每一个,看看你得到了什么样的性能。 SingleThreadedDispatcher 在调试代码或不想杀死机器上的处理器时非常有用。

编辑:我忘了添加 ThreadPoolWorker 的代码:

public class ThreadPoolWorker
{
    private IJob job;
    private ManualResetEvent doneEvent;

    public ThreadPoolWorker(IJob job, ManualResetEvent doneEvent)
    {
        this.job = job;
        this.doneEvent = doneEvent;
    }

    public void ThreadPoolCallback(object state)
    {
        try
        {
            job.Execute();
        }
        finally
        {
            doneEvent.Set();
        }
    }
}

【讨论】:

    【解决方案3】:

    最简单有效的解决方案是使用计数器并使其线程安全。这将消耗更少的内存,并且可以扩展到更多的线程

    这是一个示例

    int itemCount = 0;
    for (int i = 0; i < 5000; i++)
    {
        Interlocked.Increment(ref itemCount);
    
        ThreadPool.QueueUserWorkItem(x=>{
            try
            {
                //code logic here.. sleep is just for demo
                Thread.Sleep(100);
            }
            finally
            {
                Interlocked.Decrement(ref itemCount);
            }
        });
    }
    
    while (itemCount > 0)
    {
        Console.WriteLine("Waiting for " + itemCount + " threads...");
        Thread.Sleep(100);
    }
    Console.WriteLine("All Done!");
    

    【讨论】:

    • +1。我使用这种模式,尽管我检查了 Interlocked.Decrement() 的结果以查看我们是否已经达到零,如果是,则设置一个事件以指示所有项目都已完成。这样你就不需要对 itemCount 进行投票了。
    • 我喜欢这种方法,但我想知道如果我们正在递增的变量被标记为 volatile,是否需要这样做。
    • volatile 保证变量不被缓存,访问器是原子的,但不保证读+改+写的全部操作都是原子的
    • @BFree。使用您最初的方法,您完成了 total++,这等于 total = total + 1。现在考虑多个线程这样做。 Pratap 就在现场 - volatile 不会让事情变得线程安全,这完全是关于缓存一致性。
    【解决方案4】:

    我使用SmartThreadPool 成功解决了这个问题。还有一个Codeplex 网站关于de assembly。

    SmartThreadPool 可以帮助解决其他问题,例如某些线程不能同时运行而其他线程可以。

    【讨论】:

      【解决方案5】:

      我使用静态实用程序方法检查所有单独的等待句柄..

          public static void WaitAll(WaitHandle[] handles)
          {
              if (handles == null)
                  throw new ArgumentNullException("handles",
                      "WaitHandle[] handles was null");
              foreach (WaitHandle wh in handles) wh.WaitOne();
          }
      

      然后在我的主线程中,我创建了一个包含这些等待句柄的列表,并且对于我放入 ThreadPool 队列中的每个委托,我将等待句柄添加到列表中......

       List<WaitHandle> waitHndls = new List<WaitHandle>();
       foreach (iterator logic )
       {
            ManualResetEvent txEvnt = new ManualResetEvent(false);
      
            ThreadPool.QueueUserWorkItem(
                 delegate
                     {
                         try { // Code to process each task... }
                         // Finally, set each wait handle when done
                         finally { lock (locker) txEvnt.Set(); } 
                     });
            waitHndls.Add(txEvnt);  // Add wait handle to List
       }
       util.WaitAll(waitHndls.ToArray());   // Check all wait Handles in List
      

      【讨论】:

        【解决方案6】:

        .Net 4.0 使多线程变得更加容易(尽管您仍然可以用副作用来打自己)。

        【讨论】:

        • 好吧 ... C 中的多线程更难。理论上图灵等效或你有什么,如果不一样,概念上相似,但额外的代码行确实妨碍了理解核心。
        【解决方案7】:

        另一种选择是使用管道。

        您将所有要完成的工作发布到管道,然后从每个线程的管道中读取数据。当管道为空时,你就完成了,线程自己结束了,每个人都很高兴(当然要确保你首先产生所有的工作,然后消费它)

        【讨论】:

          【解决方案8】:

          我建议将未触及的图像放入队列中,当您从队列中读取时启动一个线程并将其System.Threading.Thread.ManagedThreadId 属性与文件名一起插入到字典中。这样,您的 UI 可以同时列出待处理和活动文件。

          当每个线程完成时,它会调用一个回调例程,并传回其 ManagedThreadId。此回调(作为委托传递给线程)从字典中删除线程的 id,从队列中启动另一个线程,并更新 UI。

          当队列和字典都为空时,您就完成了。

          稍微复杂一些,但这样可以获得响应式 UI,您可以轻松控制活动线程的数量,并且可以查看正在运行的内容。收集统计数据。使用 WPF 并为每个文件设置进度条。她不禁为之感动。

          【讨论】:

            猜你喜欢
            • 1970-01-01
            • 2011-05-10
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            • 2012-04-09
            • 1970-01-01
            相关资源
            最近更新 更多