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

.NET Framework下异步加载多文件至缓存,如何限制并发读取数?

.NET Framework下文件异步预加载的替代方案优化

你的Rx实现本身是可行的,但如果团队对Rx不熟悉,或者想追求更轻量的依赖,下面两种方案可以作为替代:

方案一:基于BlockingCollection + Task.WhenAny的经典并发控制

这种方案不需要额外引入Rx依赖,用.NET Framework原生的Task API实现,逻辑直观易懂:

private static async Task PreloadBinariesAsync(BlockingCollection<IProjectEx> inputs, int mdop, ConcurrentDictionary<string, byte[]> binaries, Progress progress)
{
    var activeTasks = new HashSet<Task<(string FilePath, byte[] Bytes)>>();

    foreach (var project in inputs.GetConsumingEnumerable())
    {
        // 当活跃任务数达到MDOP时,等待任一任务完成
        while (activeTasks.Count >= mdop)
        {
            var completedTask = await Task.WhenAny(activeTasks).ConfigureAwait(false);
            activeTasks.Remove(completedTask);
            var result = await completedTask.ConfigureAwait(false);
            progress.IncAsmFileCachedCount();
            binaries[result.FilePath] = result.Bytes;
        }

        // 启动新的文件读取任务
        activeTasks.Add(ReadAllBytesAsync(project.TargetPath));
    }

    // 等待剩余所有任务完成
    foreach (var task in activeTasks)
    {
        var result = await task.ConfigureAwait(false);
        progress.IncAsmFileCachedCount();
        binaries[result.FilePath] = result.Bytes;
    }
}

private static async Task<(string FilePath, byte[] Bytes)> ReadAllBytesAsync(string filePath) 
    => (filePath, await File.ReadAllBytesAsync(filePath).ConfigureAwait(false));

方案二:替换为Channel实现异步生产者消费者

如果可以引入NuGet包(需安装System.Threading.Tasks.Channel,.NET Framework原生不带此组件),Channel比BlockingCollection更适合异步场景,它的异步API设计更自然,避免不必要的线程阻塞:

// 核心预加载逻辑,基于Channel实现
private static async Task PreloadBinariesAsync(Channel<IProjectEx> inputChannel, int mdop, ConcurrentDictionary<string, byte[]> binaries, Progress progress)
{
    var reader = inputChannel.Reader;
    var semaphore = new SemaphoreSlim(mdop);

    async Task ProcessProject(IProjectEx project)
    {
        await semaphore.WaitAsync().ConfigureAwait(false);
        try
        {
            var result = await ReadAllBytesAsync(project.TargetPath).ConfigureAwait(false);
            progress.IncAsmFileCachedCount();
            binaries[result.FilePath] = result.Bytes;
        }
        finally
        {
            semaphore.Release();
        }
    }

    var tasks = new List<Task>();
    while (await reader.WaitToReadAsync().ConfigureAwait(false))
    {
        while (reader.TryRead(out var project))
        {
            tasks.Add(ProcessProject(project));
        }
    }

    await Task.WhenAll(tasks).ConfigureAwait(false);
}

// 辅助方法:将原有BlockingCollection的数据导入Channel(若生产者端可直接写入Channel则可省略)
private static Channel<IProjectEx> BlockingCollectionToChannel(BlockingCollection<IProjectEx> collection)
{
    var channel = Channel.CreateUnbounded<IProjectEx>();
    _ = Task.Run(async () =>
    {
        foreach (var item in collection.GetConsumingEnumerable())
        {
            await channel.Writer.WriteAsync(item).ConfigureAwait(false);
        }
        channel.Writer.Complete();
    });
    return channel;
}

private static async Task<(string FilePath, byte[] Bytes)> ReadAllBytesAsync(string filePath) 
    => (filePath, await File.ReadAllBytesAsync(filePath).ConfigureAwait(false));

方案对比

  • 原Rx方案:适合复杂的数据流转换场景(比如需要过滤、映射、合并等多步操作),但引入了Rx依赖,学习曲线较高。
  • Task.WhenAny方案:原生无依赖,逻辑直观,适合团队熟悉传统异步编程的场景,并发控制逻辑清晰。
  • Channel方案:异步友好,性能更优(尤其是高并发场景),但需要引入NuGet包,适合可以升级依赖的项目。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 22:03:11