diff --git a/src/Numerics/Control.cs b/src/Numerics/Control.cs index aff59c25..89418049 100644 --- a/src/Numerics/Control.cs +++ b/src/Numerics/Control.cs @@ -43,7 +43,6 @@ namespace MathNet.Numerics CheckDistributionParameters = true; ThreadSafeRandomNumberGenerators = true; DisableParallelization = false; - InitialThreadBlockSize = 2; } /// @@ -76,12 +75,5 @@ namespace MathNet.Numerics /// Gets or sets a value indicating whether parallelization shall be disabled globally. /// public static bool DisableParallelization { get; set; } - - /// - /// Gets or sets the initial size of a - /// processing block (the number of elements the first thread should process). - /// - /// The initial size of the thread processing bloc. - public static int InitialThreadBlockSize { get; set; } } } diff --git a/src/Numerics/LinearAlgebra/Double/DenseVector.cs b/src/Numerics/LinearAlgebra/Double/DenseVector.cs index 58c9b59a..ffa5e643 100644 --- a/src/Numerics/LinearAlgebra/Double/DenseVector.cs +++ b/src/Numerics/LinearAlgebra/Double/DenseVector.cs @@ -577,7 +577,9 @@ namespace MathNet.Numerics.LinearAlgebra.Double var sum = 0.0; var syncLock = new object(); - Parallel.For(0, Count, + Parallel.For( + 0, + Count, () => 0.0, (index, localData) => { @@ -620,7 +622,9 @@ namespace MathNet.Numerics.LinearAlgebra.Double var sum = 0.0; var syncLock = new object(); - Parallel.For(0, Count, + Parallel.For( + 0, + Count, () => 0.0, (index, localData) => { @@ -648,7 +652,9 @@ namespace MathNet.Numerics.LinearAlgebra.Double var syncLock = new object(); - Parallel.For(0, Count, + Parallel.For( + 0, + Count, () => 0.0, (index, localData) => { diff --git a/src/Numerics/LinearAlgebra/Double/Vector.cs b/src/Numerics/LinearAlgebra/Double/Vector.cs index 3b0084f7..43ad3f46 100644 --- a/src/Numerics/LinearAlgebra/Double/Vector.cs +++ b/src/Numerics/LinearAlgebra/Double/Vector.cs @@ -609,20 +609,22 @@ namespace MathNet.Numerics.LinearAlgebra.Double var sum = 0.0; var syncLock = new object(); - Parallel.For(0, Count, - ()=> 0.0, + Parallel.For( + 0, + Count, + () => 0.0, (index, localData) => { localData += Math.Pow(Math.Abs(this[index]), p); return localData; }, - localResult=> + localResult => { lock (syncLock) { sum += localResult; } - } ); + }); return Math.Pow(sum, 1.0 / p); } @@ -637,7 +639,9 @@ namespace MathNet.Numerics.LinearAlgebra.Double { var max = 0.0; var syncLock = new object(); - Parallel.For(0, Count, + Parallel.For( + 0, + Count, () => 0.0, (index, localData) => { @@ -716,6 +720,7 @@ namespace MathNet.Numerics.LinearAlgebra.Double { return; } + Parallel.For(0, Count, index => target[index] = this[index]); } @@ -769,9 +774,7 @@ namespace MathNet.Numerics.LinearAlgebra.Double } else { - Parallel.For(0, count, - index => destination[destinationOffset + index] = this[offset + index] - ); + Parallel.For(0, count, index => destination[destinationOffset + index] = this[offset + index]); } } diff --git a/src/Numerics/Numerics.csproj b/src/Numerics/Numerics.csproj index 12a8c481..9869cd48 100644 --- a/src/Numerics/Numerics.csproj +++ b/src/Numerics/Numerics.csproj @@ -122,6 +122,7 @@ + diff --git a/src/Numerics/Threading/Parallel.cs b/src/Numerics/Threading/Parallel.cs index 8a9eb72c..ca888d94 100644 --- a/src/Numerics/Threading/Parallel.cs +++ b/src/Numerics/Threading/Parallel.cs @@ -30,17 +30,28 @@ namespace MathNet.Numerics.Threading { using System; using System.Collections.Generic; - using System.Threading; using Properties; /// - /// Provides support for parallel loops. + /// Provides support for parallel loops. /// internal static class Parallel { + /// + /// The amount to scale the foreach buffer after each iteration. + /// private const int ScalingFactor = 2; + + /// + /// The maximum size of the foreach buffer. + /// private const int MaxBlockSize = 65536; + /// + /// The initial size of the for each buffer. + /// + private const int IntialBlockSize = 1024; + /// /// Executes a for loop in which iterations may run in parallel. /// @@ -121,29 +132,26 @@ namespace MathNet.Numerics.Threading /// The end index, exclusive. /// The function delegate that returns the initial state of the local data for each thread. /// The delegate that is invoked once per iteration. - /// The delegate that is invoked once per iteration. - public static void For(int fromInclusive, int toExclusive, - Func localInit, - Func body, - Action localFinally) + /// The delegate that performs a final action on the local state of each thread. + public static void For(int fromInclusive, int toExclusive, Func localInit, Func body, Action localFinally) { var count = toExclusive - fromInclusive; var tasks = new Task[ThreadQueue.ThreadCount]; var size = count / tasks.Length; - if (count <= 0){ - // if (count <= 1) -// { -// if (count == 1) - // { - // body(fromInclusive); - // } + var intial = localInit(); + + // fast forward execution if it's only one or none items + if (count <= 1) + { + if (count == 1) + { + localFinally(body(fromInclusive, intial)); + } return; } - var intial = localInit(); - // partition the jobs into separate sets for each but the last worked thread for (var i = 0; i < tasks.Length - 1; i++) { @@ -157,8 +165,10 @@ namespace MathNet.Numerics.Threading { localresult = body(j, (T)localresult); } + return (T)localresult; - }, intial ); + }, + intial); ThreadQueue.Enqueue(tasks[i]); } @@ -171,20 +181,25 @@ namespace MathNet.Numerics.Threading { localresult = body(i, (T)localresult); } + return (T)localresult; - }, intial ); + }, + intial); ThreadQueue.Enqueue(tasks[tasks.Length - 1]); if (tasks.Length <= 0) { return; } + WaitForTasksToComplete(tasks); + for (var i = 0; i < tasks.Length; i++) { localFinally(tasks[i].Result); } - CollectExceptionsAndDisposeTasks(tasks); + + CollectExceptions(tasks); } /// @@ -221,10 +236,10 @@ namespace MathNet.Numerics.Threading return; } - var enumerator = source.GetEnumerator(); - var maxBlockSize = Control.InitialThreadBlockSize; - var scalingFactor = ScalingFactor; + var maxBlockSize = IntialBlockSize; var tasks = new List(); + + var enumerator = source.GetEnumerator(); while (enumerator.MoveNext()) { var pos = 0; @@ -249,13 +264,13 @@ namespace MathNet.Numerics.Threading ThreadQueue.Enqueue(task); tasks.Add(task); - maxBlockSize = Math.Min(MaxBlockSize, maxBlockSize * scalingFactor); + maxBlockSize = Math.Min(MaxBlockSize, maxBlockSize * ScalingFactor); } if (tasks.Count > 0) { WaitForTasksToComplete(tasks.ToArray()); - CollectExceptionsAndDisposeTasks(tasks); + CollectExceptions(tasks); } } @@ -268,21 +283,42 @@ namespace MathNet.Numerics.Threading /// The function delegate that returns the initial state of the local data for each thread. /// The delegate that is invoked once per iteration. /// The delegate that performs a final action on the local state of each thread. - public static void ForEach(IEnumerable source, Func localInit, - Func body, Action localFinally) + public static void ForEach(IEnumerable source, Func localInit, Func body, Action localFinally) { if (body == null) { throw new ArgumentNullException("body"); } - var enumerator = source.GetEnumerator(); - var maxBlockSize = Control.InitialThreadBlockSize; - var scalingFactor = ScalingFactor; + // fast forward execution in case parallelization is disabled + if (Control.DisableParallelization + || ThreadQueue.ThreadCount <= 1 + || ThreadQueue.IsInWorkerThread) + { + var localResult = localInit(); + foreach (var item in source) + { + localResult = body(item, localResult); + } + + localFinally(localResult); + return; + } + + // source is a IList, call For instead. + if (source is IList) + { + var list = (IList)source; + For(0, list.Count, localInit, (i, local) => body(list[i], local), localFinally); + return; + } + + var maxBlockSize = IntialBlockSize; var tasks = new List>(); var intial = localInit(); + var enumerator = source.GetEnumerator(); while (enumerator.MoveNext()) { var pos = 0; @@ -304,25 +340,30 @@ namespace MathNet.Numerics.Threading { localresult = body(list[i], (TLocal)localresult); } + return (TLocal)localresult; - }, intial); + }, + intial); ThreadQueue.Enqueue(task); tasks.Add(task); - maxBlockSize = Math.Min(MaxBlockSize, maxBlockSize * scalingFactor); + maxBlockSize = Math.Min(MaxBlockSize, maxBlockSize * ScalingFactor); } if (tasks.Count <= 0) { return; } + var taskArray = tasks.ToArray(); WaitForTasksToComplete(taskArray); + for (var i = 0; i < taskArray.Length; i++) { localFinally(tasks[i].Result); } - CollectExceptionsAndDisposeTasks(taskArray); + + CollectExceptions(taskArray); } /// @@ -355,7 +396,7 @@ namespace MathNet.Numerics.Threading || ThreadQueue.ThreadCount <= 1 || ThreadQueue.IsInWorkerThread) { - for (int i = 0; i < actions.Length; i++) + for (var i = 0; i < actions.Length; i++) { actions[i](); } @@ -392,7 +433,7 @@ namespace MathNet.Numerics.Threading WaitForTasksToComplete(tasks); - CollectExceptionsAndDisposeTasks(tasks); + CollectExceptions(tasks); } /// @@ -401,17 +442,17 @@ namespace MathNet.Numerics.Threading /// The tasks. private static void WaitForTasksToComplete(Task[] tasks) { - for (var i = 0; i < tasks.Length; i++) - { - tasks[i].Wait(); - } + for (var i = 0; i < tasks.Length; i++) + { + tasks[i].Wait(); + } } /// /// Collects the exceptions and dispose tasks. /// /// The tasks. - private static void CollectExceptionsAndDisposeTasks(IEnumerable tasks) + private static void CollectExceptions(IEnumerable tasks) { // collect all thrown exceptions and dispose the jobs var exceptions = new List(); diff --git a/src/Numerics/Threading/Task.cs b/src/Numerics/Threading/Task.cs index 33dd438f..c0a5ca12 100644 --- a/src/Numerics/Threading/Task.cs +++ b/src/Numerics/Threading/Task.cs @@ -113,56 +113,15 @@ namespace MathNet.Numerics.Threading _body(); } - public void Wait() - { - while(!IsCompleted && !IsFaulted) - { - Thread.Sleep(100); - } - - } - } - - /// - /// Internal Generic Parallel Task Handle. - /// - internal class Task : Task - { - /// - /// Delegate to the task's action. - /// - private readonly Func _body; - - private readonly object _state; - - /// - /// Gets the result of the task. - /// - /// The result of the task. - public TResult Result { get; private set; } - /// - /// Initializes a new instance of the Task class. + /// Waits for the task to complete execution. /// - /// An object representing data to be used by the action. - /// Delegate to the task's action. - public Task(Func body, object state) + public void Wait() { - if (body == null) + while (!IsCompleted && !IsFaulted) { - throw new ArgumentNullException("body"); + Thread.Sleep(100); } - - _state = state; - _body = body; - } - - /// - /// Runs the actual task. - /// - protected override void DoCompute() - { - Result = _body(_state); } } } diff --git a/src/Numerics/Threading/TaskOfT.cs b/src/Numerics/Threading/TaskOfT.cs new file mode 100644 index 00000000..efdfaf1f --- /dev/null +++ b/src/Numerics/Threading/TaskOfT.cs @@ -0,0 +1,79 @@ +// +// Math.NET Numerics, part of the Math.NET Project +// http://mathnet.opensourcedotnet.info +// +// Copyright (c) 2009 Math.NET +// +// Permission is hereby granted, free of charge, to any person +// obtaining a copy of this software and associated documentation +// files (the "Software"), to deal in the Software without +// restriction, including without limitation the rights to use, +// copy, modify, merge, publish, distribute, sublicense, and/or sell +// copies of the Software, and to permit persons to whom the +// Software is furnished to do so, subject to the following +// conditions: +// +// The above copyright notice and this permission notice shall be +// included in all copies or substantial portions of the Software. +// +// THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, +// EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES +// OF MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND +// NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR COPYRIGHT +// HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, +// WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING +// FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR +// OTHER DEALINGS IN THE SOFTWARE. +// + +namespace MathNet.Numerics.Threading +{ + using System; + + /// + /// Internal Generic Parallel Task Handle. + /// + /// The type of the result. + internal class Task : Task + { + /// + /// Delegate to the task's action. + /// + private readonly Func _body; + + /// + /// Variable used to hold state information between interations. + /// + private readonly object _state; + + /// + /// Gets the result of the task. + /// + /// The result of the task. + public TResult Result { get; private set; } + + /// + /// Initializes a new instance of the Task class. + /// + /// Delegate to the task's action. + /// An object representing data to be used by the action. + public Task(Func body, object state) + { + if (body == null) + { + throw new ArgumentNullException("body"); + } + + _state = state; + _body = body; + } + + /// + /// Runs the actual task. + /// + protected override void DoCompute() + { + Result = _body(_state); + } + } +} diff --git a/src/Numerics/Threading/ThreadQueue.cs b/src/Numerics/Threading/ThreadQueue.cs index 4469e915..8ef80bf6 100644 --- a/src/Numerics/Threading/ThreadQueue.cs +++ b/src/Numerics/Threading/ThreadQueue.cs @@ -94,7 +94,7 @@ namespace MathNet.Numerics.Threading /// /// Gets a value indicating whether the current thread is a parallelized worker thread. /// - internal static bool IsInWorkerThread + public static bool IsInWorkerThread { get { return _isInWorkerThread; } } @@ -103,7 +103,7 @@ namespace MathNet.Numerics.Threading /// Add a job to the queue. /// /// The job to run. - internal static void Enqueue(Task task) + public static void Enqueue(Task task) { if (!_running) { @@ -122,7 +122,7 @@ namespace MathNet.Numerics.Threading /// Add a set of jobs to the queue. /// /// The jobs to run. - internal static void Enqueue(IList tasks) + public static void Enqueue(IList tasks) { if (!_running) { @@ -176,7 +176,6 @@ namespace MathNet.Numerics.Threading // ...and run it task.Compute(); - //task.Set(); } } @@ -184,7 +183,7 @@ namespace MathNet.Numerics.Threading /// Start or restart the queue with the specified number of worker threads. /// /// Number of worker threads. - internal static void Start(int numberOfThreads) + public static void Start(int numberOfThreads) { // instead of throwing an out of range exception, simply normalize numberOfThreads = Math.Max(1, Math.Min(1024, numberOfThreads)); @@ -209,7 +208,7 @@ namespace MathNet.Numerics.Threading /// /// Start the thread queue, if it is not already running. /// - internal static void Start() + public static void Start() { lock (_stateSync) { @@ -237,7 +236,7 @@ namespace MathNet.Numerics.Threading /// /// Stop the thread queue, if it is running. /// - internal static void Shutdown() + public static void Shutdown() { lock (_stateSync) {