如何为Parallel.ForEachAsync配置NoBuffering无缓冲行为?
同步方法Parallel.ForEach提供了多个重载,其中部分支持通过EnumerablePartitionerOptions.NoBuffering配置并行循环,该选项的作用是:
创建一个分区器,从源可枚举对象中逐个获取项,不使用可被多线程高效访问的中间存储。此选项支持低延迟(项一从源可用就会被处理),并部分支持项之间的依赖关系(线程不会因等待自身负责处理的项而死锁)。
但异步方法Parallel.ForEachAsync并没有类似的选项或重载,这给我的生产者-消费者场景带来了困扰:我需要以Channel<T>作为源,让消费者仅处理自身能力范围内的项,不能提前拉取多余项存入内部缓冲区。我希望Channel<T>是系统中唯一的队列,这样才能准确监控等待处理的项数。
此前我误以为Parallel.ForEachAsync设计上不会做缓冲,但在向微软确认后得到了这样的反馈:
这是实现细节。对于
Parallel.ForEach,缓冲是为了处理可能极快的主体委托,从而尽量减少/分摊获取锁以访问共享枚举器的成本。对于ForEachAsync,预期主体委托至少会有一定复杂度,因此目前不会尝试此类分摊操作。
依赖实现细节显然不可靠,因此我需要重新寻找解决方案。
我的问题
是否可以配置Parallel.ForEachAsync API,确保它具备NoBuffering的行为?如果可以,具体该如何操作?
说明:我并非要从零实现Parallel.ForEachAsync,而是需要一个基于现有API的轻量包装,注入理想的NoBuffering行为,示例如下:
public static Task ForEachAsync_NoBuffering<TSource>( IAsyncEnumerable<TSource> source, ParallelOptions parallelOptions, Func<TSource, CancellationToken, ValueTask> body) { // 实现逻辑 return Parallel.ForEachAsync(source, parallelOptions, body); }
该包装的行为需与.NET 6上的Parallel.ForEachAsync完全一致。
更新:场景代码示例
class Processor { private readonly Channel<Item> _channel; private readonly Task _consumer; public Processor() { _channel = Channel.CreateUnbounded<Item>(); _consumer = StartConsumer(); } public int PendingItemsCount => _channel.Reader.Count; public Task Completion => _consumer; public void QueueItem(Item item) => _channel.Writer.TryWrite(item); private async Task StartConsumer() { ParallelOptions options = new() { MaxDegreeOfParallelism = 2 }; await Parallel.ForEachAsync(_channel.Reader.ReadAllAsync(), options, async (item, _) => { // 调用异步API // 将API响应持久化到关系型数据库 }); } }
我倾向于使用.NET 6的Parallel.ForEachAsync API,这是本次问题的核心。
内容的提问来源于stack exchange,提问作者Theodor Zoulias

