You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在C#中高效并行运行代码?动态适配并行数的实现方案

并行执行CPU密集型任务的优化方案

问题场景

当前单线程代码逻辑为:

var worker = new Worker();
while (true)
    worker.RandomTest();

其中RandomTest是CPU密集型方法,仅极少数情况涉及文件写入(可忽略I/O延迟);Worker类非线程安全,必须为每个执行单元创建独立实例。

尝试用Parallel.For实现并行时遇到核心问题:

  1. 若用大数值上限替代while循环,会在每次迭代重复创建Worker,无法复用实例;且无法提前关联线程与Worker实例。
  2. 改用固定次数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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.08 10:15:11