当前位置:

首页 > 编程开发 > C#TaskTaskFactory设置最大并行线程数的方法

C#TaskTaskFactory设置最大并行线程数的方法

LimitedConcurrencyLevelTaskScheduler继承TaskScheduler,通过链表存储待执行任务,设定最大并发数。QueueTask将任务加入链表尾部,若当前运行委托数小于最大值则递增计数并调用ThreadPool.UnsafeQueueUserWorkItem启动工作线程。工作线程循环从链表取出任务执行直至队列为空,减少运行计

1. LimitedConcurrencyLevelTaskScheduler 介绍

这个TaskScheduler,接触过.NET并发编程的同学应该不陌生——微软开源的一个任务调度器,代码本身确实不长,逻辑也算直白。不过,有一个问题值得琢磨:它到底是怎么实现并发数限制的?

C#TaskTaskFactory设置最大并行线程数的方法

先把源码贴出来,大家一起熟悉一下。

public class LimitedConcurrencyLevelTaskScheduler : TaskScheduler
{
    /// Whether the current thread is processing work items. 
    [ThreadStatic]
    private static bool _currentThreadIsProcessingItems;
    /// The list of tasks to be executed. 
    private readonly LinkedList _tasks = new LinkedList(); // protected by lock(_tasks) 
                                                                       /// The maximum concurrency level allowed by this scheduler. 
    private readonly int _maxDegreeOfParallelism;
    /// Whether the scheduler is currently processing work items. 
    private int _delegatesQueuedOrRunning = 0; // protected by lock(_tasks) 
    ///  
    /// Initializes an instance of the LimitedConcurrencyLevelTaskScheduler class with the 
    /// specified degree of parallelism. 
    ///  
    /// The maximum degree of parallelism provided by this scheduler. 
    public LimitedConcurrencyLevelTaskScheduler(int maxDegreeOfParallelism)
    {
        if (maxDegreeOfParallelism < 1) throw new ArgumentOutOfRangeException("maxDegreeOfParallelism");
        _maxDegreeOfParallelism = maxDegreeOfParallelism;
    }
    /// 
    /// current executing number;
    /// 
    public int CurrentCount { get; set; }
    /// Queues a task to the scheduler. 
    /// The task to be queued. 
    protected sealed override void QueueTask(Task task)
    {
        // Add the task to the list of tasks to be processed. If there aren't enough 
        // delegates currently queued or running to process tasks, schedule another. 
        lock (_tasks)
        {
            Console.WriteLine("Task Count : {0} ", _tasks.Count);
            _tasks.AddLast(task);
            if (_delegatesQueuedOrRunning < _maxDegreeOfParallelism)
            {
                ++_delegatesQueuedOrRunning;
                NotifyThreadPoolOfPendingWork();
            }
        }
    }
    int executingCount = 0;
    private static object executeLock = new object();
    ///  
    /// Informs the ThreadPool that there's work to be executed for this scheduler. 
    ///  
    private void NotifyThreadPoolOfPendingWork()
    {
        ThreadPool.UnsafeQueueUserWorkItem(_ =>
        {
            // Note that the current thread is now processing work items. 
            // This is necessary to enable inlining of tasks into this thread. 
            _currentThreadIsProcessingItems = true;
            try
            {
                // Process all a vailable items in the queue. 
                while (true)
                {
                    Task item;
                    lock (_tasks)
                    {
                        // When there are no more items to be processed, 
                        // note that we're done processing, and get out. 
                        if (_tasks.Count == 0)
                        {
                            --_delegatesQueuedOrRunning;
                            break;
                        }
                        // Get the next item from the queue 
                        item = _tasks.First.Value;
                        _tasks.RemoveFirst();
                    }
                    // Execute the task we pulled out of the queue 
                    base.TryExecuteTask(item);
                }
            }
            // We're done processing items on the current thread 
            finally { _currentThreadIsProcessingItems = false; }
        }, null);
    }
    /// Attempts to execute the specified task on the current thread. 
    /// The task to be executed. 
    ///  
    /// Whether the task could be executed on the current thread. 
    protected sealed override bool TryExecuteTaskInline(Task task, bool taskWasPreviouslyQueued)
    {
        // If this thread isn't already processing a task, we don't support inlining 
        if (!_currentThreadIsProcessingItems) return false;
        // If the task was previously queued, remove it from the queue 
        if (taskWasPreviouslyQueued) TryDequeue(task);
        // Try to run the task. 
        return base.TryExecuteTask(task);
    }
    /// Attempts to remove a previously scheduled task from the scheduler. 
    /// The task to be removed. 
    /// Whether the task could be found and removed. 
    protected sealed override bool TryDequeue(Task task)
    {
        lock (_tasks) return _tasks.Remove(task);
    }
    /// Gets the maximum concurrency level supported by this scheduler. 
    public sealed override int MaximumConcurrencyLevel { get { return _maxDegreeOfParallelism; } }
    /// Gets an enumerable of the tasks currently scheduled on this scheduler. 
    /// An enumerable of the tasks currently scheduled. 
    protected sealed override IEnumerable GetScheduledTasks()
    {
        bool lockTaken = false;
        try
        {
            Monitor.TryEnter(_tasks, ref lockTaken);
            if (lockTaken) return _tasks.ToArray();
            else throw new NotSupportedException();
        }
        finally
        {
            if (lockTaken) Monitor.Exit(_tasks);
        }
    }
}

