C#多线程树遍历问题:如何判断BlockingCollection任务全部完成?
解决C#多线程树遍历中BlockingCollection的任务完成检测问题
这个场景我太熟悉了!当消费者同时又是生产者的时候,单纯判断BlockingCollection是否为空或者调用IsCompleted根本靠不住——你永远没法确定是不是有线程刚处理完一个节点,正准备往队列里塞它的子节点呢。下面给你几个经过实践验证的可行方案:
方案一:用SemaphoreSlim跟踪活跃工作项
核心思路是用一个信号量来统计当前正在处理(或者待处理)的节点总数,当信号量计数归0且队列为空时,就代表所有遍历操作完成了。
示例代码:
using System.Collections.Concurrent; using System.Threading; class TreeTraversal { static void Main() { var root = new TreeNode { Value = "Root", Children = { new TreeNode("Child1"), new TreeNode("Child2") } }; var queue = new BlockingCollection<TreeNode>(); // 初始化信号量,初始计数0,最大不限 var semaphore = new SemaphoreSlim(0); // 启动多个消费者线程 for (int i = 0; i < Environment.ProcessorCount; i++) { ThreadPool.QueueUserWorkItem(_ => { while (true) { if (queue.TryTake(out var node, Timeout.Infinite)) { try { // 处理当前节点(这里换成你的业务逻辑) Console.WriteLine($"Processing {node.Value}"); // 把子节点加入队列,每加一个就给信号量+1 foreach (var child in node.Children) { queue.Add(child); semaphore.Release(); } } finally { // 当前节点处理完成,信号量-1 semaphore.Wait(); } } else { // 队列已完成添加,退出循环 break; } } }); } // 加入根节点,启动流程 queue.Add(root); semaphore.Release(); // 等待所有工作项完成 semaphore.Wait(); // 确认队列确实为空,然后标记完成 queue.CompleteAdding(); Console.WriteLine("All traversal completed!"); } } class TreeNode { public string Value { get; set; } public List<TreeNode> Children { get; } = new List<TreeNode>(); public TreeNode() { } public TreeNode(string value) => Value = value; }
方案二:用Interlocked计数器+CompleteAdding标记
通过原子计数器跟踪当前正在处理的任务数,当计数器归0且队列为空时,手动调用CompleteAdding(),让所有消费者线程退出循环。
示例代码:
using System.Collections.Concurrent; using System.Threading; class TreeTraversal { static int _activeTasks; static void Main() { var root = new TreeNode { Value = "Root", Children = { new TreeNode("Child1"), new TreeNode("Child2") } }; var queue = new BlockingCollection<TreeNode>(); // 启动消费者线程 for (int i = 0; i < Environment.ProcessorCount; i++) { ThreadPool.QueueUserWorkItem(_ => { while (!queue.IsCompleted) { if (queue.TryTake(out var node, 100)) { try { Interlocked.Increment(ref _activeTasks); // 处理节点 Console.WriteLine($"Processing {node.Value}"); // 添加子节点到队列 foreach (var child in node.Children) { queue.Add(child); } } finally { Interlocked.Decrement(ref _activeTasks); // 检查是否所有任务都完成且队列为空 if (_activeTasks == 0 && queue.Count == 0) { queue.CompleteAdding(); } } } } }); } // 加入根节点 queue.Add(root); Interlocked.Increment(ref _activeTasks); // 等待队列完成 queue.CompleteAdding(); queue.GetConsumingEnumerable().ToList(); // 确保所有剩余元素被处理 Console.WriteLine("All traversal completed!"); } } // TreeNode类同方案一
方案三:用TPL Dataflow(更简洁的开箱即用方案)
如果你可以使用TPL Dataflow组件,那这个场景简直是为它量身定做的——TransformBlock可以处理每个输入节点,输出子节点,然后把输出链接回自身形成循环,自带完成检测机制。
示例代码:
using System.Threading.Tasks.Dataflow; class TreeTraversal { static async Task Main() { var root = new TreeNode { Value = "Root", Children = { new TreeNode("Child1"), new TreeNode("Child2") } }; // 创建处理块:输入TreeNode,输出子节点集合 var traversalBlock = new TransformBlock<TreeNode, IEnumerable<TreeNode>>(node => { // 处理当前节点 Console.WriteLine($"Processing {node.Value}"); return node.Children; }, new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = Environment.ProcessorCount // 并行度 }); // 把输出链接回自身,形成循环处理子节点 traversalBlock.LinkTo(traversalBlock, new DataflowLinkOptions { PropagateCompletion = true }); // 发布根节点 traversalBlock.Post(root); // 标记块不再接受新输入 traversalBlock.Complete(); // 等待所有处理完成 await traversalBlock.Completion; Console.WriteLine("All traversal completed!"); } } // TreeNode类同方案一
这三个方案里,我个人最推荐TPL Dataflow的方式,因为它把线程调度、队列管理、完成检测都封装好了,不用自己手动处理同步原语,代码更简洁也更少出错。
内容的提问来源于stack exchange,提问作者SuperGSJ
相关产品推荐
相关产品推荐

