C#中如何同步三个存在依赖关系的循环执行任务?
实现方案
这个场景属于典型的多阶段流水线处理场景,C#官方提供的TPL Dataflow库是最适配的实现方案,不需要自己手动实现线程同步、任务调度逻辑。
前置依赖
需要先安装NuGet包 System.Threading.Tasks.Dataflow,这是.NET官方维护的库,稳定性和性能都有保障。
核心代码示例
using System.Threading.Tasks.Dataflow; // 这里可以根据你的实际业务替换输入输出类型,示例用string做演示 var task1Block = new TransformBlock<string, string>(input => { // 这里写Task1的业务逻辑 Console.WriteLine($"Task1 处理:{input}"); Task.Delay(1000).Wait(); // 模拟任务耗时 return $"Task1处理结果_{input}"; }, new ExecutionDataflowBlockOptions { // 配置Task1的并行度,根据任务是IO/CPU密集调整,默认是1 MaxDegreeOfParallelism = 1 }); var task2Block = new TransformBlock<string, string>(task1Result => { // 这里写Task2的业务逻辑,依赖Task1的输出 Console.WriteLine($"Task2 处理:{task1Result}"); Task.Delay(2000).Wait(); // 模拟任务耗时 return $"Task2处理结果_{task1Result}"; }, new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = 1 }); var task3Block = new ActionBlock<string>(task2Result => { // 这里写Task3的业务逻辑,依赖Task2的输出 Console.WriteLine($"Task3 处理:{task2Result}"); Task.Delay(1500).Wait(); // 模拟任务耗时 }, new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = 1 }); // 把三个块链接成流水线,自动传递处理结果和完成状态 task1Block.LinkTo(task2Block, new DataflowLinkOptions { PropagateCompletion = true }); task2Block.LinkTo(task3Block, new DataflowLinkOptions { PropagateCompletion = true }); // 模拟持续输入,你可以根据实际业务替换成你的输入源 int count = 0; while (true) { task1Block.Post($"输入_{count++}"); // 这里可以控制输入速度,如果是无限制输入可以去掉Delay await Task.Delay(500); } // 如果需要停止流水线,可以调用 // task1Block.Complete(); // await task3Block.Completion;
方案特性
- 完全匹配你的需求:Task1处理完一个输入立刻把结果传给Task2,同时马上开始处理下一个输入,不需要等Task2、Task3执行完成
- 依赖关系天然保证:Task2只会拿到Task1处理完成的结果,Task3只会拿到Task2处理完成的结果,不会出现顺序错乱
- 自带背压能力:如果后面的任务处理速度比Task1慢,Dataflow会自动限制Task1的生产速度,避免内存溢出
- 灵活调整并行度:每个任务可以独立配置并行数,比如Task1如果是IO密集型可以把
MaxDegreeOfParallelism调大提高吞吐量 - 异常处理统一:任意块抛出的异常都会传递到最终的Completion任务里,不会出现未捕获的线程异常
轻量替代方案(不引入额外NuGet包)
如果不想引入Dataflow库,也可以用.NET内置的Channel实现:
using System.Threading.Channels; // 定义两个Channel传递Task1、Task2的结果 var channel1 = Channel.CreateUnbounded<string>(); var channel2 = Channel.CreateUnbounded<string>(); // 启动Task1循环 _ = Task.Run(async () => { int count = 0; while (true) { var input = $"输入_{count++}"; Console.WriteLine($"Task1 处理:{input}"); await Task.Delay(1000); var result = $"Task1处理结果_{input}"; await channel1.Writer.WriteAsync(result); } }); // 启动Task2循环 _ = Task.Run(async () => { await foreach (var task1Result in channel1.Reader.ReadAllAsync()) { Console.WriteLine($"Task2 处理:{task1Result}"); await Task.Delay(2000); var result = $"Task2处理结果_{task1Result}"; await channel2.Writer.WriteAsync(result); } }); // 启动Task3循环 _ = Task.Run(async () => { await foreach (var task2Result in channel2.Reader.ReadAllAsync()) { Console.WriteLine($"Task3 处理:{task2Result}"); await Task.Delay(1500); } }); // 阻塞主线程避免程序退出 await Task.Delay(Timeout.Infinite);
这个方案更轻量,但是需要自己处理异常、背压、并行度控制等逻辑,适合简单场景使用。
内容的提问来源于stack exchange,提问作者Abdelsalam Hamdi
相关产品推荐
相关产品推荐