简单使用

下面是调用示例,非常简单:

static void Main(string[] args)
{
        TaskFactory fac = new TaskFactory(new LimitedConcurrencyLevelTaskScheduler(5));
        //TaskFactory fac = new TaskFactory();
        for (int i = 0; i < 1000; i++)
        {
            fac.StartNew(s => {
                Thread.Sleep(1000);
                Console.WriteLine("Current Index {0}, ThreadId {1}",s,Thread.CurrentThread.ManagedThreadId);
            }, i);
        }
        Console.ReadKey();
}

调用逻辑很清晰:用 LimitedConcurrencyLevelTaskScheduler 创建 TaskFactory,然后通过 StartNew 提交任务。从调试顺序可以看到,每次 StartNew 都会进入 QueueTask 方法。

/// Queues a task to the scheduler. 
    /// The task to be queued. 
    protected sealed override void QueueTask(Task task)
    {
        // Add the task to the list of tasks to be processed. If there aren't enough 
        // delegates currently queued or running to process tasks, schedule another. 
        lock (_tasks)
        {
            Console.WriteLine("Task Count : {0} ", _tasks.Count);
            _tasks.AddLast(task);
            if (_delegatesQueuedOrRunning < _maxDegreeOfParallelism)
            {
                ++_delegatesQueuedOrRunning;
                NotifyThreadPoolOfPendingWork();
            }
        }
    }
    

QueueTask 的步骤很简单:把新任务追加到链表尾部,然后检查当前正在运行或已排队的委托数量(_delegatesQueuedOrRunning)是否小于设定的最大并发数。如果小于,就递增计数并调用 NotifyThreadPoolOfPendingWork 去启动一个工作线程。

但真正的疑问,恰恰出在这个 NotifyThreadPoolOfPendingWork 方法上。

private void NotifyThreadPoolOfPendingWork()
    {
        ThreadPool.UnsafeQueueUserWorkItem(_ =>
        {
            // Note that the current thread is now processing work items. 
            // This is necessary to enable inlining of tasks into this thread. 
            _currentThreadIsProcessingItems = true;
            try
            {
                // Process all a vailable items in the queue. 
                while (true)
                {
                    Task item;
                    lock (_tasks)
                    {
                        // When there are no more items to be processed, 
                        // note that we're done processing, and get out. 
                        if (_tasks.Count == 0)
                        {
                            --_delegatesQueuedOrRunning;
                            break;
                        }
                        // Get the next item from the queue 
                        item = _tasks.First.Value;
                        _tasks.RemoveFirst();
                    }
                    // Execute the task we pulled out of the queue 
                    base.TryExecuteTask(item);
                }
            }
            // We're done processing items on the current thread 
            finally { _currentThreadIsProcessingItems = false; }
        }, null);
    }

看这个方法的内部逻辑:它直接丢了一个死循环到线程池里,循环内不断从 _tasks 中取出任务执行,直到队列为空才退出循环。这看起来就像是一个“无限吞噬”的过程——一旦启动,就会把所有任务吃光,根本看不到任何限制并发数的机制。

唯一能扯上“限制”的地方,是 QueueTask 里那个 if 判断:只有当前执行线程数小于最大并发度时,才调用 NotifyThreadPoolOfPendingWork。但这似乎没什么用,因为一旦调用,这个工作线程就会一直跑,直到把队列清空。这样一来,并发度不就失控了吗?

