自定义LINQ运算符SelectParallelAsync无法实现任务并发运行的问题
解决IAsyncEnumerable并发SelectAwait的问题
你的代码无法实现预期并发效果,核心有两个问题:
- 信号量初始配置错误:
new SemaphoreSlim(1, 10)将初始并发数设为1,直接限制了同时只能执行1个任务 - IAsyncEnumerable的延迟迭代特性:
enumerable.Select是延迟执行的,await foreach会逐个拉取元素,每次迭代才会创建并启动对应的任务,导致任务无法批量并发启动
以下是两种满足需求的实现方案,同时支持信号量共享:
方案1:保持输出与输入顺序一致
适合需要结果顺序和原始枚举顺序匹配的场景:
private async IAsyncEnumerable<TOut> SelectParallelAsync<T, TOut>( this IAsyncEnumerable<T> enumerable, Func<T, Task<TOut>> predicate, int maxDegreeOfParallelism = 10) { using var semaphore = new SemaphoreSlim(maxDegreeOfParallelism, maxDegreeOfParallelism); var tasks = new List<Task<(int Index, TOut Result)>>(); int index = 0; // 先遍历所有元素,批量启动受信号量控制的任务 await foreach (var item in enumerable) { int currentIndex = index++; tasks.Add(Task.Run(async () => { await semaphore.WaitAsync(); try { var result = await predicate(item); return (currentIndex, result); } finally { semaphore.Release(); } })); } // 等待所有任务完成,按原始顺序返回结果 foreach (var task in await Task.WhenAll(tasks)) { yield return task.Result; } }
方案2:任务完成即返回(不保证顺序)
适合不需要严格顺序、希望尽快输出结果的场景:
private async IAsyncEnumerable<TOut> SelectParallelUnorderedAsync<T, TOut>( this IAsyncEnumerable<T> enumerable, Func<T, Task<TOut>> predicate, int maxDegreeOfParallelism = 10) { using var semaphore = new SemaphoreSlim(maxDegreeOfParallelism, maxDegreeOfParallelism); var tasks = new ConcurrentQueue<Task<TOut>>(); // 后台线程负责遍历元素并提交任务 var enumerationTask = Task.Run(async () => { await foreach (var item in enumerable) { tasks.Enqueue(Task.Run(async () => { await semaphore.WaitAsync(); try { return await predicate(item); } finally { semaphore.Release(); } })); } }); // 持续处理完成的任务,直到全部完成 while (!enumerationTask.IsCompleted || tasks.Count > 0) { var completedTask = await Task.WhenAny(tasks); tasks.TryDequeue(out _); yield return await completedTask; } // 确保遍历过程无异常 await enumerationTask; }
支持跨方法共享信号量
如果需要在多个方法间复用同一个信号量,只需将信号量作为参数传入即可:
private async IAsyncEnumerable<TOut> SelectParallelAsync<T, TOut>( this IAsyncEnumerable<T> enumerable, Func<T, Task<TOut>> predicate, SemaphoreSlim sharedSemaphore) { var tasks = new List<Task<(int Index, TOut Result)>>(); int index = 0; await foreach (var item in enumerable) { int currentIndex = index++; tasks.Add(Task.Run(async () => { await sharedSemaphore.WaitAsync(); try { return (currentIndex, await predicate(item)); } finally { sharedSemaphore.Release(); } })); } foreach (var task in await Task.WhenAll(tasks)) { yield return task.Result; } }
调用示例:
// 创建共享信号量,控制全局并发数为10 var sharedSem = new SemaphoreSlim(10, 10); // 多个方法共用该信号量 await foreach (var result in enumerable.SelectParallelAsync(async i => { Console.WriteLine($"In Select : {i}"); await Task.Delay(1000); return i + 5; }, sharedSem)) { // 处理结果 }
内容的提问来源于stack exchange,提问作者Grimson
相关产品推荐
相关产品推荐

