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

如何在Parallel.ForEach中等待异步函数执行完成?

异步并行任务等待问题的修复

问题重现

现有以下用于在字典中收集计算结果的函数:

void Func()
{
    // ...
    var dtDict = await HandleComputeBooster();
    // ...
}

async private static Task DoBooster(..., ConcurrentDictionary<string, DataTable> dtDict, ...)
{
    DataTable dt = ...;
    // ... 数据处理逻辑
    dtDict[symbol] = dt;
    // ...
}

HandleComputeBooster函数存在提前返回的问题:

async private Task<ConcurrentDictionary<string, DataTable>> HandleComputeBooster()
{
    var dtDict = new ConcurrentDictionary<string, DataTable>();

    // ... 初始化逻辑
    var chunks = listOfBoosterSymbols.ChunkBy(8);
    var pcCount = Environment.ProcessorCount;

    Parallel.ForEach(chunks, new ParallelOptions
    { MaxDegreeOfParallelism = pcCount - 2 }, async listStr =>
        {
            var symbol = listStr[0];
            await DoBooster(..., dtDict, ...);
        }
    );
    // ...

    return dtDict;
}

问题在于:HandleComputeBooster会在dtDict的所有计算任务完成前就返回。虽然最终所有数据都会存入字典,但需要让函数等待所有块处理完成后再返回。

问题原因

Parallel.ForEach不支持异步委托:传入的async lambda会被当作async void执行。Parallel.ForEach仅会等待线程池中的任务启动完成,不会等待委托内部的await操作结束,因此函数会提前返回,此时字典中的数据还未全部生成。

解决方案

改用Task.WhenAll结合异步遍历的方式,来正确等待所有异步任务完成:

async private Task<ConcurrentDictionary<string, DataTable>> HandleComputeBooster()
{
    var dtDict = new ConcurrentDictionary<string, DataTable>();

    // ... 初始化逻辑
    var chunks = listOfBoosterSymbols.ChunkBy(8);

    // 将每个chunk的处理包装为Task
    var tasks = chunks.Select(async listStr =>
        {
            var symbol = listStr[0];
            await DoBooster(..., dtDict, ...);
        });

    // 等待所有异步任务完成
    await Task.WhenAll(tasks);

    // ...

    return dtDict;
}

可选:限制并发数

如果需要像原代码一样限制最大并发数,可以使用SemaphoreSlim来控制:

async private Task<ConcurrentDictionary<string, DataTable>> HandleComputeBooster()
{
    var dtDict = new ConcurrentDictionary<string, DataTable>();

    // ... 初始化逻辑
    var chunks = listOfBoosterSymbols.ChunkBy(8);
    var pcCount = Environment.ProcessorCount;
    var semaphore = new SemaphoreSlim(pcCount - 2);

    var tasks = chunks.Select(async listStr =>
        {
            await semaphore.WaitAsync();
            try
            {
                var symbol = listStr[0];
                await DoBooster(..., dtDict, ...);
            }
            finally
            {
                semaphore.Release();
            }
        });

    await Task.WhenAll(tasks);
    semaphore.Dispose();

    // ...

    return dtDict;
}

说明

  • 因为DoBooster是异步方法,属于IO密集型操作,使用异步并行(Task.WhenAll)比Parallel.ForEach更合适,能更高效利用线程资源。
  • ConcurrentDictionary本身是线程安全的,多个异步任务同时写入不会有问题。

内容的提问来源于stack exchange,提问作者Ivan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 19:01:01