如何并行请求分页API并将结果实时写入返回的IAsyncEnumerable
问题解答
核心问题回复
- 你的需求完全可以实现
- 不允许多个异步Task同时直接写入同一个
IAsyncEnumerable:yield return是编译器生成的单线程状态机,本身不支持跨线程并发调用,直接操作会引发未知异常或数据错乱。
实现方案
我们可以借助System.Threading.Channels(.NET官方提供的异步生产者消费者队列实现,天生线程安全,支持多生产者并发写入、单消费者读取)来做中转,实现“拿到一页就立刻返回”的效果:
- 先创建一个Channel作为临时数据缓存
- 后台启动所有分页请求任务,每拿到一页的结果就直接写入Channel
- 方法直接返回Channel读取器的异步枚举,消费者枚举时只要有数据就会立刻收到,无需等待所有请求完成
- 所有分页请求处理完成后,标记Channel写入结束,枚举自动终止
代码实现
如果是低于.NET Core 3.0的版本,需要先安装System.Threading.Channels NuGet包,更高版本框架已内置该组件:
using System.Threading.Channels; public async IAsyncEnumerable<ApiObject> GetResults() { int totalPages = _client.MagicGetTotalPagesMethod(); // 可根据实际内存限制选择CreateBounded设置队列上限,避免接口响应过慢导致内存暴涨 var resultChannel = Channel.CreateUnbounded<ApiObject>(); // 后台启动所有分页拉取任务,不阻塞当前返回 _ = Task.Run(async () => { try { // 所有分页请求的任务集合 var pageTasks = Enumerable.Range(0, totalPages).Select(async pageIndex => { var response = await _client.GetResults(_endpoint, pageIndex * _pageSize, _pageSize); // 当前页数据写入队列,写入后消费者立刻就能收到 foreach (var item in response.Items) { await resultChannel.Writer.WriteAsync(item); } }); await Task.WhenAll(pageTasks); } finally { // 所有请求处理完成,标记队列没有更多数据 resultChannel.Writer.Complete(); } }); // 读取队列数据并返回给消费者 await foreach (var item in resultChannel.Reader.ReadAllAsync()) { yield return item; } }
可选优化:限制并发请求数
如果不希望同时发起所有分页请求避免给API造成过大压力,可以添加SemaphoreSlim控制并行度,示例如下:
public async IAsyncEnumerable<ApiObject> GetResults(int maxConcurrentRequests = 5) { int totalPages = _client.MagicGetTotalPagesMethod(); var resultChannel = Channel.CreateUnbounded<ApiObject>(); // 控制最多同时发起maxConcurrentRequests个请求 using var semaphore = new SemaphoreSlim(maxConcurrentRequests); _ = Task.Run(async () => { try { var pageTasks = Enumerable.Range(0, totalPages).Select(async pageIndex => { await semaphore.WaitAsync(); try { var response = await _client.GetResults(_endpoint, pageIndex * _pageSize, _pageSize); foreach (var item in response.Items) { await resultChannel.Writer.WriteAsync(item); } } finally { semaphore.Release(); } }); await Task.WhenAll(pageTasks); } finally { resultChannel.Writer.Complete(); } }); await foreach (var item in resultChannel.Reader.ReadAllAsync()) { yield return item; } }
原有实现的问题
你之前的实现还有一个隐藏风险:List<T>的AddRange方法不是线程安全的,多个Task并发调用会出现数据丢失、索引越界等异常,Channel的并发写入是线程安全的,完全规避了这个问题。
内容的提问来源于stack exchange,提问作者ewok
相关产品推荐
相关产品推荐

