Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
139 changes: 139 additions & 0 deletions EnumerableAsyncProcessor.UnitTests/AsyncEnumerableProcessorTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -215,6 +215,145 @@ public async Task ForEachAsync_WithException_PropagatesException()
var exception = await Assert.ThrowsAsync<InvalidOperationException>(async () => await task);
await Assert.That(exception!.Message).IsEqualTo("Test exception");
}

[Test]
public async Task ForEachAsync_ProcessInParallel_UnboundedConcurrency_ProcessesAllItems()
{
var processedItems = new List<int>();
var asyncEnumerable = GenerateAsyncEnumerable(50);

await asyncEnumerable
.ForEachAsync(async item =>
{
await Task.Delay(5);
lock (processedItems)
{
processedItems.Add(item);
}
})
.ProcessInParallel() // Unbounded concurrency
.ExecuteAsync();

await Assert.That(processedItems.Count).IsEqualTo(50);
await Assert.That(processedItems.OrderBy(x => x)).IsEquivalentTo(Enumerable.Range(1, 50));
}

[Test]
public async Task SelectAsync_ProcessInParallel_UnboundedConcurrency_ReturnsAllResults()
{
var asyncEnumerable = GenerateAsyncEnumerable(30);

var results = await asyncEnumerable
.SelectAsync(async item =>
{
await Task.Delay(5);
return item * 2;
})
.ProcessInParallel() // Unbounded concurrency
.ExecuteAsync()
.ToListAsync();

await Assert.That(results.Count).IsEqualTo(30);
await Assert.That(results.OrderBy(x => x)).IsEquivalentTo(Enumerable.Range(1, 30).Select(x => x * 2));
}

[Test]
public async Task ForEachAsync_ProcessInParallel_WithThreadPoolScheduling_ProcessesAllItems()
{
var processedItems = new List<int>();
var asyncEnumerable = GenerateAsyncEnumerable(20);

await asyncEnumerable
.ForEachAsync(async item =>
{
await Task.Delay(5);
lock (processedItems)
{
processedItems.Add(item);
}
})
.ProcessInParallel(scheduleOnThreadPool: true)
.ExecuteAsync();

await Assert.That(processedItems.Count).IsEqualTo(20);
await Assert.That(processedItems.OrderBy(x => x)).IsEquivalentTo(Enumerable.Range(1, 20));
}

[Test]
public async Task ForEachAsync_ProcessInBatches_ProcessesAllItemsInBatches()
{
var processedBatches = new List<int>();
var asyncEnumerable = GenerateAsyncEnumerable(25);

await asyncEnumerable
.ForEachAsync(async item =>
{
await Task.Delay(5);
lock (processedBatches)
{
processedBatches.Add(item);
}
})
.ProcessInBatches(5)
.ExecuteAsync();

await Assert.That(processedBatches.Count).IsEqualTo(25);
await Assert.That(processedBatches.OrderBy(x => x)).IsEquivalentTo(Enumerable.Range(1, 25));
}

[Test]
public async Task SelectAsync_ProcessInBatches_ReturnsAllResultsInBatches()
{
var asyncEnumerable = GenerateAsyncEnumerable(23);

var results = await asyncEnumerable
.SelectAsync(async item =>
{
await Task.Delay(5);
return item * 3;
})
.ProcessInBatches(5)
.ExecuteAsync()
.ToListAsync();

await Assert.That(results.Count).IsEqualTo(23);
// Batches maintain order within batch, so results should be in order
await Assert.That(results).IsEquivalentTo(Enumerable.Range(1, 23).Select(x => x * 3));
}

[Test]
public async Task ProcessInParallel_NullableConcurrency_WorksCorrectly()
{
var asyncEnumerable = GenerateAsyncEnumerable(15);
var processedCount = 0;

// Test with null concurrency (unbounded)
await asyncEnumerable
.ForEachAsync(async item =>
{
await Task.Delay(5);
Interlocked.Increment(ref processedCount);
})
.ProcessInParallel((int?)null)
.ExecuteAsync();

await Assert.That(processedCount).IsEqualTo(15);

// Reset and test with specified concurrency
processedCount = 0;
asyncEnumerable = GenerateAsyncEnumerable(15);

await asyncEnumerable
.ForEachAsync(async item =>
{
await Task.Delay(5);
Interlocked.Increment(ref processedCount);
})
.ProcessInParallel((int?)5)
.ExecuteAsync();

await Assert.That(processedCount).IsEqualTo(15);
}
}

