using System; using System.Collections.Generic; using System.Linq; using System.Text; using System.Threading; using System.Threading.Tasks; namespace Lskj.Util { public class TaskSchedulerUtil : TaskScheduler { private readonly LinkedList _tasks = new LinkedList(); //任务队列 private readonly int _maxDegreeOfParallelism; //最大并发线程数 private int _runningTasks = 0; //当前正在运行的线程数 private readonly object _lock = new object(); //用于线程同步 public TaskSchedulerUtil(int maxDegreeOfParallelism) { if (maxDegreeOfParallelism < 1) throw new ArgumentOutOfRangeException(nameof(maxDegreeOfParallelism)); _maxDegreeOfParallelism = maxDegreeOfParallelism; } // 将任务加入队列 protected override void QueueTask(Task task) { lock (_lock) { _tasks.AddLast(task); // 将任务加入队列 TryExecuteNextTask(); // 尝试执行下一任务 } } // 尝试执行下一个任务 private void TryExecuteNextTask() { lock (_lock) { // 如果当前运行的线程数小于最大并发数,且队列中还有任务 while (_runningTasks < _maxDegreeOfParallelism && _tasks.Count > 0) { var task = _tasks.First.Value; _tasks.RemoveFirst(); // 取出一个任务 _runningTasks++; // 增加正在运行的任务计数 // 在线程池中运行任务 ThreadPool.QueueUserWorkItem(_ => { try { base.TryExecuteTask(task); // 执行任务 } finally { lock (_lock) { _runningTasks--; // 减少正在运行的任务计数 TryExecuteNextTask(); // 尝试执行下一个任务 } } }); } } } // 尝试在线程中直接运行任务(不支持) protected override bool TryExecuteTaskInline(Task task, bool taskWasPreviouslyQueued) { // 不允许直接在线程中运行任务(可以根据需求修改) return false; } // 返回已计划的任务(用于调试) protected override IEnumerable GetScheduledTasks() { lock (_lock) { return _tasks.ToList(); } } } }