You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

自定义LINQ运算符SelectParallelAsync无法实现任务并发运行的问题

解决IAsyncEnumerable并发SelectAwait的问题

你的代码无法实现预期并发效果,核心有两个问题:

  1. 信号量初始配置错误:new SemaphoreSlim(1, 10) 将初始并发数设为1,直接限制了同时只能执行1个任务
  2. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.03 04:25:37