Skip to content

Commit 0694eab

Browse files
committed
Added quick implementation on top of Parallel.ForEachAsync, but left it disabled for now as benchmark results are mixed
1 parent ec3134e commit 0694eab

2 files changed

Lines changed: 104 additions & 1 deletion

File tree

src/CSRakowski.Parallel/ParallelAsync.Unordered.cs

Lines changed: 100 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,34 @@ private static async Task<IEnumerable<TResult>> ForEachAsyncImplUnordered<TResul
3030
long runId = ParallelAsyncEventSource.Log.GetRunId();
3131
ParallelAsyncEventSource.Log.RunStart(runId, batchSize, true, estimatedResultSize);
3232

33+
34+
#if false && NET6_0_OR_GREATER
35+
36+
var concurrentResult = new System.Collections.Concurrent.ConcurrentBag<TResult>();
37+
38+
try
39+
{
40+
var options = new System.Threading.Tasks.ParallelOptions
41+
{
42+
CancellationToken = cancellationToken,
43+
MaxDegreeOfParallelism = batchSize
44+
};
45+
46+
await System.Threading.Tasks.Parallel.ForEachAsync<TIn>(collection, options, async (i, ct) =>
47+
{
48+
var r = await func(i, ct).ConfigureAwait(false);
49+
concurrentResult.Add(r);
50+
}).ConfigureAwait(false);
51+
}
52+
catch (TaskCanceledException)
53+
{
54+
// Expected
55+
}
56+
57+
result.AddRange(concurrentResult);
58+
59+
#else
60+
3361
using (var enumerator = collection.GetEnumerator())
3462
{
3563
var hasNext = true;
@@ -81,6 +109,8 @@ private static async Task<IEnumerable<TResult>> ForEachAsyncImplUnordered<TResul
81109
}
82110
}
83111

112+
#endif
113+
84114
ParallelAsyncEventSource.Log.RunStop(runId);
85115

86116
return result;
@@ -100,6 +130,25 @@ private static async Task ForEachAsyncImplUnordered<TIn>(IEnumerable<TIn> collec
100130
long runId = ParallelAsyncEventSource.Log.GetRunId();
101131
ParallelAsyncEventSource.Log.RunStart(runId, batchSize, true, 0);
102132

133+
#if false && NET6_0_OR_GREATER
134+
135+
try
136+
{
137+
var options = new System.Threading.Tasks.ParallelOptions
138+
{
139+
CancellationToken = cancellationToken,
140+
MaxDegreeOfParallelism = batchSize
141+
};
142+
143+
await System.Threading.Tasks.Parallel.ForEachAsync<TIn>(collection, options, (i, ct) => new ValueTask(func(i, ct))).ConfigureAwait(false);
144+
}
145+
catch (TaskCanceledException)
146+
{
147+
// Expected
148+
}
149+
150+
#else
151+
103152
using (var enumerator = collection.GetEnumerator())
104153
{
105154
var hasNext = true;
@@ -140,6 +189,8 @@ private static async Task ForEachAsyncImplUnordered<TIn>(IEnumerable<TIn> collec
140189
}
141190
}
142191

192+
#endif
193+
143194
ParallelAsyncEventSource.Log.RunStop(runId);
144195
}
145196

@@ -154,6 +205,32 @@ private static async Task<IEnumerable<TResult>> ForEachAsyncImplUnordered<TResul
154205
long runId = ParallelAsyncEventSource.Log.GetRunId();
155206
ParallelAsyncEventSource.Log.RunStart(runId, batchSize, true, estimatedResultSize);
156207

208+
#if false && NET6_0_OR_GREATER
209+
210+
var concurrentResult = new System.Collections.Concurrent.ConcurrentBag<TResult>();
211+
212+
try
213+
{
214+
var options = new System.Threading.Tasks.ParallelOptions
215+
{
216+
CancellationToken = cancellationToken,
217+
MaxDegreeOfParallelism = batchSize
218+
};
219+
220+
await System.Threading.Tasks.Parallel.ForEachAsync<TIn>(collection, options, async (i, ct) =>
221+
{
222+
var r = await func(i, ct).ConfigureAwait(false);
223+
concurrentResult.Add(r);
224+
}).ConfigureAwait(false);
225+
}
226+
catch (TaskCanceledException)
227+
{
228+
// Expected
229+
}
230+
231+
result.AddRange(concurrentResult);
232+
#else
233+
157234
var enumerator = collection.GetAsyncEnumerator(cancellationToken);
158235
try
159236
{
@@ -210,6 +287,8 @@ private static async Task<IEnumerable<TResult>> ForEachAsyncImplUnordered<TResul
210287
await enumerator.DisposeAsync().ConfigureAwait(false);
211288
}
212289

290+
#endif
291+
213292
ParallelAsyncEventSource.Log.RunStop(runId);
214293

215294
return result;
@@ -220,6 +299,25 @@ private static async Task ForEachAsyncImplUnordered<TIn>(IAsyncEnumerable<TIn> c
220299
long runId = ParallelAsyncEventSource.Log.GetRunId();
221300
ParallelAsyncEventSource.Log.RunStart(runId, batchSize, true, 0);
222301

302+
#if false && NET6_0_OR_GREATER
303+
304+
try
305+
{
306+
var options = new System.Threading.Tasks.ParallelOptions
307+
{
308+
CancellationToken = cancellationToken,
309+
MaxDegreeOfParallelism = batchSize
310+
};
311+
312+
await System.Threading.Tasks.Parallel.ForEachAsync<TIn>(collection, options, (i, ct) => new ValueTask(func(i, ct))).ConfigureAwait(false);
313+
}
314+
catch (TaskCanceledException)
315+
{
316+
// Expected
317+
}
318+
319+
#else
320+
223321
var enumerator = collection.GetAsyncEnumerator(cancellationToken);
224322
try
225323
{
@@ -265,6 +363,8 @@ private static async Task ForEachAsyncImplUnordered<TIn>(IAsyncEnumerable<TIn> c
265363
await enumerator.DisposeAsync().ConfigureAwait(false);
266364
}
267365

366+
#endif
367+
268368
ParallelAsyncEventSource.Log.RunStop(runId);
269369
}
270370

tests/CSRakowski.Parallel.Benchmarks/Program.cs

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,10 @@ public static class Program
1616
{
1717
public static void Main(string[] args)
1818
{
19-
var summary = BenchmarkRunner.Run<ParallelAsyncBenchmarks_AsyncStreams>();
19+
#if NET6_0_OR_GREATER
20+
21+
var summary = BenchmarkRunner.Run<CompareWith_Parallel_ForEachAsync>();
22+
#endif
2023
}
2124
}
2225
}

0 commit comments

Comments
 (0)