那么问题来了:LimitedConcurrencyLevelTaskScheduler 到底是如何实现并发数限制的?

是不是哪里理解有偏差?比如,NotifyThreadPoolOfPendingWork 中 while 循环每次取任务时,会不会因为锁的竞争或其他机制而自然阻塞?但实际上锁只会保护 _tasks 的访问,并不控制线程数量。更关键的是,QueueTask 中的 if 条件保证了同时只有 _maxDegreeOfParallelism 个线程被启动,但每个线程都是“死循环”,这会不会导致任务被一个线程全部执行完,其他线程根本拿不到任务?

仔细想想,死循环本身并不占用多个线程——它只用当前这一个线程。但问题在于,当多个任务同时被提交时,QueueTask 可能被多次调用(来自不同的调用线程),而每次调用如果满足条件都会启动一个新的工作线程。假设并发数设为5,在任务提交的瞬间,如果同时有10个线程调用 QueueTask,前5个会启动工作线程,后5个不会。但前5个工作线程各自进入死循环,彼此独立地从 _tasks 中取任务——这确实实现了5个线程同时消费任务。真正的限制在于:不会启动超过5个工作线程。而死循环保证了每个工作线程会持续消费任务,而不是执行一个就退出,这样即便后续有新任务加入,也无需再启动新线程(因为现有工作线程还在循环中)。

换句话说,这个设计的精巧之处在于:工作线程采用“持续消费”模式,而不是“消费一次就结束”。QueueTask 中的 if 判断确保了最多只有 _maxDegreeOfParallelism 个工作线程存在,而这些线程会一直循环直到队列为空,从而实现了并发度的硬限制。

当然,这个实现有一个潜在的缺陷:如果任务生产速度远大于消费速度,工作线程会一直忙,但一旦队列为空,工作线程退出,后续新任务到来时,如果此时 _delegatesQueuedOrRunning 已经减到小于 _maxDegreeOfParallelism,就会重新启动新工作线程。这个计数更新是在 while 循环退出时(_tasks.Count == 0)进行的,所以是安全的。

所以,回过头来看,这个调度器对并发数的限制,本质上是通过控制“同时运行的工作线程数量”来实现的。虽然每个工作线程内部是个死循环,但死循环的线程数量是固定的,因此并发度也就固定了。

以上是个人理解,不知是否完全准确。如果有不同看法或者更深入的分析,欢迎交流讨论。

本站声明:本文内容由网友自发贡献,版权归原作者所有,本站不承担相应法律责任。如您发现有涉嫌抄袭侵权的内容,请联系bd@zhengruan.com
作者最新文章
编程开发
相关文章 更多
ServBay安装配置详细教程与操作指南
ServBay安装配置详细教程与操作指南

新手入门 ServBay 本地开发环境,详解安装包下载、Dashboard 状态监控、Packages 组件安装、Services 服务控制及 Websites 项目配置。掌握 .servbay.config 版本管理与日志排查技巧,快速搭建稳定的 PHP、Node.js 等多语言开发环境。

codekit环境配置指南从安装到环境搭建完整教程
codekit环境配置指南从安装到环境搭建完整教程

详解 CodeKit 在 macOS 下的安装步骤、项目导入方法、Sass与JavaScript编译设置及浏览器自动刷新功能,助您快速搭建高效的前端开发环境。

codex安装windows 命令行完整操作教程
codex安装windows 命令行完整操作教程

详解Windows环境下安装OpenAI Codex CLI的步骤,包括WSL环境检查、Node.js/npm配置、npm全局安装命令及首次启动验证,适合开发者快速上手。

NativeRest环境配置要求与完整操作教程
NativeRest环境配置要求与完整操作教程

学习如何配置 NativeRest REST API 客户端。涵盖 Windows/macOS/Linux 安装后的工作区创建、环境变量管理、请求编辑及响应查看步骤,帮助开发者快速完成基础环境搭建与连通性测试。

CSS设置透明度的注意事项有哪些?opacity属性详解
CSS设置透明度的注意事项有哪些?opacity属性详解

深入解析CSS中设置透明度的核心属性opacity,剖析子元素继承、事件穿透、层叠上下文等关键注意事项,并提供与rgba、hsla的实用选型对比。

flutter页面传值到后台的方法及示例代码
flutter页面传值到后台的方法及示例代码

