From 960717c37444cff043c7246b4c644e5886668c22 Mon Sep 17 00:00:00 2001 From: Marcus Cuda Date: Fri, 21 Aug 2009 21:01:33 +0800 Subject: [PATCH] Parallel: first take on aggregate parallel methods Signed-off-by: Marcus Cuda --- .../LinearAlgebra/Double/DenseVector.cs | 65 +++++++-- src/Numerics/LinearAlgebra/Double/Vector.cs | 39 +++-- src/Numerics/Threading/Parallel.cs | 133 +++++++++++++++++- src/Numerics/Threading/Task.cs | 53 ++++++- 4 files changed, 262 insertions(+), 28 deletions(-) diff --git a/src/Numerics/LinearAlgebra/Double/DenseVector.cs b/src/Numerics/LinearAlgebra/Double/DenseVector.cs index a0ec57b7..08d0751a 100644 --- a/src/Numerics/LinearAlgebra/Double/DenseVector.cs +++ b/src/Numerics/LinearAlgebra/Double/DenseVector.cs @@ -573,11 +573,23 @@ namespace MathNet.Numerics.LinearAlgebra.Double /// Scalar ret = sum(abs(this[i])) public override double Norm1() { - double sum = 0; - for (var i = 0; i < Data.Length; i++) - { - sum += Math.Abs(Data[i]); - } + var sum = 0.0; + var syncLock = new object(); + + Parallel.For(0, Count, + () => 0.0, + (index, localData) => + { + localData += Math.Abs(Data[index]); + return localData; + }, + localResult => + { + lock (syncLock) + { + sum += localResult; + } + }); return sum; } @@ -605,10 +617,22 @@ namespace MathNet.Numerics.LinearAlgebra.Double } var sum = 0.0; - for (var i = 0; i < Data.Length; i++) - { - sum += Math.Pow(Math.Abs(Data[i]), p); - } + var syncLock = new object(); + + Parallel.For(0, Count, + () => 0.0, + (index, localData) => + { + localData += Math.Pow(Math.Abs(Data[index]), p); + return localData; + }, + localResult => + { + lock (syncLock) + { + sum += localResult; + } + }); return Math.Pow(sum, 1.0 / p); } @@ -620,10 +644,23 @@ namespace MathNet.Numerics.LinearAlgebra.Double public override double NormInfinity() { double max = 0; - for (int i = 0; i < Data.Length; i++) - { - max = Math.Max(max, Math.Abs(Data[i])); - } + + var syncLock = new object(); + + Parallel.For(0, Count, + () => 0.0, + (index, localData) => + { + localData = Math.Max(localData, Math.Abs(Data[index])); + return localData; + }, + localResult => + { + lock (syncLock) + { + max = Math.Max(max, localResult); + } + }); return max; } @@ -720,7 +757,7 @@ namespace MathNet.Numerics.LinearAlgebra.Double data[i] = Double.Parse(token.Value, NumberStyles.Any, formatProvider); token = token.Next; - if(token != null) + if (token != null) { token = token.Next; } diff --git a/src/Numerics/LinearAlgebra/Double/Vector.cs b/src/Numerics/LinearAlgebra/Double/Vector.cs index 1df0044b..73a397ed 100644 --- a/src/Numerics/LinearAlgebra/Double/Vector.cs +++ b/src/Numerics/LinearAlgebra/Double/Vector.cs @@ -680,11 +680,21 @@ namespace MathNet.Numerics.LinearAlgebra.Double } var sum = 0.0; - - foreach (var pair in GetIndexedEnumerator()) - { - sum += Math.Pow(Math.Abs(pair.Value), p); - } + var syncLock = new object(); + Parallel.ForEach(GetIndexedEnumerator(), + ()=> 0.0, + (pair, localData) => + { + localData += Math.Pow(Math.Abs(pair.Value), p); + return localData; + }, + localResult=> + { + lock (syncLock) + { + sum += localResult; + } + } ); return Math.Pow(sum, 1.0 / p); } @@ -698,10 +708,21 @@ namespace MathNet.Numerics.LinearAlgebra.Double public virtual double NormInfinity() { var max = 0.0; - foreach (var pair in GetIndexedEnumerator()) - { - max = Math.Max(max, Math.Abs(pair.Value)); - } + var syncLock = new object(); + Parallel.ForEach(GetIndexedEnumerator(), + () => 0.0, + (pair, localData) => + { + localData = Math.Max(localData, Math.Abs(pair.Value)); + return localData; + }, + localResult => + { + lock (syncLock) + { + max = Math.Max(localResult, max); + } + }); return max; } diff --git a/src/Numerics/Threading/Parallel.cs b/src/Numerics/Threading/Parallel.cs index 5d04a9ed..cbe6b3e6 100644 --- a/src/Numerics/Threading/Parallel.cs +++ b/src/Numerics/Threading/Parallel.cs @@ -110,6 +110,69 @@ namespace MathNet.Numerics.Threading Invoke(actions); } + /// + /// Executes a for loop in which iterations may run in parallel. + /// + /// The type of the thread-local data. + /// The start index, inclusive. + /// 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) + { + var count = toExclusive - fromInclusive; + var tasks = new Task[ThreadQueue.ThreadCount]; + var size = count / tasks.Length; + + 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++) + { + var start = fromInclusive + (i * size); + var stop = fromInclusive + ((i + 1) * size); + tasks[i] = new Task(intial, + localData => + { + var localresult = localData; + for (var j = start; j < stop; j++) + { + localresult = body(j, localresult); + } + return localresult; + } ); + ThreadQueue.Enqueue(tasks[i]); + } + + // add another set for last worker thread + tasks[tasks.Length - 1] = new Task(intial, + localData => + { + var localresult = localData; + for (var i = fromInclusive + ((tasks.Length - 1) * size); i < toExclusive; i++) + { + localresult = body(i, localresult); + } + return localresult; + } ); + + 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); + } + /// /// Executes a for each operation on an IEnumerable{T} in which iterations may run in parallel. /// @@ -160,7 +223,7 @@ namespace MathNet.Numerics.Threading list[pos++] = enumerator.Current; count++; } - + var task = new Task( () => { @@ -182,6 +245,72 @@ namespace MathNet.Numerics.Threading } } + /// + /// Executes a for each operation on an IEnumerable in which iterations may run in parallel. + /// + /// The type of the data in the source. + /// The type of the thread-local data. + /// An enumerable data source. + /// 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) + { + if (body == null) + { + throw new ArgumentNullException("body"); + } + + var enumerator = source.GetEnumerator(); + var maxBlockSize = Control.InitialThreadBlockSize; + var scalingFactor = Control.BlockScalingFactor; + var tasks = new List>(); + + var intial = localInit(); + + while (enumerator.MoveNext()) + { + var pos = 0; + var list = new TSource[maxBlockSize]; + list[pos++] = enumerator.Current; + + var count = 1; + while (count < maxBlockSize && enumerator.MoveNext()) + { + list[pos++] = enumerator.Current; + count++; + } + + var task = new Task(intial, + localData => + { + var localresult = localData; + for (var i = 0; i < pos; i++) + { + localresult = body(list[i], localresult); + } + return localresult; + }); + + ThreadQueue.Enqueue(task); + tasks.Add(task); + maxBlockSize = Math.Min(Control.MaximumBlockSize, 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); + } + /// /// Executes each of the provided actions inside a discrete, asynchronous task. /// @@ -233,7 +362,7 @@ namespace MathNet.Numerics.Threading { // create a job for each action var tasks = new Task[actions.Length]; - for (int i = 0; i < tasks.Length; i++) + for (var i = 0; i < tasks.Length; i++) { Action action = actions[i]; if (action == null) diff --git a/src/Numerics/Threading/Task.cs b/src/Numerics/Threading/Task.cs index 544a1cb3..72ace2c3 100644 --- a/src/Numerics/Threading/Task.cs +++ b/src/Numerics/Threading/Task.cs @@ -46,7 +46,7 @@ namespace MathNet.Numerics.Threading /// /// Delegate to the task's action. internal Task(Action body) - : base(false, EventResetMode.ManualReset) + : this() { if (body == null) { @@ -56,6 +56,8 @@ namespace MathNet.Numerics.Threading _body = body; } + protected Task() : base(false, EventResetMode.ManualReset) { } + /// /// Gets a value indicating whether the task has thrown one or more exceptions while executing. /// @@ -67,12 +69,12 @@ namespace MathNet.Numerics.Threading /// /// Gets the exception thrown by the task, if any. /// - internal Exception Exception { get; private set; } + protected internal Exception Exception { get; set; } /// /// Run the task. /// - internal void Compute() + internal virtual void Compute() { try { @@ -84,4 +86,49 @@ namespace MathNet.Numerics.Threading } } } + + /// + /// Internal Generic Parallel Task Handle. + /// + internal class Task : Task + { + /// + /// Delegate to the task's action. + /// + private readonly Func _body; + + //private T _initialValue; + + public T Result { get; private set; } + + /// + /// Initializes a new instance of the Task class. + /// + /// The initial value. + /// Delegate to the task's action. + internal Task(T intialValue, Func body) + { + if (body == null) + { + throw new ArgumentNullException("body"); + } + Result = intialValue; + _body = body; + } + + /// + /// Run the task. + /// + internal override void Compute() + { + try + { + Result = _body(Result); + } + catch (Exception e) + { + Exception = e; + } + } + } }