Параллельное программирование для многоядерных процессоров

Программирование с использованием Task Parallel Library (TPL)

Показывать лекцию целиком

Как было отмечено во введении, в библиотеке PFX поддерживаются три уровня параллелизма:

  • уровень декларативной обработки данных (PLINQ, см. Лекцию 7),
  • уровень императивной обработки данных (конструкции Parallel.For/ForEach/Invoke, см. Лекции 2 и 4), а также
  • уровень императивной работы с задачами (tasks).
  • Последний уровень является базовым в PFX, имеет специальное название - Task Parallel Library (TPL), и на его основе могут быть реализованы два других уровня. Программирование на уровне задач требует больших усилий со стороны программиста, однако, этот уровень является универсальным, а потому, более гибким - с помощью задач можно построить любое параллельное приложение в отличие от специальных шаблонов Parallel.For/ForEach/Invoke.

    На первый взгляд, класс System.Threading.Tasks.Task является аналогом пула потоков в .NET (имеется в виду класс System.Threading.ThreadPool ). Действительно, общность принципов работы с этими классами налицо - например, запросить у пула поток для выполнения программа может следующим образом:

    ThreadPool.QueueUserWorkItem(delegate { ... });

    Работа с классом Task в этом случае будет выглядеть так:

    Task.Create(delegate { ... });

    Однако, класс Task и в целом библиотека TPL представляет несравнимо большие возможности, чем стандартный пул потоков. Например, если вы хотите параллельно выполнить три задачи и дождаться завершения их выполнения, то используя пул потоков вы можете написать:

    //создадим три события "Выполнение задачи завершено"
    using(ManualResetEvent mre1 = new ManualResetEvent(false))
    using(ManualResetEvent mre2 = new ManualResetEvent(false))
    using(ManualResetEvent mre3 = new ManualResetEvent(false))
          {
    //запросим у пула потоков параллельное исполнение 3-х задач
              ThreadPool.QueueUserWorkItem(delegate
              {
                  A();
                  mre1.Set();
              });
              ThreadPool.QueueUserWorkItem(delegate
              {
                  B();
                  mre2.Set();
              });
              ThreadPool.QueueUserWorkItem(delegate
              {
                  C();
                  mre3.Set();
              });
     	//дождемся выполнения всех задач
              WaitHandle.WaitAll(new WaitHandle[]{mre1, mre2, mre3});
          }

    Аналогичный код с использованием TPL будет выглядеть следующим образом:

    Task t1 = Task.Create(delegate { A(); });
    Task t2 = Task.Create(delegate { B(); });
    Task t3 = Task.Create(delegate { C(); });
    t1.Wait();
    t2.Wait();
    t3.Wait();

    Будем называть код делегата типа Action<Object>, передаваемого при создании задачи, телом задачи.

    Код выше можно переписать немного элегантнее:

    Task t1 = Task.Create(delegate { A(); });
    Task t2 = Task.Create(delegate { B(); });
    Task t3 = Task.Create(delegate { C(); });
    Task.WaitAll(t1, t2, t3);

    При этом нужно понимать, что реально исполнение задачи одним из рабочих потоков начнется не сразу же после вызова метода Task.Create, а через некоторое время. То есть, так же как и в случае с ThreadPool, при своем создании задача просто помещается в очередь планировщика. Решение о моменте запуска задачи на исполнение потоком принимает планировщик в соответствии с дисциплиной планирования исполнения задач. Более подробно о дисциплинах и принципах планирования см. Лекцию 3.

    Внимательный читатель может заметить, что библиотека PFX позволяет еще проще реализовать параллельный запуск задач с помощью метода Parallel.Invoke:

    Parallel.Invoke( ()=>A() , ()=>B() , ()=>C() );

    Иногда между запусками задач необходимо выполнить какие-то дополнительные действия, в этом случае удобнее работать непосредственно с классом Task:

    ... // подготовительная работа для задачи А
    Task t1 = Task.Create(delegate { A(); });
    ... // подготовительная работа для задачи B
    Task t2 = Task.Create(delegate { B(); });
    ... // подготовительная работа для задачи C
    Task t3 = Task.Create(delegate { C(); });
    Task.WaitAll(t1, t2, t3);

    Выше мы рассмотрели примеры порождения новых задач и ожидания их завершения. Существуют и другие возможности управления исполнением задач. Так, например, для того чтобы отменить выполнение задачи можно использовать метод Task.Cancel. При вызове данного метода возможны два варианта:

  • Если задача не была отправлена планировщиком на исполнение каким-либо рабочим потоком, то произойдет изъятие задачи из очереди планировщка
  • Если задача уже исполняется каким-либо рабочим потоком, то из-за соображений надежности ее исполнение не будет прервано, однако свойство Task.IsCanceled примет значение True, что позволит коду задачи, проанализировав значение этого свойства, завершить свое исполнение.
  • Аналогичный метод Task.CancelAndWait позволит дождаться остановки исполнения задачи (т.е., тогда как метод Task.Cancel является неблокирующим, то метод Task.CancelAndWait заблокирует вызвавший поток до тех пор, пока соответствующая задача не будет реально снята).

    Для того чтобы из тела задачи получить доступ к объекту Task можно воспользоваться статическим свойством Task.Current:

    static void Main(string[] args)
    {
          Task a = Task.Current;
          //a==null
          Task.Create(delegate { A(); });
    }
    
    //выведет "False"
    static void A() { Console.WriteLine(Task.Current.IsCompleted); }

    Рассмотрим свойства задачи Task.Parent и Task.Creator. Эти свойства отвечают за иерархические связи на множестве задач. Свойство Task.Creator указывает на задачу в теле которой была создана данная задача, другими словами Task.Creator равно Task.Current материнской задачи. Свойство Task.Parent совпадает со свойством Task.Creator, за исключением того случая, когда при создании задачи была указана опция TaskCreationOptions.Detached (в последнем случае, у задачи нет "родителя"). Отметим, что задача завершает свое исполнение тогда и только тогда, когда все ее дочерние задачи также завершают свое исполнение:

    Task p = Task.Create(delegate 
    { 
        Task c1 = Task.Create(...); 
        Task c2 = Task.Create(...); 
        Task c3 = Task.Create(...); 
    }); 
    ... 
    p.Wait(); // ожидание завершения задач p, c1, c2 и c3

    В текущей версии библиотеки PFX work-stealing планировщик представлен классом System.Threading.Tasks.TaskManager. При создании объекта TaskManager программст может указать ряд значений параметров, которые будут влиять на планирование исполнения задач, переданных данному планировщику. Программист может создавать несколько планировщиков, а при создании задачи указывать какой планировщик будет контролировать исполнение задачи.

    Параметры планировщика описываются классом TaskManagerPolicy. Перечислим эти параметры:

  • minProcessors - минимальное число процессоров, используемое планировщиком. По умолчанию - 1 процессор;
  • idealProcessors - оптимальное число процессоров, используемое планировщиком. По умолчанию - количество доступных процессоров в системе;
  • idealThreadsPerProcessor - оптимальное число потоков, создаваемых планировщиком для каждого процессора. По умолчанию - один поток на один процессор;
  • maxStackSize - максимальный объем стека для рабочих потоков;
  • threadPriority - приоритет рабочих потоков. По умолчанию - ThreadPriority.Normal ;
  • Задачи являются одним из базовых элементов библиотеки PFX, другим базовым элементом являются Future, работа с которыми будет описана в лекции 6.

    Замечание:

    Отметим, что начиная с версии .NET 4 CTP TPL входит в состав сборки mscorlib.dll, поэтому теперь любое приложение сможет иметь доступ к TPL без необходимости включения в свой состав дополнительных сборок.

    Семинарское занятие № 5. Рекурсия и параллелизм (часть 2)

    На предыдущем семинарском занятии было рассмотрено несколько методов рекурсивной параллельной обработки (бинарных) деревьев. Хотя последний из приведенных методов обладает достаточной степенью параллелизма, он имеет ряд серьезных недостатков:

  • вызов метода Process в рамках потока из ThreadPool блокирует этот поток до тех пор, пока не будут обработаны вершины-потомки для данной вершины-родителя; таким образом, данная реализация требует одного потока из ThreadPool на одну вершину обрабатываемого дерева; однако, класс ThreadPool имеет ограниченное количество потоков, а именно 25 потоков на ядро (процессор) для .NET 1.x/2.0, и 250 потоков на процессор для .NET 2.0 SP1; таким образом, если из пула выбраны все потоки, то новые потоки не могут быть созданы, а потому приложение может перейти в состояние дедлока;
  • создание нового потока в ThreadPool, когда количество уже созданных потоков больше или равно количеству доступных ядер (процессоров), занимает около 500 мсек; а потому обработка дерева, например, из 250 вершин займет более 2-х минут, потраченных только на создание потоков;
  • для каждого вновь создаваемого потока в .NET отводится около 1 Мб (виртуальной) памяти;
  • для обработки каждой вершины дерева создается один объект класса ManualResetEvent, операции над которым Set и WaitOne выполняются ядром операционной системы, а потому являются дорогостоящими в смысле используемых ресурсов.
  • Чтобы избавиться от некоторых из этих недостатков, ниже показана реализация, которая основана на последовательном проходе по дереву и записи в очередь на обработку в пул потоков задания на обработку каждой вершины в дереве.

    public static void Process<T> (Tree<T> tree, Action<T> action)
    {
        if (tree == null) return;
    
        // Использование события для ожидания завершения 
        // обработки всех вершин дерева
    
        using (var mre = new ManualResetEvent(false))
        {
            int count = 1;
            
            // Рекурсивный делегат для прохода по дереву
    
            Action<Tree<T> > processNode = null;
            processNode = node =>
            {
                if (node == null) return;
     
                // Асинхронный запуск обработки текущей вершины
    
                Interlocked.Increment(ref count);
                ThreadPool.QueueUserWorkItem(delegate
                {
                    action(node.Data);
                    if (Interlocked.Decrement(ref count) == 0) 
                        mre.Set();
                });
     
                // Обработка потомков
    
                processNode(node.Left);
                processNode(node.Right);
            };
     
            // Запуск обработки, начиная с корневой вершины
    
            processNode(tree);
     
            // Сигнал о том, что заданий на обработку больше
            // создаваться не будет
    
            if (Interlocked.Decrement(ref count) == 0) mre.Set();
    
            // Ожидание завершения обработки всех вершин
    
            mre.WaitOne();
        }
    }

    Задача 5.

  • Объясните почему данная реализация свободна от дедлоков в отличие от предыдущей реализации.
  • Объясните возможен ли вариант данной реализации, в котором исходным значением счетчика является 0, т.е.,

    int   count  =  0;

    а в методе Process последний оператор

    if ( Interlocked.Decrement ( ref count ) == 0 ) mre.Set();

    опущен.

  • Рекурсивный подход, приведенный в предыдущем примере, можно реализовать итеративно:

    public static void Process<T> (Tree<T> tree, Action<T> action)
    {
        if (tree == null) return;
    
        // Использование события для ожидания завершения
        // обработки всех вершин дерева
    
        using (var mre = new ManualResetEvent(false))
        {
            int count = 1;
    
            // Запуск обработки, начиная с корневой вершины
    
            var toExplore = new Stack<Tree<T> > ();
            toExplore.Push(tree);
    
            // Обработка всех вершин 
    
            while (toExplore.Count > 0)
            {
                // Извлечь текущую вершину и поместить в стек
                // ее потомков
    
                var current = toExplore.Pop();
                if (current.Left != null) 
                    toExplore.Push(current.Left);
                if (current.Right != null) 
                    toExplore.Push(current.Right); 
    
                // Асинхронная обработка данных
    
                Interlocked.Increment(ref count);
                ThreadPool.QueueUserWorkItem(delegate
                {
                    action(current.Data);
                    if (Interlocked.Decrement(ref count) == 0) 
                        mre.Set();
                });
            }
    
            // Сигнал о том, что больше заданий на обработку
            // создаваться не будет
    
            if (Interlocked.Decrement(ref count) == 0) mre.Set();
    
            // Ожидание завершения обработки всех верщин
    
            mre.WaitOne();
        }
    }

    Недостаток предыдущего варианта состоит в том, что для каждой вершины дерева в ThreadPool помещается одно задание для обработки этой вершины. Гораздо экономичнее, а потому, эффективнее будет создать только $$N$$ заданий на обработку (где $$N$$ равно числу доступных процессоров на машине), и где каждое задание будет состоять в обработке примерно $$1/N$$ -ой части всех вершин дерева. Такой подход можно реализовать, например, сохранив все вершины дерева в виде списка, и разделив его на $$N$$ частей, для обработки каждой части в виде параллельного задания:

    public static void Process<T> (Tree<T> tree, Action<T> action)
    {
        if (tree == null) return;
     
        // Создать список всех вершин дерева
    
        var nodes = new List<Tree<T> > ();
        var toExplore = new Stack<Tree<T> > ();
        toExplore.Push(tree);
        while (toExplore.Count > 0)
        {
            var current = toExplore.Pop();
            nodes.Add(current);
            if (current.Left != null) 
                toExplore.Push(current.Left);
            if (current.Right != null) 
                toExplore.Push(current.Right);
        }
     
        // Разбиение списка на части
    
        int workItems = Environment.ProcessorCount;
        int chunkSize = Math.Max( nodes.Count / workItems, 1);
        int count = workItems;
    
        // Использование события для ожидания завершения
        // обработки всех заданий
    
        using (var mre = new ManualResetEvent(false))
        {
            // В каждом задании обрабатывается примерно 1/N-я
            // часть вершин
    
            WaitCallback callback = state =>
            {
                int iteration = (int)state;
                int from = chunkSize * iteration;
                int to = iteration == workItems - 1 ?
                    nodes.Count : chunkSize * (iteration + 1);
    
                while (from < to) action(nodes[from++].Data);
    
                if (Interlocked.Decrement(ref count) == 0) 
                    mre.Set();
            };
    
            // ThreadPool используется для обработки N - 1 части; 
            // для обработки последней части используется
            // текущий поток 
            
            for (int i = 0; i<workItems; i++)
            {
                if (i < workItems-1) 
                    ThreadPool.QueueUserWorkItem(callback, i);
                else 
                    callback(i);
            }
            // Ожидание завершения обработки всех заданий
    
            mre.WaitOne();
    
        }
    }

    Задача 6.

    Проверьте, какой величины получатся задания для каждого процессора, если обрабатываемое дерево имеет 63 вершины, а количество процессоров равно 8.

    Придумайте и реализуйте более оптимальный способ распределения вершин по заданиям, такой, чтобы размер любой пары заданий отличался не более, чем на 1.

    Статическое разбиение списка на части также может быть не оптимальным способом деления работы между потоками. Альтернативным вариантом является динамическое разбиение, когда потоки конкурируют между собой за получение задания на обработку очередной вершины дерева после окончания обработки предыдущей вершины. Подход на основе динамического разбиения можно реализовать следующим способом:

    public static void Process<T> (Tree<T> tree, Action<T> action)
    {
        if (tree == null) return;
     
        // Получение конструкции перечисления для дерева
    
        IEnumerator<T> enumerator = 
            GetNodes(tree).GetEnumerator();
        int workItems = Environment.ProcessorCount;
        int count = workItems;
     
        // Использование события для ожидания завершения
        // работы всех потоков
    
        using (var mre = new ManualResetEvent(false))
        {
            // Каждый поток будет получать данные для обработки, 
            // используя механизм перечисления до тех пор,
            // пока не будут исчерпаны все вершины дерева
    
            WaitCallback callback = delegate
            {
                while (true)
                {
                    T data;
                    lock (enumerator)
                    {
                        if (!enumerator.MoveNext()) break;
                        data = enumerator.Current;
                    }
                    action(data);
                }
                if (Interlocked.Decrement(ref count) == 0) 
                    mre.Set();
            };
     
            // Из ThreadPool'а берется всего N - 1 потоков;
            // кроме того, для обработки используется  
            // текущий поток
    
            for (int i = 0; i < workItems; i++)
            {
                if (i < workItems-1)
                    ThreadPool.QueueUserWorkItem(callback, i);
                else
                    callback(i);
            }
    
            // Ожидание завершения работы всех потоков
    
            mre.WaitOne();
        }
    }
    
    // Конструкция перечисления вершин дерева
    
    public static IEnumerable<T> GetNodes<T> (Tree<T> tree)
    {
        if (tree != null)
        {
            yield return tree.Data;
            foreach (var data in GetNodes(tree.Left)) 
                yield return data;
            foreach (var data in GetNodes(tree.Right)) 
                yield return data;
        }
    }

    Недостаток последнего варианта состоит в том, что если обработка каждой вершины состоит из небольшого количества вычислительных операций, то общая производительность приложения будет низка из-за частого использования конструкции запирания lock. Поэтому, в некоторых случаях более оптимальным может оказаться вариант при котором при каждом применении lock к enumerator потоком извлекается сразу несколько вершин для обработки.

    Задача 7.

    Пусть для предыдущего примера дополнительным аргументом метода Process является целочисленный параметр $$p\le1$$, задающий количество вершин, извлекаемых из enumerator'а, при каждом применении к нему оператора lock. Реализуйте метод Process с указанным дополнением.

    Возможны и некоторые другие способы распределения вершин (бинарного) дерева между потоками для обработки.

    Задача 8.

    Пусть $$P$$ есть количество доступных ядер/процессоров на машине. Предположим, что $$P = 2^n$$ для некоторого $$n\le0$$. Напишите программу обработки всех вершин бинарного дерева, в которой главный поток сам обрабатывает все вершины дерева, начиная с корня, и до глубины $$d = log_2 P$$, и порождает дополнительные потоки для обработки соответствующих поддеревьев на этой глубине.

    Пример с обходом вершин дерева был достаточно прост в том смысле, что не требовалось блокировать обработку родительской вершины до тех пор, пока не будут обработаны ее вершины-потомки, и, наоборот, обработку вершин-потомков можно было начинать, не дожидаясь окончания обработки родительской вершины.

    Однако, для многих рекурсивных задач такая автономная схема обработки невозможна. Рассмотрим, для примера, типичный рекурсивный алгоритм сортировки, а именно, алгоритм быстрой сортировки (дополнительную информацию об этом алгоритме можно найти на странице http://en.wikipedia.org/wiki/Quicksort):

    public static void Quicksort<T> (T[] arr, int left, int right)
    		where T : IComparable <T>
    {
    if (right > left)
    		{
    		 int pivot = Partition ( arr, left, right );
    		 Quicksort ( arr, left, pivot - 1 );
    		 Quicksort ( arr, pivot + 1, right );
    		}
    }

    Из этого примера видно, что стартовать асинхронное выполнение вызовов Quicksort можно только после завершения шага Partition.

    Задача 9.

    Реализуйте параллельный вариант рекурсивного алгоритма Quicksort с помощью класса ThreadPool. Оцените количество порождаемых потоков при сортировке массивов большой длины. Придумайте способ увеличения эффективности данного алгоритма и реализуйте его.

    Противоположная проблема, по сравнению с алгоритмом Quicksort, имеет место в рекурсивном варианте программы сортировки слиянием:

    public static void Mergesort<T> (T[] arr, int left, int right)
    		where T : IComparable <T>
    {
    if (right > left)
    		{
    		 int mid =  ( right + left ) / 2;
    		 Mergesort ( arr, left, mid );
    		 Mergesort ( arr, mid + 1, right );
    		 Merge ( arr, left, mid + 1, right );
    		}
    }

    В этом варианте, мы получили алгоритм, аналогичный рассматривавшимся в семинаре к лекции 4. А именно - поток, выполняющий функцию Mergesort, блокируется перед выполнением шага Merge до тех пор, пока не будут (асинхронно) выполнены рекурсивные вызовы Mergesort. Другими словами, такой вид обработки очень трудно эффективно реализовать средствами класса ThreadPool.

    Вернуться к учебному плану