123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184 |
- using System;
- using System.Collections.Generic;
- using System.Linq;
- using System.Text;
- using System.Threading;
- using System.Threading.Tasks;
- namespace HttpDownload
- {
- class LimitedConcurrencyLevelTaskScheduler : TaskScheduler
- {
-
- [ThreadStatic]
- private static bool _currentThreadIsProcessingItems;
-
- private readonly LinkedList<Task> _tasks = new LinkedList<Task>();
-
- private readonly int _maxDegreeOfParallelism;
-
- private int _delegatesQueuedOrRunning = 0;
-
-
-
-
-
- public LimitedConcurrencyLevelTaskScheduler(int maxDegreeOfParallelism)
- {
- if (maxDegreeOfParallelism < 1) throw new ArgumentOutOfRangeException("maxDegreeOfParallelism");
- _maxDegreeOfParallelism = maxDegreeOfParallelism;
- }
-
-
- protected sealed override void QueueTask(Task task)
- {
-
-
- lock (_tasks)
- {
- _tasks.AddLast(task);
- if (_delegatesQueuedOrRunning < _maxDegreeOfParallelism)
- {
- ++_delegatesQueuedOrRunning;
- NotifyThreadPoolOfPendingWork();
- }
- }
- }
-
-
-
- private void NotifyThreadPoolOfPendingWork()
- {
- ThreadPool.UnsafeQueueUserWorkItem(_ =>
- {
-
-
- _currentThreadIsProcessingItems = true;
- try
- {
-
- while (true)
- {
- Task item;
- lock (_tasks)
- {
-
-
- if (_tasks.Count == 0)
- {
- --_delegatesQueuedOrRunning;
- break;
- }
-
- item = _tasks.First.Value;
- _tasks.RemoveFirst();
- }
-
- base.TryExecuteTask(item);
- }
- }
-
- finally { _currentThreadIsProcessingItems = false; }
- }, null);
- }
-
-
-
-
- protected sealed override bool TryExecuteTaskInline(Task task, bool taskWasPreviouslyQueued)
- {
-
- if (!_currentThreadIsProcessingItems) return false;
-
- if (taskWasPreviouslyQueued) TryDequeue(task);
-
- return base.TryExecuteTask(task);
- }
-
-
-
- protected sealed override bool TryDequeue(Task task)
- {
- lock (_tasks) return _tasks.Remove(task);
- }
-
- public sealed override int MaximumConcurrencyLevel { get { return _maxDegreeOfParallelism; } }
-
-
- protected sealed override IEnumerable<Task> 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);
- }
- }
- }
- class TaskFactoryMananger
- {
-
- public static void Run()
- {
- try
- {
- while (true)
- {
- LimitedConcurrencyLevelTaskScheduler lcts = new LimitedConcurrencyLevelTaskScheduler(10);
- TaskFactory factory = new TaskFactory(lcts);
- Task[] spiderTask = new Task[] {
- factory.StartNew(() =>
- {
-
-
- }),
- factory.StartNew(() =>
- {
- Thread.Sleep(TimeSpan.FromSeconds(3));
-
-
- }),
- factory.StartNew(() =>
- {
- Thread.Sleep(TimeSpan.FromSeconds(5));
-
-
- })
- };
- Task.WaitAll(spiderTask);
- Thread.Sleep(TimeSpan.FromMinutes(1));
- }
- }
- catch (AggregateException ex)
- {
- foreach (Exception inner in ex.InnerExceptions)
- {
-
- }
- }
- }
-
-
-
-
- }
- }
|