From 65ac0273f9a14d9dc605c70d81e1a332baa9a3e6 Mon Sep 17 00:00:00 2001 From: Marcus Cuda Date: Sun, 29 Nov 2009 05:22:11 +0800 Subject: [PATCH] modified the thread queue to use wait/pulse instead of semaphores --- build/build.proj | 5 ++ src/Numerics/Threading/Parallel.cs | 8 +-- src/Numerics/Threading/ThreadQueue.cs | 53 ++++++++----------- src/UnitTests/PrecisionTest.cs | 2 - .../StatisticsTests/HistogramTests.cs | 3 -- 5 files changed, 31 insertions(+), 40 deletions(-) diff --git a/build/build.proj b/build/build.proj index 0f4c4da7..5975e2a0 100644 --- a/build/build.proj +++ b/build/build.proj @@ -19,6 +19,11 @@ + + + + + diff --git a/src/Numerics/Threading/Parallel.cs b/src/Numerics/Threading/Parallel.cs index 44bc4a43..7d87b68d 100644 --- a/src/Numerics/Threading/Parallel.cs +++ b/src/Numerics/Threading/Parallel.cs @@ -207,9 +207,9 @@ namespace MathNet.Numerics.Threading WaitForTasksToComplete(tasks); - for (var i = 0; i < tasks.Length; i++) + foreach (var t in tasks) { - localFinally(tasks[i].Result); + localFinally(t.Result); } CollectExceptions(tasks); @@ -453,9 +453,9 @@ namespace MathNet.Numerics.Threading /// The tasks. private static void WaitForTasksToComplete(Task[] tasks) { - for (var i = 0; i < tasks.Length; i++) + foreach (var task in tasks) { - tasks[i].Wait(); + task.Wait(); } } diff --git a/src/Numerics/Threading/ThreadQueue.cs b/src/Numerics/Threading/ThreadQueue.cs index 8ef80bf6..5fb23253 100644 --- a/src/Numerics/Threading/ThreadQueue.cs +++ b/src/Numerics/Threading/ThreadQueue.cs @@ -52,16 +52,6 @@ namespace MathNet.Numerics.Threading /// private static readonly Queue _queue = new Queue(); - /// - /// Maximum number of jobs that can be in the queue at the same time. - /// - private const int MaximumQueueLength = 4096; - - /// - /// Counting Semaphore to make the worker thread wait for jobs - /// - private static Semaphore _tasksAvailableSemaphore; - /// /// Running flag, used to signal worker threads to stop cleanly. /// @@ -113,9 +103,8 @@ namespace MathNet.Numerics.Threading lock (_queueSync) { _queue.Enqueue(task); + Monitor.Pulse(_queueSync); } - - _tasksAvailableSemaphore.Release(); } /// @@ -135,9 +124,9 @@ namespace MathNet.Numerics.Threading { _queue.Enqueue(task); } - } - _tasksAvailableSemaphore.Release(tasks.Count); + Monitor.PulseAll(_queueSync); + } } /// @@ -149,13 +138,9 @@ namespace MathNet.Numerics.Threading while (_running) { - // Wait until a job is available, or we should shut down - _tasksAvailableSemaphore.WaitOne(); - // Check whether we should shut down if (!_running) { - _tasksAvailableSemaphore.Release(); break; } @@ -167,13 +152,17 @@ namespace MathNet.Numerics.Threading { task = _queue.Dequeue(); } + else + { + Monitor.Wait(_queueSync); + } } if (task == null) { continue; } - + // ...and run it task.Compute(); } @@ -185,11 +174,11 @@ namespace MathNet.Numerics.Threading /// Number of worker threads. public static void Start(int numberOfThreads) { - // instead of throwing an out of range exception, simply normalize - numberOfThreads = Math.Max(1, Math.Min(1024, numberOfThreads)); - lock (_stateSync) { + // instead of throwing an out of range exception, simply normalize + numberOfThreads = Math.Max(1, Math.Min(1024, numberOfThreads)); + if (_threads != null) { if (_threads.Length == numberOfThreads) @@ -203,7 +192,7 @@ namespace MathNet.Numerics.Threading ThreadCount = numberOfThreads; Start(); } - } + } /// /// Start the thread queue, if it is not already running. @@ -217,8 +206,8 @@ namespace MathNet.Numerics.Threading return; } - _tasksAvailableSemaphore = new Semaphore(_queue.Count, MaximumQueueLength); - _running = true; + _running = true; + _threads = new Thread[ThreadCount]; for (var i = 0; i < _threads.Length; i++) @@ -238,25 +227,27 @@ namespace MathNet.Numerics.Threading /// public static void Shutdown() { + // try to stop the worker threads cleanly lock (_stateSync) { if (_threads == null) { return; } - - // try to stop the worker threads cleanly + _running = false; - _tasksAvailableSemaphore.Release(); + + lock (_queueSync) + { + Monitor.PulseAll(_queueSync); + } // wait until all threads have stopped foreach (var thread in _threads) { thread.Join(); } - - _tasksAvailableSemaphore.Close(); - _tasksAvailableSemaphore = null; + _threads = null; } } diff --git a/src/UnitTests/PrecisionTest.cs b/src/UnitTests/PrecisionTest.cs index 02e00e08..589ccc4f 100644 --- a/src/UnitTests/PrecisionTest.cs +++ b/src/UnitTests/PrecisionTest.cs @@ -184,8 +184,6 @@ namespace MathNet.Numerics.UnitTests public void CoerceZero() { Assert.AreEqual(0.0, Precision.CoerceZero(0d)); - Console.WriteLine(0.0.EpsilonOf()); - Console.WriteLine(Precision.Increment(0.0)); Assert.AreEqual(0.0, Precision.CoerceZero(Precision.Increment(0.0))); Assert.AreEqual(0.0, Precision.CoerceZero(Precision.Decrement(0.0))); diff --git a/src/UnitTests/StatisticsTests/HistogramTests.cs b/src/UnitTests/StatisticsTests/HistogramTests.cs index d06cdc22..45bc17bd 100644 --- a/src/UnitTests/StatisticsTests/HistogramTests.cs +++ b/src/UnitTests/StatisticsTests/HistogramTests.cs @@ -269,11 +269,8 @@ namespace MathNet.Numerics.UnitTests.StatisticsTests Assert.AreEqual(9, hist.BucketCount); - Console.WriteLine("{0}", hist); - for (int i = 1; i < 9; i++) { - Console.WriteLine("{0} : {1}", i, hist[i].Count); Assert.AreEqual(1.0, hist[i].Count); } Assert.AreEqual(2.0, hist[0].Count);