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
相关产品推荐
相关产品推荐

