Skip to content

Commit d494f49

Browse files
committed
perf: worker-pool processing, array hot paths, O(N) streaming
Measured at 100k items with completed-task selectors (pure overhead, net9): - unbounded parallel: 61ms -> 16ms, 444 -> 162 B/item - maxConcurrency 64: 57ms -> 10ms, 431 -> 164 B/item - rate-limited 64: 163ms -> 9ms, 1062 -> 156 B/item Changes: - Rate-limited, timed and throttled parallel processors now run on a fixed pool of P worker loops (WorkerPool) instead of queueing one task per item: P Task.Run calls and one Interlocked increment per item replace N Task.Run tasks, N closures and N semaphore waits. The internal ITaskWrapper interface lets one generic helper serve all four wrapper structs via constrained calls without boxing. - TaskWrappers fields are typed as arrays so iteration uses struct enumerators and LINQ array fast paths instead of interface-dispatched enumerators. - ToIAsyncEnumerable pre-net9 fallback replaced the O(N^2) WhenAny loop (rebuilding the task list every completion) with completion-order buckets: O(N), verified on net8 - 100k results stream in 73ms. - Throttled ProcessInParallel now validates maxConcurrency at build time; previously 0 surfaced as a semaphore ArgumentOutOfRangeException at process time.
1 parent eda8d64 commit d494f49

19 files changed

Lines changed: 148 additions & 225 deletions

EnumerableAsyncProcessor/Extensions/EnumerableExtensions.cs

Lines changed: 20 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -140,38 +140,29 @@ internal static async IAsyncEnumerable<T> ToIAsyncEnumerable<T>(this IEnumerable
140140
yield return task.Result;
141141
}
142142
#else
143-
var managedTasksList = tasks.ToList();
143+
// Interleaving via completion-order buckets: each task's continuation claims the next
144+
// bucket, so streaming N tasks is O(N) rather than the O(N^2) of a WhenAny loop.
145+
var inputTasks = tasks.ToList();
144146

145-
// Create a cancellation task that will complete when cancellation is requested
146-
using var cancellationTcs = new CancellationTokenSource();
147-
var cancellationTask = Task.Delay(Timeout.Infinite, cancellationTcs.Token);
148-
149-
// Register callback to trigger the cancellation task
150-
using var registration = cancellationToken.Register(() => cancellationTcs.Cancel());
147+
var buckets = new TaskCompletionSource<Task<T>>[inputTasks.Count];
148+
for (var i = 0; i < buckets.Length; i++)
149+
{
150+
buckets[i] = new TaskCompletionSource<Task<T>>(TaskCreationOptions.RunContinuationsAsynchronously);
151+
}
151152

152-
while (managedTasksList.Count != 0)
153+
var nextBucketIndex = -1;
154+
foreach (var task in inputTasks)
153155
{
154-
// Check for cancellation before each iteration
155-
cancellationToken.ThrowIfCancellationRequested();
156-
157-
// Include the cancellation task in WhenAny
158-
var allTasks = new List<Task>(managedTasksList.Count + 1);
159-
allTasks.AddRange(managedTasksList);
160-
allTasks.Add(cancellationTask);
161-
162-
var finishedTask = await Task.WhenAny(allTasks).ConfigureAwait(false);
163-
164-
// If the cancellation task completed, throw cancellation
165-
if (finishedTask == cancellationTask)
166-
{
167-
cancellationToken.ThrowIfCancellationRequested();
168-
// This should not happen as cancellation should throw above, but as a safety measure:
169-
throw new OperationCanceledException(cancellationToken);
170-
}
171-
172-
// Remove and yield the completed task
173-
var completedTask = (Task<T>)finishedTask;
174-
managedTasksList.Remove(completedTask);
156+
_ = task.ContinueWith(
157+
completedTask => buckets[Interlocked.Increment(ref nextBucketIndex)].TrySetResult(completedTask),
158+
CancellationToken.None,
159+
TaskContinuationOptions.ExecuteSynchronously,
160+
TaskScheduler.Default);
161+
}
162+
163+
foreach (var bucket in buckets)
164+
{
165+
var completedTask = await bucket.Task.WaitAsync(cancellationToken).ConfigureAwait(false);
175166
yield return await completedTask.ConfigureAwait(false);
176167
}
177168
#endif

EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessor.cs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.Abstract;
44

55
public abstract class AbstractAsyncProcessor : AbstractAsyncProcessorBase
66
{
7-
protected readonly IReadOnlyList<ActionTaskWrapper> TaskWrappers;
7+
protected readonly ActionTaskWrapper[] TaskWrappers;
88

99
private readonly TaskCompletionSource[] _taskCompletionSources;
1010

EnumerableAsyncProcessor/RunnableProcessors/Abstract/AbstractAsyncProcessor_1.cs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.Abstract;
44

55
public abstract class AbstractAsyncProcessor<TInput> : AbstractAsyncProcessorBase
66
{
7-
protected readonly IReadOnlyList<ItemTaskWrapper<TInput>> TaskWrappers;
7+
protected readonly ItemTaskWrapper<TInput>[] TaskWrappers;
88

99
private readonly TaskCompletionSource[] _taskCompletionSources;
1010

Lines changed: 10 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -1,14 +1,20 @@
11
using EnumerableAsyncProcessor.RunnableProcessors.Abstract;
2+
using EnumerableAsyncProcessor.Validation;
23

34
namespace EnumerableAsyncProcessor.RunnableProcessors;
45

56
public class ParallelAsyncProcessor : AbstractAsyncProcessor
67
{
78
private readonly int? _maxConcurrency;
89
private readonly bool _scheduleOnThreadPool;
9-
10+
1011
internal ParallelAsyncProcessor(int count, Func<Task> taskSelector, CancellationTokenSource cancellationTokenSource, int? maxConcurrency = null, bool scheduleOnThreadPool = false) : base(count, taskSelector, cancellationTokenSource)
1112
{
13+
if (maxConcurrency is { } concurrencyLimit)
14+
{
15+
ValidationHelper.ThrowIfNegativeOrZero(concurrencyLimit, nameof(maxConcurrency));
16+
}
17+
1218
_maxConcurrency = maxConcurrency;
1319
_scheduleOnThreadPool = scheduleOnThreadPool;
1420
}
@@ -38,38 +44,8 @@ await Task.WhenAll(TaskWrappers.Select(taskWrapper =>
3844
return;
3945
}
4046

41-
// Use semaphore for concurrency throttling
42-
using var semaphore = new SemaphoreSlim(_maxConcurrency.Value, _maxConcurrency.Value);
43-
44-
// Materialize tasks immediately to ensure they all start in parallel (up to concurrency limit)
45-
var tasks = _scheduleOnThreadPool
46-
? // Use Task.Run to prevent synchronous code from blocking thread pool threads
47-
TaskWrappers.Select(taskWrapper => Task.Run(async () =>
48-
{
49-
await semaphore.WaitAsync(CancellationToken).ConfigureAwait(false);
50-
try
51-
{
52-
await taskWrapper.Process(CancellationToken).ConfigureAwait(false);
53-
}
54-
finally
55-
{
56-
semaphore.Release();
57-
}
58-
}, CancellationToken)).ToList()
59-
: // Direct execution for maximum performance
60-
TaskWrappers.Select(async taskWrapper =>
61-
{
62-
await semaphore.WaitAsync(CancellationToken).ConfigureAwait(false);
63-
try
64-
{
65-
await taskWrapper.Process(CancellationToken).ConfigureAwait(false);
66-
}
67-
finally
68-
{
69-
semaphore.Release();
70-
}
71-
}).ToList(); // Force immediate task creation
72-
73-
await Task.WhenAll(tasks).ConfigureAwait(false);
47+
// Throttled processing runs on a fixed worker pool: P worker tasks instead of
48+
// one queued task, closure and semaphore wait per item
49+
await WorkerPool.ProcessAsync(TaskWrappers, _maxConcurrency.Value, minimumIterationTime: null, CancellationToken).ConfigureAwait(false);
7450
}
7551
}
Lines changed: 10 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -1,14 +1,20 @@
11
using EnumerableAsyncProcessor.RunnableProcessors.Abstract;
2+
using EnumerableAsyncProcessor.Validation;
23

34
namespace EnumerableAsyncProcessor.RunnableProcessors;
45

56
public class ParallelAsyncProcessor<TInput> : AbstractAsyncProcessor<TInput>
67
{
78
private readonly int? _maxConcurrency;
89
private readonly bool _scheduleOnThreadPool;
9-
10+
1011
internal ParallelAsyncProcessor(IEnumerable<TInput> items, Func<TInput, Task> taskSelector, CancellationTokenSource cancellationTokenSource, int? maxConcurrency = null, bool scheduleOnThreadPool = false) : base(items, taskSelector, cancellationTokenSource)
1112
{
13+
if (maxConcurrency is { } concurrencyLimit)
14+
{
15+
ValidationHelper.ThrowIfNegativeOrZero(concurrencyLimit, nameof(maxConcurrency));
16+
}
17+
1218
_maxConcurrency = maxConcurrency;
1319
_scheduleOnThreadPool = scheduleOnThreadPool;
1420
}
@@ -38,38 +44,8 @@ await Task.WhenAll(TaskWrappers.Select(taskWrapper =>
3844
return;
3945
}
4046

41-
// Use semaphore for concurrency throttling
42-
using var semaphore = new SemaphoreSlim(_maxConcurrency.Value, _maxConcurrency.Value);
43-
44-
// Materialize tasks immediately to ensure they all start in parallel (up to concurrency limit)
45-
var tasks = _scheduleOnThreadPool
46-
? // Use Task.Run to prevent synchronous code from blocking thread pool threads
47-
TaskWrappers.Select(taskWrapper => Task.Run(async () =>
48-
{
49-
await semaphore.WaitAsync(CancellationToken).ConfigureAwait(false);
50-
try
51-
{
52-
await taskWrapper.Process(CancellationToken).ConfigureAwait(false);
53-
}
54-
finally
55-
{
56-
semaphore.Release();
57-
}
58-
}, CancellationToken)).ToList()
59-
: // Direct execution for maximum performance
60-
TaskWrappers.Select(async taskWrapper =>
61-
{
62-
await semaphore.WaitAsync(CancellationToken).ConfigureAwait(false);
63-
try
64-
{
65-
await taskWrapper.Process(CancellationToken).ConfigureAwait(false);
66-
}
67-
finally
68-
{
69-
semaphore.Release();
70-
}
71-
}).ToList(); // Force immediate task creation
72-
73-
await Task.WhenAll(tasks).ConfigureAwait(false);
47+
// Throttled processing runs on a fixed worker pool: P worker tasks instead of
48+
// one queued task, closure and semaphore wait per item
49+
await WorkerPool.ProcessAsync(TaskWrappers, _maxConcurrency.Value, minimumIterationTime: null, CancellationToken).ConfigureAwait(false);
7450
}
7551
}

EnumerableAsyncProcessor/RunnableProcessors/RateLimitedParallelAsyncProcessor.cs

Lines changed: 1 addition & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,3 @@
1-
using EnumerableAsyncProcessor.Extensions;
21
using EnumerableAsyncProcessor.RunnableProcessors.Abstract;
32
using EnumerableAsyncProcessor.Validation;
43

@@ -17,9 +16,6 @@ internal RateLimitedParallelAsyncProcessor(int count, Func<Task> taskSelector, i
1716

1817
internal override Task Process()
1918
{
20-
// Task.Run guards the shared worker slots against synchronous code in user delegates
21-
return TaskWrappers.InParallelAsync(_levelsOfParallelism,
22-
taskWrapper => Task.Run(() => taskWrapper.Process(CancellationToken), CancellationToken),
23-
CancellationToken);
19+
return WorkerPool.ProcessAsync(TaskWrappers, _levelsOfParallelism, minimumIterationTime: null, CancellationToken);
2420
}
2521
}

EnumerableAsyncProcessor/RunnableProcessors/RateLimitedParallelAsyncProcessor_1.cs

Lines changed: 1 addition & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,3 @@
1-
using EnumerableAsyncProcessor.Extensions;
21
using EnumerableAsyncProcessor.RunnableProcessors.Abstract;
32
using EnumerableAsyncProcessor.Validation;
43

@@ -17,9 +16,6 @@ internal RateLimitedParallelAsyncProcessor(IEnumerable<TInput> items, Func<TInpu
1716

1817
internal override Task Process()
1918
{
20-
// Task.Run guards the shared worker slots against synchronous code in user delegates
21-
return TaskWrappers.InParallelAsync(_levelsOfParallelism,
22-
taskWrapper => Task.Run(() => taskWrapper.Process(CancellationToken), CancellationToken),
23-
CancellationToken);
19+
return WorkerPool.ProcessAsync(TaskWrappers, _levelsOfParallelism, minimumIterationTime: null, CancellationToken);
2420
}
2521
}

EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessor_1.cs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract;
44

55
public abstract class ResultAbstractAsyncProcessor<TOutput> : ResultAbstractAsyncProcessorBase<TOutput>
66
{
7-
protected readonly IReadOnlyList<ActionTaskWrapper<TOutput>> TaskWrappers;
7+
protected readonly ActionTaskWrapper<TOutput>[] TaskWrappers;
88

99
private readonly TaskCompletionSource<TOutput>[] _taskCompletionSources;
1010

EnumerableAsyncProcessor/RunnableProcessors/ResultProcessors/Abstract/ResultAbstractAsyncProcessor_2.cs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@ namespace EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract;
44

55
public abstract class ResultAbstractAsyncProcessor<TInput, TOutput> : ResultAbstractAsyncProcessorBase<TOutput>
66
{
7-
protected readonly IReadOnlyList<ItemTaskWrapper<TInput, TOutput>> TaskWrappers;
7+
protected readonly ItemTaskWrapper<TInput, TOutput>[] TaskWrappers;
88

99
private readonly TaskCompletionSource<TOutput>[] _taskCompletionSources;
1010

Lines changed: 10 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -1,14 +1,20 @@
11
using EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors.Abstract;
2+
using EnumerableAsyncProcessor.Validation;
23

34
namespace EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors;
45

56
public class ResultParallelAsyncProcessor<TOutput> : ResultAbstractAsyncProcessor<TOutput>
67
{
78
private readonly int? _maxConcurrency;
89
private readonly bool _scheduleOnThreadPool;
9-
10+
1011
internal ResultParallelAsyncProcessor(int count, Func<Task<TOutput>> taskSelector, CancellationTokenSource cancellationTokenSource, int? maxConcurrency = null, bool scheduleOnThreadPool = false) : base(count, taskSelector, cancellationTokenSource)
1112
{
13+
if (maxConcurrency is { } concurrencyLimit)
14+
{
15+
ValidationHelper.ThrowIfNegativeOrZero(concurrencyLimit, nameof(maxConcurrency));
16+
}
17+
1218
_maxConcurrency = maxConcurrency;
1319
_scheduleOnThreadPool = scheduleOnThreadPool;
1420
}
@@ -38,38 +44,8 @@ await Task.WhenAll(TaskWrappers.Select(taskWrapper =>
3844
return;
3945
}
4046

41-
// Use semaphore for concurrency throttling
42-
using var semaphore = new SemaphoreSlim(_maxConcurrency.Value, _maxConcurrency.Value);
43-
44-
// Materialize tasks immediately to ensure they all start in parallel (up to concurrency limit)
45-
var tasks = _scheduleOnThreadPool
46-
? // Use Task.Run to prevent synchronous code from blocking thread pool threads
47-
TaskWrappers.Select(taskWrapper => Task.Run(async () =>
48-
{
49-
await semaphore.WaitAsync(CancellationToken).ConfigureAwait(false);
50-
try
51-
{
52-
await taskWrapper.Process(CancellationToken).ConfigureAwait(false);
53-
}
54-
finally
55-
{
56-
semaphore.Release();
57-
}
58-
}, CancellationToken)).ToList()
59-
: // Direct execution for maximum performance
60-
TaskWrappers.Select(async taskWrapper =>
61-
{
62-
await semaphore.WaitAsync(CancellationToken).ConfigureAwait(false);
63-
try
64-
{
65-
await taskWrapper.Process(CancellationToken).ConfigureAwait(false);
66-
}
67-
finally
68-
{
69-
semaphore.Release();
70-
}
71-
}).ToList(); // Force immediate task creation
72-
73-
await Task.WhenAll(tasks).ConfigureAwait(false);
47+
// Throttled processing runs on a fixed worker pool: P worker tasks instead of
48+
// one queued task, closure and semaphore wait per item
49+
await WorkerPool.ProcessAsync(TaskWrappers, _maxConcurrency.Value, minimumIterationTime: null, CancellationToken).ConfigureAwait(false);
7450
}
7551
}

0 commit comments

Comments
 (0)