Parallel.ForEach结合BlockingCollection导致线程池线程持续增长的排查与解决
现象描述
测试发现,在单独线程中执行以下代码时,线程池线程数会持续增长:
var queue = new BlockingCollection<int>(); Parallel.ForEach(queue.GetConsumingEnumerable(), _ => { });
此时队列是空的,预期Parallel.ForEach会空闲等待新元素,但实际ThreadPool.ThreadCount每秒都会增加,最小复现代码及输出如下:
最小复现代码
public class Program { public static void Main() { new Thread(() => { var queue = new BlockingCollection<int>(); Parallel.ForEach(queue.GetConsumingEnumerable(), _ => { }); }) { IsBackground = true }.Start(); Stopwatch stopwatch = Stopwatch.StartNew(); while (true) { Console.WriteLine($"ThreadCount: {ThreadPool.ThreadCount}"); if (stopwatch.ElapsedMilliseconds > 8000) break; Thread.Sleep(1000); } Console.WriteLine("Finished"); } }
输出结果
ThreadCount: 0 ThreadCount: 4 ThreadCount: 5 ThreadCount: 6 ThreadCount: 7 ThreadCount: 8 ThreadCount: 10 ThreadCount: 11 ThreadCount: 12 Finished
原因分析
Parallel.ForEach的默认调度逻辑是试探性扩容线程数:它会不断尝试增加线程,直到任务吞吐量不再提升,以此最大化CPU利用率。但BlockingCollection.GetConsumingEnumerable()在队列空时,会让调用线程进入阻塞状态(等待生产者添加元素)。
Parallel.ForEach的调度器无法区分"任务真的在执行"和"线程被阻塞等待",会把阻塞判定为"任务执行缓慢",于是持续从线程池创建新线程来尝试处理"更多任务"。这些新线程最终也会因为队列空而阻塞,导致线程池线程数无限制增长。
简单来说:Parallel.ForEach的设计目标是处理CPU密集型、无阻塞的枚举源,不适合处理会阻塞等待的生产者-消费者队列。
解决方案(适用于.NET Core 3.1+)
方案1:限制Parallel.ForEach的最大并行度
通过ParallelOptions指定MaxDegreeOfParallelism,直接限制最多使用的线程数。如果要求空队列时仅占用1个线程,设为1即可:
var options = new ParallelOptions { MaxDegreeOfParallelism = 1 }; Parallel.ForEach(queue.GetConsumingEnumerable(), options, _ => { });
如果需要支持并行处理,可以把值设为Environment.ProcessorCount或其他合适的数值,但要注意:空队列时所有并行线程都会阻塞,会占用对应数量的线程池线程。
方案2:手动控制消费逻辑(推荐)
改用"单线程监听队列+线程池并行处理"的模式,空队列时仅1个监听线程,有元素时才按需提交任务到线程池:
var queue = new BlockingCollection<int>(); // 限制并行度的信号量,可根据需求调整数量 var semaphore = new SemaphoreSlim(Environment.ProcessorCount); new Thread(() => { foreach (var item in queue.GetConsumingEnumerable()) { semaphore.Wait(); // 提交任务到线程池 _ = Task.Run(() => { try { // 这里编写处理item的业务逻辑 } finally { semaphore.Release(); } }); } }) { IsBackground = true }.Start();
这种方式完全满足需求:空队列时仅占用1个线程,有元素时按设定的并行度处理,不会无限制扩容线程池。
方案3:使用Channel替代(更优雅)
.NET Core 3.1及以上提供的System.Threading.Channels是专门为生产者-消费者场景设计的,配合Parallel.ForEachAsync(.NET 6+)可以更优雅地实现需求,且不会出现线程池无限制增长的问题:
using System.Threading.Channels; var channel = Channel.CreateUnbounded<int>(); // 消费者逻辑 _ = Task.Run(async () => { var options = new ParallelOptions { MaxDegreeOfParallelism = Environment.ProcessorCount }; await Parallel.ForEachAsync(channel.Reader.ReadAllAsync(), options, async (item, token) => { // 处理item的逻辑,支持异步操作 await Task.Delay(100, token); }); }); // 生产者逻辑示例 _ = Task.Run(async () => { for (int i = 0; i < 100; i++) { await channel.Writer.WriteAsync(i); await Task.Delay(50); } channel.Writer.Complete(); });
Channel的读取逻辑不会让线程无意义阻塞,Parallel.ForEachAsync的调度更适配异步场景,能高效利用线程池资源。
内容的提问来源于stack exchange,提问作者Theodor Zoulias

