ab56a9bcf7
SVN-Revision: r240
79 lines
2.8 KiB
C#
79 lines
2.8 KiB
C#
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<Task> _tasks = new LinkedList<Task>(); //任务队列
|
|
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<Task> GetScheduledTasks()
|
|
{
|
|
lock (_lock)
|
|
{
|
|
return _tasks.ToList();
|
|
}
|
|
}
|
|
}
|
|
}
|