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

