-
-
Notifications
You must be signed in to change notification settings - Fork 4
Expand file tree
/
Copy pathActionAsyncProcessorBuilder_1.cs
More file actions
88 lines (76 loc) · 4.33 KB
/
Copy pathActionAsyncProcessorBuilder_1.cs
File metadata and controls
88 lines (76 loc) · 4.33 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
using EnumerableAsyncProcessor.Interfaces;
using EnumerableAsyncProcessor.RunnableProcessors.ResultProcessors;
using EnumerableAsyncProcessor.Extensions;
namespace EnumerableAsyncProcessor.Builders;
public class ActionAsyncProcessorBuilder<TOutput>
{
private readonly int _count;
private readonly Func<Task<TOutput>> _taskSelector;
private readonly CancellationTokenSource _cancellationTokenSource;
internal ActionAsyncProcessorBuilder(int count, Func<Task<TOutput>> taskSelector, CancellationToken cancellationToken)
{
_count = count;
_taskSelector = taskSelector;
_cancellationTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
}
public IAsyncProcessor<TOutput> ProcessInBatches(int batchSize)
{
return new ResultBatchAsyncProcessor<TOutput>(batchSize, _count, _taskSelector, _cancellationTokenSource).StartProcessing();
}
public IAsyncProcessor<TOutput> ProcessInParallel(int levelOfParallelism)
{
return new ResultRateLimitedParallelAsyncProcessor<TOutput>(_count, _taskSelector, levelOfParallelism, _cancellationTokenSource).StartProcessing();
}
public IAsyncProcessor<TOutput> ProcessInParallel(int levelOfParallelism, TimeSpan timeSpan)
{
return new ResultTimedRateLimitedParallelAsyncProcessor<TOutput>(_count, _taskSelector, levelOfParallelism, timeSpan, _cancellationTokenSource).StartProcessing();
}
public IAsyncProcessor<TOutput> ProcessInParallel()
{
return new ResultParallelAsyncProcessor<TOutput>(_count, _taskSelector, _cancellationTokenSource).StartProcessing();
}
/// <summary>
/// Process tasks in parallel with optimizations for I/O-bound operations and return results.
/// Removes Task.Run overhead and allows higher concurrency levels.
/// </summary>
/// <param name="maxConcurrency">Maximum concurrent operations. If null, defaults to 10x processor count or minimum 100 for I/O-bound tasks.</param>
/// <returns>An async processor optimized for I/O operations that returns results.</returns>
public IAsyncProcessor<TOutput> ProcessInParallelForIO(int? maxConcurrency = null)
{
return new ResultIOBoundParallelAsyncProcessor<TOutput>(_count, _taskSelector, _cancellationTokenSource, maxConcurrency).StartProcessing();
}
/// <summary>
/// Process tasks in parallel with explicit I/O vs CPU-bound configuration and return results.
/// </summary>
/// <param name="isIOBound">True for I/O-bound tasks (removes Task.Run overhead), false for CPU-bound tasks.</param>
/// <returns>An async processor configured for the specified workload type that returns results.</returns>
public IAsyncProcessor<TOutput> ProcessInParallel(bool isIOBound)
{
return new ResultParallelAsyncProcessor<TOutput>(_count, _taskSelector, _cancellationTokenSource, isIOBound).StartProcessing();
}
public IAsyncProcessor<TOutput> ProcessOneAtATime()
{
return new ResultOneAtATimeAsyncProcessor<TOutput>(_count, _taskSelector, _cancellationTokenSource).StartProcessing();
}
/// <summary>
/// Process ALL tasks in parallel without any concurrency limits and return results.
/// WARNING: Use with caution - can overwhelm system resources with large task counts.
/// Ideal for scenarios requiring maximum parallelism like running thousands of unit tests.
/// </summary>
/// <returns>An async processor with unbounded parallelism that returns results.</returns>
public IAsyncProcessor<TOutput> ProcessInParallelUnbounded()
{
return new ResultUnboundedParallelAsyncProcessor<TOutput>(_count, _taskSelector, _cancellationTokenSource).StartProcessing();
}
#if NET6_0_OR_GREATER
/// <summary>
/// Process tasks using a channel-based approach with producer-consumer pattern and return results.
/// </summary>
/// <param name="options">Channel configuration options. If null, uses unbounded channel with single consumer.</param>
/// <returns>An async processor that processes tasks through a channel and returns results.</returns>
public IAsyncProcessor<TOutput> ProcessWithChannel(ChannelProcessorOptions? options = null)
{
return new ResultChannelBasedBatchAsyncProcessor<TOutput>(_count, _taskSelector, _cancellationTokenSource, options).StartProcessing();
}
#endif
}