如何在C#中高效并行运行代码?动态适配并行数的实现方案
并行执行CPU密集型任务的优化方案
问题场景
当前单线程代码逻辑为:
var worker = new Worker(); while (true) worker.RandomTest();
其中RandomTest是CPU密集型方法,仅极少数情况涉及文件写入(可忽略I/O延迟);Worker类非线程安全,必须为每个执行单元创建独立实例。
尝试用Parallel.For实现并行时遇到核心问题:
- 若用大数值上限替代
while循环,会在每次迭代重复创建Worker,无法复用实例;且无法提前关联线程与Worker实例。 - 改用固定次数
Parallel.For包裹内部无限循环时:
Parallel.For(0, 1000, (i, _) => { var worker = new Worker(); while (true) worker.RandomTest(); });
Parallel会无限制创建迭代任务,导致线程池饱和,已有任务陷入停滞,无法匹配CPU核心数执行。
需求:实现适配当前机器/环境(如后台存在其他CPU密集进程)的合理并行度,每个并行任务持有独立Worker实例,持续执行无限循环的CPU密集任务。
解决方案
1. 基于Task.Run动态匹配CPU核心数
利用Environment.ProcessorCount获取当前机器逻辑核心数,创建对应数量的长期任务,每个任务仅初始化一次Worker并执行无限循环:
// 可根据后台负载调整,比如取Environment.ProcessorCount - 1 var parallelCount = Environment.ProcessorCount; var tasks = new Task[parallelCount]; for (int i = 0; i < parallelCount; i++) { tasks[i] = Task.Run(() => { var worker = new Worker(); while (true) { worker.RandomTest(); } }); } // 无限期等待所有任务执行 Task.WaitAll(tasks);
- 优势:直接匹配硬件核心数,避免
Parallel的过度调度问题;每个Worker实例仅创建一次,无重复初始化开销。 - 适配优化:若后台有其他CPU密集进程,可动态降低
parallelCount,或通过Process.GetCurrentProcess().ProcessorAffinity绑定特定核心进一步优化。
2. 自定义并行调度器(进阶)
若需要更精细的并发控制,可自定义TaskScheduler限制并发数,替代默认线程池的动态调度:
public class LimitedConcurrencyTaskScheduler : TaskScheduler { private readonly LinkedList<Task> _tasks = new LinkedList<Task>(); private readonly int _maxConcurrencyLevel; private int _currentWorkers = 0; public LimitedConcurrencyTaskScheduler(int maxConcurrencyLevel) { if (maxConcurrencyLevel < 1) throw new ArgumentOutOfRangeException(nameof(maxConcurrencyLevel)); _maxConcurrencyLevel = maxConcurrencyLevel; } protected override IEnumerable<Task> GetScheduledTasks() { lock (_tasks) { return _tasks.ToArray(); } } protected override void QueueTask(Task task) { lock (_tasks) { _tasks.AddLast(task); if (_currentWorkers < _maxConcurrencyLevel) { _currentWorkers++; NotifyThreadPoolOfPendingWork(); } } } protected override bool TryExecuteTaskInline(Task task, bool taskWasPreviouslyQueued) { if (_currentWorkers >= _maxConcurrencyLevel) return false; return TryExecuteTask(task); } private void NotifyThreadPoolOfPendingWork() { ThreadPool.UnsafeQueueUserWorkItem(_ => { while (true) { Task task; lock (_tasks) { if (_tasks.Count == 0) { _currentWorkers--; break; } task = _tasks.First.Value; _tasks.RemoveFirst(); } TryExecuteTask(task); } }, null); } public override int MaximumConcurrencyLevel => _maxConcurrencyLevel; }
使用方式:
var scheduler = new LimitedConcurrencyTaskScheduler(Environment.ProcessorCount); var tasks = new List<Task>(); using (var taskFactory = new TaskFactory(scheduler)) { for (int i = 0; i < scheduler.MaximumConcurrencyLevel; i++) { tasks.Add(taskFactory.StartNew(() => { var worker = new Worker(); while (true) { worker.RandomTest(); } })); } } Task.WaitAll(tasks.ToArray());
- 优势:完全掌控并发数,避免线程池动态调度的干扰;适合对并行度有严格要求的场景。
3. 理解Parallel的设计误区
Parallel.For/Parallel.ForEach的设计目标是处理有限数量的独立短任务,而非长期运行的无限循环。当内部包含while(true)时,Parallel的调度逻辑会持续尝试创建新迭代任务,导致线程池资源耗尽,进而阻塞已有任务执行——这是你之前代码出现停滞的核心原因。
内容的提问来源于stack exchange,提问作者greenoldman
相关产品推荐
相关产品推荐

