.NET Core 3.1中如何将AsyncPageable的await foreach转为Parallel.ForEach
.NET Core 3.1 下将AsyncPageable的await foreach转为Parallel.ForEach的实现
针对遍历AsyncPageable<ChatParticipant>的异步循环转并行循环的需求,首先明确几个框架原生限制:
Parallel.ForEach仅支持同步可枚举的IEnumerable<T>数据源,无法直接传入实现IAsyncEnumerable<T>的AsyncPageable<T>- 普通
List<T>不是线程安全类型,多线程并行调用Add会出现数据丢失、内部结构损坏的异常,必须替换为线程安全集合或对写入操作加锁 Parallel.ForEach设计用于CPU密集型的同步并行逻辑,不支持异步委托,不要在循环内写asynclambda,否则会生成async void,引发异常无法捕获、线程池调度混乱的问题
方案1:全量拉取后并行(实现简单,适合数据量可控场景)
先通过异步枚举把所有分页数据加载到内存,再传入Parallel.ForEach处理,结果用线程安全的ConcurrentBag<T>存储:
using System.Collections.Concurrent; // 拉取全量分页数据 AsyncPageable<ChatParticipant> pagedResult = chatThreadClient.GetParticipantsAsync(); var allParticipants = new List<ChatParticipant>(); await foreach (var participant in pagedResult) { allParticipants.Add(participant); } var participants = new ConcurrentBag<string>(); // 执行并行处理 Parallel.ForEach(allParticipants, participant => { // 此处放需要CPU并行计算的业务逻辑 participants.Add(participant.User.ToString()); }); // 若需要List类型结果直接转换即可 List<string> finalResult = participants.ToList();
方案2:边拉取边并行(适合数据量大、不想一次性加载全量数据场景)
按分页拉取数据,每拿到一页就立刻提交并行处理,避免全量数据占用过多内存:
using System.Collections.Concurrent; AsyncPageable<ChatParticipant> pagedResult = chatThreadClient.GetParticipantsAsync(); var participants = new ConcurrentBag<string>(); // 按页枚举数据源 await foreach (var page in pagedResult.AsPages()) { // 并行处理当前页的所有数据 Parallel.ForEach(page.Values, participant => { participants.Add(participant.User.ToString()); }); } List<string> finalResult = participants.ToList();
补充提示
如果循环内包含IO绑定的异步操作(比如接口调用、数据库读写),不要强行用Parallel.ForEach。.NET Core 3.1没有内置的异步并行ForEach方法,可以通过Task.WhenAll搭配信号量限流实现异步并行处理,避免并发过高压垮下游服务。
内容的提问来源于stack exchange,提问作者user1037747
相关产品推荐
相关产品推荐