internal static class AsyncEnumerableExtensionsForTests
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,20 +21,54 @@ public AsyncEnumerableActionAsyncProcessorBuilder(
}

/// <summary>
/// Process items in parallel with a specified level of parallelism.
/// Process items in parallel without concurrency limits.
/// </summary>
/// <returns>An async processor configured for parallel execution.</returns>
public IAsyncEnumerableProcessor ProcessInParallel()
{
return ProcessInParallel(null, false);
}

/// <summary>
/// Process items in parallel without concurrency limits.
/// </summary>
/// <param name="scheduleOnThreadPool">If true, schedules tasks on thread pool to prevent blocking. Default is false for maximum performance.</param>
/// <returns>An async processor configured for parallel execution.</returns>
public IAsyncEnumerableProcessor ProcessInParallel(bool scheduleOnThreadPool)
{
return ProcessInParallel(null, scheduleOnThreadPool);
}

/// <summary>
/// Process items in parallel with specified concurrency limit.
/// </summary>
/// <param name="maxConcurrency">Maximum concurrent operations.</param>
/// <returns>An async processor configured for parallel execution.</returns>
public IAsyncEnumerableProcessor ProcessInParallel(int maxConcurrency)
{
return new AsyncEnumerableParallelProcessor<TInput>(
_items, _taskSelector, maxConcurrency, _cancellationTokenSource);
return ProcessInParallel((int?)maxConcurrency, false);
}

/// <summary>
/// Process items in parallel with default concurrency (processor count).
/// Process items in parallel with specified concurrency limit.
/// </summary>
public IAsyncEnumerableProcessor ProcessInParallel()
/// <param name="maxConcurrency">Maximum concurrent operations.</param>
/// <returns>An async processor configured for parallel execution.</returns>
public IAsyncEnumerableProcessor ProcessInParallel(int? maxConcurrency)
{
return ProcessInParallel(maxConcurrency, false);
}

/// <summary>
/// Process items in parallel with specified concurrency limit.
/// </summary>
/// <param name="maxConcurrency">Maximum concurrent operations.</param>
/// <param name="scheduleOnThreadPool">If true, schedules tasks on thread pool to prevent blocking.</param>
/// <returns>An async processor configured for parallel execution.</returns>
public IAsyncEnumerableProcessor ProcessInParallel(int? maxConcurrency, bool scheduleOnThreadPool)
{
return ProcessInParallel(Environment.ProcessorCount);
return new AsyncEnumerableParallelProcessor<TInput>(
_items, _taskSelector, maxConcurrency, scheduleOnThreadPool, _cancellationTokenSource);
}


Expand All @@ -46,6 +80,17 @@ public IAsyncEnumerableProcessor ProcessOneAtATime()
return new AsyncEnumerableOneAtATimeProcessor<TInput>(
_items, _taskSelector, _cancellationTokenSource);
}

/// <summary>
/// Process items in batches.
/// </summary>
/// <param name="batchSize">The size of each batch.</param>
/// <returns>An async processor configured for batch execution.</returns>
public IAsyncEnumerableProcessor ProcessInBatches(int batchSize)
{
return new AsyncEnumerableBatchProcessor<TInput>(
_items, _taskSelector, batchSize, _cancellationTokenSource);
}

}
#endif
Original file line number Diff line number Diff line change
Expand Up @@ -21,20 +21,54 @@ public AsyncEnumerableActionAsyncProcessorBuilder(
}

/// <summary>
/// Process items in parallel with a specified level of parallelism and return results.
/// Process items in parallel without concurrency limits and return results.
/// </summary>
/// <returns>An async processor configured for parallel execution.</returns>
public IAsyncEnumerableProcessor<TOutput> ProcessInParallel()
{
return ProcessInParallel(null, false);
}

/// <summary>
/// Process items in parallel without concurrency limits and return results.
/// </summary>
/// <param name="scheduleOnThreadPool">If true, schedules tasks on thread pool to prevent blocking. Default is false for maximum performance.</param>
/// <returns>An async processor configured for parallel execution.</returns>
public IAsyncEnumerableProcessor<TOutput> ProcessInParallel(bool scheduleOnThreadPool)
{
return ProcessInParallel(null, scheduleOnThreadPool);
}

