如何实现TPL Dataflow聚合器:多源完成后发送单条消息
TPL Dataflow 实现未知数量源的聚合器方案
需求回顾
- 可变数量的源Block持有状态,接收消息后修改状态并向下游发送
- 聚合器需收集所有源的消息,检查错误,等待所有源完成后发送单条聚合结果到Releaser
- Releaser根据聚合结果更新状态,并发送最终消息
解决方案代码
using System; using System.Collections.Generic; using System.Linq; using System.Threading.Tasks; using System.Threading.Tasks.Dataflow; public static class TplDataflowAggregatorExample { public static void Run() { // 模拟动态数量的源Block(可根据需求动态添加) var sourceBlocks = new List<TransformBlock<int, int>> { new TransformBlock<int, int>(x => x * 2), new TransformBlock<int, int>(x => x * 3) }; // 暂存所有源消息的缓冲区 var messageBuffer = new BufferBlock<int>(); // 聚合器:接收消息数组,执行错误检查后传递给下游 var aggregator = new TransformBlock<int[], int[]>(messages => { // 自定义错误检查逻辑(示例:判断是否存在负数) if (messages.Any(m => m < 0)) { throw new InvalidOperationException("检测到错误消息"); } return messages; }); // Releaser:持有状态引用,根据聚合结果更新状态并输出 var state = new { TotalSum = 0, OperationSuccess = false }; var releaser = new TransformBlock<int[], (object CurrentState, bool IsSuccess, string Info)>(xs => { var sum = xs.Sum(); state = state with { TotalSum = sum, OperationSuccess = true }; return (state, true, $"操作完成,总和为{sum}"); }); // 链接所有源到消息缓冲区,自动传播完成状态 foreach (var source in sourceBlocks) { source.LinkTo(messageBuffer, new DataflowLinkOptions { PropagateCompletion = true }); } // 启动聚合任务:等待所有源完成后,聚合消息并发送给aggregator _ = Task.Run(async () => { // 等待所有源Block处理完成 await Task.WhenAll(sourceBlocks.Select(s => s.Completion)); // 取出缓冲区中所有消息 var collectedMessages = new List<int>(); while (messageBuffer.TryReceive(out var msg)) { collectedMessages.Add(msg); } // 发送聚合后的消息数组到aggregator,并标记完成 await aggregator.SendAsync(collectedMessages.ToArray()); aggregator.Complete(); }); // 链接聚合器到Releaser,传播完成与错误状态 aggregator.LinkTo(releaser, new DataflowLinkOptions { PropagateCompletion = true }); // 向源发送测试消息 sourceBlocks[0].Post(10); sourceBlocks[1].Post(20); // 等待整个数据流完成并处理结果/错误 try { releaser.Completion.Wait(); if (releaser.TryReceive(out var result)) { Console.WriteLine($"{result.Info},当前状态:{result.CurrentState}"); } } catch (AggregateException ex) { Console.WriteLine($"流程执行失败:{ex.InnerException?.Message ?? ex.Message}"); } } }
关键实现要点
消息暂存机制
使用BufferBlock<int>作为中间缓冲区,自动处理多源并发发送的消息,无需手动编写同步逻辑,确保所有消息都能被收集。跟踪所有源的完成状态
通过Task.WhenAll(sourceBlocks.Select(s => s.Completion))等待所有源Block的处理任务完成,避免提前聚合导致消息遗漏。聚合与错误检查
在所有源完成后,一次性取出缓冲区的所有消息,执行自定义错误检查(如消息合法性校验),若存在错误则抛出异常或返回错误标记,确保只有合法的聚合结果进入Releaser。动态源适配
源Block数量可变时,只需将新创建的Block加入sourceBlocks列表并链接到messageBuffer即可,聚合逻辑无需修改。状态传播与错误处理
链接Block时设置PropagateCompletion = true,确保上游的完成/错误状态能传递到下游,保证整个数据流的生命周期一致;若任何源或聚合器出错,错误会被传递到Releaser,最终可通过捕获AggregateException处理。
内容的提问来源于stack exchange,提问作者DavidY
相关产品推荐
相关产品推荐

