Browse Source

Parallel: first take on aggregate parallel methods

Signed-off-by: Marcus Cuda <marcus@cuda.net>
la-knuth
Marcus Cuda 17 years ago
parent
commit
960717c374
  1. 65
      src/Numerics/LinearAlgebra/Double/DenseVector.cs
  2. 39
      src/Numerics/LinearAlgebra/Double/Vector.cs
  3. 133
      src/Numerics/Threading/Parallel.cs
  4. 53
      src/Numerics/Threading/Task.cs

65
src/Numerics/LinearAlgebra/Double/DenseVector.cs

@ -573,11 +573,23 @@ namespace MathNet.Numerics.LinearAlgebra.Double
/// <returns>Scalar ret = sum(abs(this[i]))</returns>
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;
}

39
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;
}

133
src/Numerics/Threading/Parallel.cs

@ -110,6 +110,69 @@ namespace MathNet.Numerics.Threading
Invoke(actions);
}
/// <summary>
/// Executes a for loop in which iterations may run in parallel.
/// </summary>
/// <typeparam name="T">The type of the thread-local data.</typeparam>
/// <param name="fromInclusive">The start index, inclusive.</param>
/// <param name="toExclusive">The end index, exclusive.</param>
/// <param name="localInit">The function delegate that returns the initial state of the local data for each thread.</param>
/// <param name="body">The delegate that is invoked once per iteration.</param>
/// <param name="localFinally">The delegate that is invoked once per iteration.</param>
public static void For<T>(int fromInclusive, int toExclusive,
Func<T> localInit,
Func<int, T, T> body,
Action<T> localFinally)
{
var count = toExclusive - fromInclusive;
var tasks = new Task<T>[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<T>(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<T>(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);
}
/// <summary>
/// Executes a for each operation on an IEnumerable{T} in which iterations may run in parallel.
/// </summary>
@ -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
}
}
/// <summary>
/// Executes a for each operation on an IEnumerable<TSource> in which iterations may run in parallel.
/// </summary>
/// <typeparam name="TSource">The type of the data in the source.</typeparam>
/// <typeparam name="TLocal">The type of the thread-local data.</typeparam>
/// <param name="source">An enumerable data source.</param>
/// <param name="localInit">The function delegate that returns the initial state of the local data for each thread.</param>
/// <param name="body">The delegate that is invoked once per iteration.</param>
/// <param name="localFinally">The delegate that performs a final action on the local state of each thread.</param>
public static void ForEach<TSource, TLocal>(IEnumerable<TSource> source, Func<TLocal> localInit,
Func<TSource, TLocal, TLocal> body, Action<TLocal> localFinally)
{
if (body == null)
{
throw new ArgumentNullException("body");
}
var enumerator = source.GetEnumerator();
var maxBlockSize = Control.InitialThreadBlockSize;
var scalingFactor = Control.BlockScalingFactor;
var tasks = new List<Task<TLocal>>();
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<TLocal>(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);
}
/// <summary>
/// Executes each of the provided actions inside a discrete, asynchronous task.
/// </summary>
@ -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)

53
src/Numerics/Threading/Task.cs

@ -46,7 +46,7 @@ namespace MathNet.Numerics.Threading
/// </summary>
/// <param name="body">Delegate to the task's action.</param>
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) { }
/// <summary>
/// Gets a value indicating whether the task has thrown one or more exceptions while executing.
/// </summary>
@ -67,12 +69,12 @@ namespace MathNet.Numerics.Threading
/// <summary>
/// Gets the exception thrown by the task, if any.
/// </summary>
internal Exception Exception { get; private set; }
protected internal Exception Exception { get; set; }
/// <summary>
/// Run the task.
/// </summary>
internal void Compute()
internal virtual void Compute()
{
try
{
@ -84,4 +86,49 @@ namespace MathNet.Numerics.Threading
}
}
}
/// <summary>
/// Internal Generic Parallel Task Handle.
/// </summary>
internal class Task<T> : Task
{
/// <summary>
/// Delegate to the task's action.
/// </summary>
private readonly Func<T, T> _body;
//private T _initialValue;
public T Result { get; private set; }
/// <summary>
/// Initializes a new instance of the Task class.
/// </summary>
/// <param name="intialValue">The initial value.</param>
/// <param name="body">Delegate to the task's action.</param>
internal Task(T intialValue, Func<T, T> body)
{
if (body == null)
{
throw new ArgumentNullException("body");
}
Result = intialValue;
_body = body;
}
/// <summary>
/// Run the task.
/// </summary>
internal override void Compute()
{
try
{
Result = _body(Result);
}
catch (Exception e)
{
Exception = e;
}
}
}
}

Loading…
Cancel
Save