flutter页面传值到后台的完整实现方法及示例代码,帮助读者快速掌握相关技术要点。

Java 8至21新特性代码写法对比:Lambda、Record与Switch
Java 8至21新特性代码写法对比:Lambda、Record与Switch

本文通过具体的旧版与新版代码对比,详细剖析Java 8引入的Lambda表达式、Java 14/16引入的Record类,以及Java 12至21逐步演进完善的Switch表达式与模式匹配,展示代码简化路径与避坑要点。

AI智能体开发培训课程学什么及实战内容介绍
AI智能体开发培训课程学什么及实战内容介绍

系统梳理AI智能体开发培训的核心知识模块、技术栈选型与典型实战项目,解析低代码平台与纯代码框架的差异,提供从零构建可落地智能体的完整学习与实施路径。

Java子类未实现抽象方法编译错误修复指南
Java子类未实现抽象方法编译错误修复指南

针对Java开发中常见的“子类未实现抽象方法”编译错误,深入分析报错原因,提供重写实现、声明抽象子类两种标准修复路径,并总结参数签名、访问修饰符等典型避坑要点。

解决PHP递归报错:max_nesting_level限制与内存溢出处理
解决PHP递归报错:max_nesting_level限制与内存溢出处理

遇到PHP递归报错时,不要盲目调大max_nesting_level。本文教你区分Xdebug限制、内存耗尽和正则递归错误,提供代码级的终止条件优化与迭代替代方案,彻底解决栈溢出问题。

查看更多
精品专题 更多
装机必备
装机必备

正软商城装机必备专区,精选办公、浏览器、安全防护、影音播放、压缩解压、设计创作和系统工具等电脑常用正版软件,帮助用户快速完成新电脑软件配置。

Windows
Windows

正软商城Windows软件专区,汇集适用于Windows电脑的办公、设计、安全防护、影音播放、开发工具和系统优化软件,提供软件介绍、系统要求、正版授权及购买下载服务。

PDF教程
PDF教程

正软商城PDF教程频道提供PDF编辑、转换、合并、拆分、压缩及格式处理方法,同时介绍常用PDF软件和工具的使用技巧。

Mac软件 更多
Shapr3D macOS版
Shapr3D macOS版
Mac

Shapr3D是一款面向工业设计、机械工程、建筑概念和三维打印工作流的CAD软件。Mac版采用Parasolid建模内核,支持草图约束、实体建模、工程图、可视化渲染及常见CAD格式交换,并可通过账户在多台设备之间同步项目。

REAPER macOS版
REAPER macOS版
Mac

REAPER是Cockos开发的数字音频工作站,提供多轨音频与MIDI录制、剪辑、处理、混音和母带制作工具。Mac版兼容Intel与Apple芯片,支持AU、VST、VST3、CLAP等插件格式,并提供高度可定制的工作流程。

Ableton Live macOS版
Ableton Live macOS版
Mac

Ableton Live 是面向音乐制作人与现场表演者的数字音频工作站,提供编曲视图、独具特色的现场视图、音频录制、MIDI创作、实时变速、乐器及效果器。Mac版原生支持Apple芯片,并可连接音频接口、MIDI控制器和第三方插件。

WINDOWS 更多
3dmax(3ds max)
3dmax(3ds max)
Windows

Autodesk 3ds Max 是一款专业的三维建模、动画与渲染软件,广泛应用于建筑可视化、游戏开发、影视动画、广告设计和产品展示等领域。

photoshop
photoshop
Windows、macOS 、 iPad

Photoshop 2026 是 Adobe 推出的专业图像处理与视觉设计软件,支持 Windows、macOS 和 iPad 等平台,广泛应用于摄影修图、电商设计、平面海报、数字绘画及视觉合成等创作场景。

Blender
Blender
Windows、macOS 和 Linux

Blender 是一款免费开源、跨平台的专业 3D 创作软件,集建模、动画、渲染、视频编辑与视觉合成等功能于一体,广泛应用于影视动画、游戏设计和建筑可视化等领域。软件支持 Cycles 物理渲染器与 Eevee 实时渲染引擎,并提供多边形建模、骨骼绑定、物理模拟等专业工具。Blender 兼容 Windows、macOS 和 Linux 系统,安装包轻巧、运行流畅,依托活跃的全球开发者社区持续更新,是从初学者到专业创作者都值得选择的正版 3D 创作工具。