自定义Dataflow转换块:动态MaxDegreeOfParallelism实现问询
问题描述
我们有如下定义的类型:
public record Item(string Name); public record ItemGroup(Item[] Items, int MaxDegreeOfParallelism);
需要处理ItemGroup实例序列,规则是:
- ItemGroup必须按顺序处理:前一个ItemGroup的所有Item处理完成后,才能开始处理下一个ItemGroup
- 每个ItemGroup内部可并行处理Item:并行度由该ItemGroup的
MaxDegreeOfParallelism指定
示例场景:
var groups = new[] { new ItemGroup(new[] { new Item("A0"), new Item("A1"), new Item("A2") }, 1), // 串行处理A组所有项 new ItemGroup(new[] { new Item("B0"), new Item("B1"), new Item("B2") }, 3) // 并行处理B组所有项 };
考虑基于IPropagatorBlock<ItemGroup, Item>实现自定义块,但不清楚当生产者发布ItemGroup时,如何正确等待内部动态创建的TransformManyBlock完成处理。
实现方案
不需要直接继承TransformManyBlock,通过组合TPL Dataflow块即可实现需求:用外层块控制ItemGroup的顺序处理,为每个ItemGroup动态创建带指定并行度的内部块,等待内部块处理完成后再推进到下一个ItemGroup。
以下是完整实现:
using System.Threading.Tasks.Dataflow; public class OrderedGroupParallelProcessor : IPropagatorBlock<ItemGroup, Item> { private readonly ActionBlock<ItemGroup> _groupProcessingBlock; private readonly BroadcastBlock<Item> _outputBlock; public OrderedGroupParallelProcessor() { _outputBlock = new BroadcastBlock<Item>(item => item); _groupProcessingBlock = new ActionBlock<ItemGroup>(async group => { // 为当前ItemGroup创建带指定并行度的转换块 var itemProcessor = new TransformManyBlock<ItemGroup, Item>( g => g.Items, new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = group.MaxDegreeOfParallelism }); // 链接内部块到输出块,自动传播完成状态 itemProcessor.LinkTo(_outputBlock, new DataflowLinkOptions { PropagateCompletion = true }); // 提交当前ItemGroup到内部块 await itemProcessor.SendAsync(group); // 标记内部块停止接受新输入 itemProcessor.Complete(); // 等待内部块处理完成,确保当前ItemGroup所有Item处理完毕再进入下一个 await itemProcessor.Completion; }); } public Task Completion => _groupProcessingBlock.Completion.ContinueWith(_ => _outputBlock.Completion); public void Complete() { _groupProcessingBlock.Complete(); } public void Fault(Exception exception) { ((IDataflowBlock)_groupProcessingBlock).Fault(exception); ((IDataflowBlock)_outputBlock).Fault(exception); } public IDisposable LinkTo(ITargetBlock<Item> target, DataflowLinkOptions linkOptions) { return _outputBlock.LinkTo(target, linkOptions); } public DataflowMessageStatus OfferMessage(DataflowMessageHeader messageHeader, ItemGroup messageValue, ISourceBlock<ItemGroup> source, bool consumeToAccept) { return ((ISourceBlock<ItemGroup>)_groupProcessingBlock).OfferMessage(messageHeader, messageValue, source, consumeToAccept); } }
关键逻辑说明
- 顺序控制:外层
ActionBlock<ItemGroup>默认并行度为1,确保ItemGroup按输入顺序依次处理,前一个处理完成后才会启动下一个。 - 内部并行处理:每个ItemGroup对应独立的
TransformManyBlock,使用该ItemGroup指定的MaxDegreeOfParallelism,实现组内Item的并行处理。 - 等待内部块完成:通过
await itemProcessor.Completion强制等待当前ItemGroup的所有Item处理结束,再进入下一个ItemGroup的处理流程。 - 输出传播:内部块链接到
BroadcastBlock,确保处理后的Item能传递到下游块,同时通过PropagateCompletion保证完成状态正确传递到整个管道。
使用示例
var processor = new OrderedGroupParallelProcessor(); // 下游消费块,处理输出的Item var consumer = new ActionBlock<Item>(item => { Console.WriteLine($"处理Item: {item.Name}, 线程ID: {Environment.CurrentManagedThreadId}"); // 模拟处理耗时 Task.Delay(100).Wait(); }); processor.LinkTo(consumer, new DataflowLinkOptions { PropagateCompletion = true }); // 发布ItemGroup序列 var groups = new[] { new ItemGroup(new[] { new Item("A0"), new Item("A1"), new Item("A2") }, 1), new ItemGroup(new[] { new Item("B0"), new Item("B1"), new Item("B2") }, 3) }; foreach (var group in groups) { processor.SendAsync(group).Wait(); } processor.Complete(); consumer.Completion.Wait();
运行后会观察到:A组Item串行处理(线程ID一致),B组Item并行处理(线程ID不同),且A组全部处理完成后才会启动B组的处理。
内容的提问来源于stack exchange,提问作者Ritmo2k
相关产品推荐
相关产品推荐

