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

C#中如何等待await foreach内部创建的多个并行异步任务

问题背景

在使用AsAsyncEnumerable配合await foreach做懒加载遍历集合时,需求为:对每个遍历到的元素调用返回Task的DoSomething方法,所有方法调用并行执行,全部执行完成后方法返回Unit类型。
初始的问题实现如下:

public async Task<Unit> Handle(SomeRequest request, CancellationToken cancellationToken)
{
    var entities = _dbContext.MyModel
        .Where(...)
        .AsAsyncEnumerable();

    await foreach (var item in entities)
    {
        // 直接丢弃任务引用,不等待执行
        _ = item.DoSomething();
    }

    return Unit.Value;
}

这个实现存在明显缺陷:遍历完所有元素后方法会立刻返回,不会等待任何DoSomething任务执行完成,属于典型的“发后不理”(fire-and-forget)写法,任务执行中抛出的异常会被直接吞没,业务逻辑根本无法保证执行完成。


收集任务后调用Task.WhenAll的方案可行性

你给出的如下写法,核心逻辑是正确的,可以满足所有任务并行执行、全部完成后再返回的基础要求:

public async Task<Unit> Handle(SomeRequest request, CancellationToken cancellationToken)
{
    var entities = _dbContext.MyModel
        .Where(...)
        .AsAsyncEnumerable();

    List<Task> tasks = new();

    await foreach (var item in entities)
    {
        tasks.Add(item.DoSomething());
    }

    await Task.WhenAll(tasks);

    return Unit.Value;
}

但这个基础实现存在三个容易引发线上问题的隐患:

  • 无并发度控制:如果懒加载返回的元素量级很大(比如数千、数万条),所有DoSomething任务会在遍历过程中几乎同时启动,极容易打满数据库连接池、触发下游服务限流,极端情况下还会引发内存溢出。
  • 未传递取消令牌:方法入参传入的cancellationToken没有透传给异步遍历和DoSomething方法,请求触发取消时无法及时终止正在执行的任务,会造成不必要的资源浪费。
  • 异常定位成本高:如果多个任务同时抛出异常,Task.WhenAll会把所有异常包装到AggregateException中,排查问题时需要额外提取内部异常才能定位根因。

生产环境优化实现

如果要规避上述隐患,推荐用SemaphoreSlim控制最大并发度,同时正确传递取消令牌,参考代码如下:

public async Task<Unit> Handle(SomeRequest request, CancellationToken cancellationToken)
{
    // 根据业务场景、下游承载能力设置最大并行数,IO密集型场景可适当调高,CPU密集型场景建议和CPU核心数对齐
    var maxParallelCount = 10;
    using var concurrencySemaphore = new SemaphoreSlim(maxParallelCount);
    var entities = _dbContext.MyModel
        .Where(...)
        .AsAsyncEnumerable();

    var tasks = new List<Task>();
    // 给异步遍历也传入取消令牌
    await foreach (var item in entities.WithCancellation(cancellationToken))
    {
        tasks.Add(ProcessSingleItemAsync(item, concurrencySemaphore, cancellationToken));
    }

    await Task.WhenAll(tasks);
    return Unit.Value;
}

private async Task ProcessSingleItemAsync(MyModel item, SemaphoreSlim concurrencySemaphore, CancellationToken cancellationToken)
{
    await concurrencySemaphore.WaitAsync(cancellationToken);
    try
    {
        // 给DoSomething也传递取消令牌
        await item.DoSomething(cancellationToken);
    }
    finally
    {
        concurrencySemaphore.Release();
    }
}

提示:如果你的数据量特别大,不想一次性把所有任务对象都存在列表里占内存,也可以结合Channel实现生产消费模式,一边遍历元素往通道里写,一边用固定数量的消费者并行处理元素,内存占用会更稳定,但核心等待逻辑依然是等待所有处理任务完成后再返回,不要直接丢弃任务引用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 14:18:27