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

如何并行请求分页API并将结果实时写入返回的IAsyncEnumerable

问题解答

核心问题回复

  • 你的需求完全可以实现
  • 不允许多个异步Task同时直接写入同一个IAsyncEnumerable:yield return是编译器生成的单线程状态机,本身不支持跨线程并发调用,直接操作会引发未知异常或数据错乱。

实现方案

我们可以借助System.Threading.Channels(.NET官方提供的异步生产者消费者队列实现,天生线程安全,支持多生产者并发写入、单消费者读取)来做中转,实现“拿到一页就立刻返回”的效果:

  1. 先创建一个Channel作为临时数据缓存
  2. 后台启动所有分页请求任务,每拿到一页的结果就直接写入Channel
  3. 方法直接返回Channel读取器的异步枚举,消费者枚举时只要有数据就会立刻收到,无需等待所有请求完成
  4. 所有分页请求处理完成后,标记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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 17:15:04