Fix for ForEach-Object -Parallel performance problem with many runspaces (#10455)

This commit is contained in:
Paul Higinbotham
2019-09-05 11:27:20 -07:00
committed by Aditya Patwardhan
parent fbf4f6c11b
commit cfdbd71888
3 changed files with 148 additions and 83 deletions
@@ -70,6 +70,8 @@ namespace Microsoft.PowerShell.Commands
#endregion
#region Common Parameters
/// <summary>
/// This parameter specifies the current pipeline object.
/// </summary>
@@ -85,6 +87,8 @@ namespace Microsoft.PowerShell.Commands
private PSObject _inputObject = AutomationNull.Value;
#endregion
#region ScriptBlockSet
private List<ScriptBlock> _scripts = new List<ScriptBlock>();
@@ -355,6 +359,7 @@ namespace Microsoft.PowerShell.Commands
_taskTimer?.Dispose();
_taskDataStreamWriter?.Dispose();
_taskPool?.Dispose();
_taskCollection?.Dispose();
}
#endregion
@@ -368,6 +373,8 @@ namespace Microsoft.PowerShell.Commands
private Dictionary<string, object> _usingValuesMap;
private Timer _taskTimer;
private PSTaskJob _taskJob;
private PSDataCollection<System.Management.Automation.PSTasks.PSTask> _taskCollection;
private Exception _taskCollectionException;
private void InitParallelParameterSet()
{
@@ -411,6 +418,7 @@ namespace Microsoft.PowerShell.Commands
if (AsJob)
{
// Set up for returning a job object.
if (MyInvocation.BoundParameters.ContainsKey(nameof(TimeoutSeconds)))
{
ThrowTerminatingError(
@@ -424,24 +432,78 @@ namespace Microsoft.PowerShell.Commands
_taskJob = new PSTaskJob(
Parallel.ToString(),
ThrottleLimit);
return;
}
else
// Set up for synchronous processing and data streaming.
_taskCollection = new PSDataCollection<System.Management.Automation.PSTasks.PSTask>();
_taskDataStreamWriter = new PSTaskDataStreamWriter(this);
_taskPool = new PSTaskPool(ThrottleLimit);
_taskPool.PoolComplete += (sender, args) =>
{
_taskDataStreamWriter = new PSTaskDataStreamWriter(this);
_taskPool = new PSTaskPool(ThrottleLimit);
_taskPool.PoolComplete += (sender, args) =>
{
_taskDataStreamWriter.Close();
};
if (TimeoutSeconds != 0)
{
_taskTimer = new Timer(
(_) => _taskPool.StopAll(),
null,
TimeoutSeconds * 1000,
Timeout.Infinite);
}
_taskDataStreamWriter.Close();
};
// Create timeout timer if requested.
if (TimeoutSeconds != 0)
{
_taskTimer = new Timer(
callback: (_) => { _taskCollection.Complete(); _taskPool.StopAll(); },
state: null,
dueTime: TimeoutSeconds * 1000,
period: Timeout.Infinite);
}
// Task collection handler.
System.Threading.ThreadPool.QueueUserWorkItem(
(_) =>
{
// As piped input are converted to PSTasks and added to the _taskCollection,
// transfer the task to the _taskPool on this dedicated thread.
// The _taskPool will block this thread when it is full, and allow more tasks to
// be added only when a currently running task completes and makes space in the pool.
// Continue adding any tasks appearing in _taskCollection until the collection is closed.
while (true)
{
// This handle will unblock the thread when a new task is available or the _taskCollection
// is closed.
_taskCollection.WaitHandle.WaitOne();
// Task collection open state is volatile.
// Record current task collection open state here, to be checked after processing.
bool isOpen = _taskCollection.IsOpen;
try
{
// Read all tasks in the collection.
foreach (var task in _taskCollection.ReadAll())
{
// This _taskPool method will block if the pool is full and will unblock
// only after a task completes making more space.
_taskPool.Add(task);
}
}
catch (Exception ex)
{
// Close the _taskCollection on an unexpected exception so the pool closes and
// lets any running tasks complete.
_taskCollection.Complete();
_taskCollectionException = ex;
break;
}
// Loop is exited only when task collection is closed and all task
// collection tasks are processed.
if (!isOpen)
{
break;
}
}
// We are done adding tasks and can close the task pool.
_taskPool.Close();
});
}
private void ProcessParallelParameterSet()
@@ -461,27 +523,39 @@ namespace Microsoft.PowerShell.Commands
if (AsJob)
{
// Add child task job.
var taskChildJob = new PSTaskChildJob(
Parallel,
_usingValuesMap,
InputObject);
_taskJob.AddJob(taskChildJob);
return;
}
else
// Write any streaming data
_taskDataStreamWriter.WriteImmediate();
// Add to task collection for processing.
if (_taskCollection.IsOpen)
{
// Write any streaming data
_taskDataStreamWriter.WriteImmediate();
var task = new System.Management.Automation.PSTasks.PSTask(
Parallel,
_usingValuesMap,
InputObject,
_taskDataStreamWriter);
// Add task to task pool.
// Block if the pool is full and wait until task can be added.
_taskPool.Add(task, _taskDataStreamWriter);
try
{
// Create a PSTask based on this piped input and add it to the task collection.
// A dedicated thread will add it to the PSTask pool in a performant manner.
_taskCollection.Add(
new System.Management.Automation.PSTasks.PSTask(
Parallel,
_usingValuesMap,
InputObject,
_taskDataStreamWriter));
}
catch (InvalidOperationException)
{
// This exception is thrown if the task collection is closed, which should not happen.
Dbg.Assert(false, "Should not add to a closed PSTask collection");
}
}
}
@@ -489,25 +563,37 @@ namespace Microsoft.PowerShell.Commands
{
if (AsJob)
{
// Start and return parent job object.
_taskJob.Start();
JobRepository.Add(_taskJob);
WriteObject(_taskJob);
}
else
{
_taskDataStreamWriter.WriteImmediate();
_taskPool.Close();
_taskDataStreamWriter.WaitAndWrite();
return;
}
// Close task collection and wait for processing to complete while streaming data.
_taskDataStreamWriter.WriteImmediate();
_taskCollection.Complete();
_taskDataStreamWriter.WaitAndWrite();
// Check for an unexpected error from the _taskCollection handler thread and report here.
var ex = _taskCollectionException;
if (ex != null)
{
var msg = string.Format(CultureInfo.InvariantCulture, InternalCommandStrings.ParallelPipedInputProcessingError, ex);
WriteError(
new ErrorRecord(
exception: new InvalidOperationException(msg),
errorId: "ParallelPipedInputProcessingError",
errorCategory: ErrorCategory.InvalidOperation,
targetObject: this));
}
}
private void StopParallelProcessing()
{
if (!AsJob)
{
_taskPool.StopAll();
}
_taskCollection?.Complete();
_taskPool?.StopAll();
}
#endregion
@@ -602,12 +602,16 @@ namespace System.Management.Automation.PSTasks
#region Members
private readonly ManualResetEvent _addAvailable;
private readonly ManualResetEvent _stopAll;
private readonly Dictionary<int, PSTaskBase> _taskPool;
private readonly int _sizeLimit;
private readonly ManualResetEvent _stopAll;
private readonly object _syncObject;
private readonly Dictionary<int, PSTaskBase> _taskPool;
private readonly WaitHandle[] _waitHandles;
private bool _isOpen;
private const int AddAvailable = 0;
private const int Stop = 1;
#endregion
#region Constructor
@@ -625,6 +629,11 @@ namespace System.Management.Automation.PSTasks
_syncObject = new object();
_addAvailable = new ManualResetEvent(true);
_stopAll = new ManualResetEvent(false);
_waitHandles = new WaitHandle[]
{
_addAvailable, // index 0
_stopAll, // index 1
};
_taskPool = new Dictionary<int, PSTaskBase>(size);
}
@@ -672,46 +681,21 @@ namespace System.Management.Automation.PSTasks
/// This method is not multi-thread safe and assumes only one thread waits and adds tasks.
/// </summary>
/// <param name="task">Task to be added to pool.</param>
/// <param name="dataStreamWriter">Optional cmdlet data stream writer.</param>
/// <returns>True when task is successfully added.</returns>
public bool Add(
PSTaskBase task,
PSTaskDataStreamWriter dataStreamWriter = null)
public bool Add(PSTaskBase task)
{
if (!_isOpen)
{
return false;
}
WaitHandle[] waitHandles;
if (dataStreamWriter != null)
{
waitHandles = new WaitHandle[]
{
_addAvailable, // index 0
_stopAll, // index 1
dataStreamWriter.DataAddedWaitHandle // index 2
};
}
else
{
waitHandles = new WaitHandle[]
{
_addAvailable, // index 0
_stopAll, // index 1
};
}
// Block until either space is available, or a stop is commanded
var index = WaitHandle.WaitAny(_waitHandles);
// Block until either room is available, data is ready for writing, or a stop command
while (true)
switch (index)
{
var index = WaitHandle.WaitAny(waitHandles);
// Add new task
if (index == 0)
{
case AddAvailable:
task.StateChanged += HandleTaskStateChangedDelegate;
lock (_syncObject)
{
if (!_isOpen)
@@ -727,21 +711,13 @@ namespace System.Management.Automation.PSTasks
task.Start();
}
return true;
}
// Stop all
if (index == 1)
{
case Stop:
return false;
default:
return false;
}
// Data ready for writing
if (index == 2)
{
dataStreamWriter.WriteImmediate();
}
}
}
@@ -977,7 +953,7 @@ namespace System.Management.Automation.PSTasks
// This thread will end once all jobs reach a finished state by either running
// to completion, terminating with error, or stopped.
System.Threading.ThreadPool.QueueUserWorkItem(
(state) =>
(_) =>
{
foreach (var childJob in ChildJobs)
{
@@ -181,4 +181,7 @@
<value>The following common parameters are not currently supported in the Parallel parameter set:
ErrorAction, WarningAction, InformationAction, PipelineVariable</value>
</data>
<data name="ParallelPipedInputProcessingError" xml:space="preserve">
<value>An unexpected error has occurred while processing ForEach-Object -Parallel input. This may mean that some of the piped input did not get processed. Error: {0}.</value>
</data>
</root>