/// <summary>
/// Process items in parallel with specified concurrency limit and return results.
/// </summary>
/// <param name="maxConcurrency">Maximum concurrent operations.</param>
/// <returns>An async processor configured for parallel execution.</returns>
public IAsyncEnumerableProcessor<TOutput> ProcessInParallel(int maxConcurrency)
{
return new ResultAsyncEnumerableParallelProcessor<TInput, TOutput>(
_items, _taskSelector, maxConcurrency, _cancellationTokenSource);
return ProcessInParallel((int?)maxConcurrency, false);
}

/// <summary>
/// Process items in parallel with default concurrency and return results.
/// Process items in parallel with specified concurrency limit and return results.
/// </summary>
public IAsyncEnumerableProcessor<TOutput> ProcessInParallel()
/// <param name="maxConcurrency">Maximum concurrent operations.</param>
/// <returns>An async processor configured for parallel execution.</returns>
public IAsyncEnumerableProcessor<TOutput> ProcessInParallel(int? maxConcurrency)
{
return ProcessInParallel(maxConcurrency, false);
}

/// <summary>
/// Process items in parallel with specified concurrency limit and return results.
/// </summary>
/// <param name="maxConcurrency">Maximum concurrent operations.</param>
/// <param name="scheduleOnThreadPool">If true, schedules tasks on thread pool to prevent blocking.</param>
/// <returns>An async processor configured for parallel execution.</returns>
public IAsyncEnumerableProcessor<TOutput> ProcessInParallel(int? maxConcurrency, bool scheduleOnThreadPool)
{
return ProcessInParallel(Environment.ProcessorCount);
return new ResultAsyncEnumerableParallelProcessor<TInput, TOutput>(
_items, _taskSelector, maxConcurrency, scheduleOnThreadPool, _cancellationTokenSource);
}


Expand All @@ -46,6 +80,17 @@ public IAsyncEnumerableProcessor<TOutput> ProcessOneAtATime()
return new ResultAsyncEnumerableOneAtATimeProcessor<TInput, TOutput>(
_items, _taskSelector, _cancellationTokenSource);
}

/// <summary>
/// Process items in batches and return results.
/// </summary>
/// <param name="batchSize">The size of each batch.</param>
/// <returns>An async processor configured for batch execution.</returns>
public IAsyncEnumerableProcessor<TOutput> ProcessInBatches(int batchSize)
{
return new ResultAsyncEnumerableBatchProcessor<TInput, TOutput>(
_items, _taskSelector, batchSize, _cancellationTokenSource);
}

}
#endif
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
#if NET6_0_OR_GREATER
using EnumerableAsyncProcessor.Extensions;

namespace EnumerableAsyncProcessor.RunnableProcessors.AsyncEnumerable;

public class AsyncEnumerableBatchProcessor<TInput> : IAsyncEnumerableProcessor
{
private readonly IAsyncEnumerable<TInput> _items;
private readonly Func<TInput, Task> _taskSelector;
private readonly int _batchSize;
private readonly CancellationTokenSource _cancellationTokenSource;

internal AsyncEnumerableBatchProcessor(
IAsyncEnumerable<TInput> items,
Func<TInput, Task> taskSelector,
int batchSize,
CancellationTokenSource cancellationTokenSource)
{
_items = items;
_taskSelector = taskSelector;
_batchSize = batchSize;
_cancellationTokenSource = cancellationTokenSource;
}

public async Task ExecuteAsync()
{
var cancellationToken = _cancellationTokenSource.Token;
var batch = new List<TInput>(_batchSize);

await foreach (var item in _items.WithCancellation(cancellationToken).ConfigureAwait(false))
{
batch.Add(item);

if (batch.Count >= _batchSize)
{
await ProcessBatch(batch, cancellationToken).ConfigureAwait(false);
batch = new List<TInput>(_batchSize);
}
}

// Process any remaining items in the final batch
if (batch.Count > 0)
{
await ProcessBatch(batch, cancellationToken).ConfigureAwait(false);
}
}

private async Task ProcessBatch(List<TInput> batch, CancellationToken cancellationToken)
{
var tasks = batch.Select(item => _taskSelector(item)).ToArray();
await Task.WhenAll(tasks).ConfigureAwait(false);
}
}
#endif
Loading
Loading