В предыдущих лекциях были рассмотрены различные классы библиотеки PFX и способы работы с ними. В данном разделе будет представлен ряд примеров программирования с использованием механизма задач и некоторых конструкций из библиотеки PFX, связанных с ними.
Библиотека PFX, а именно, ее составная часть Task Task и Future<T>, и суть его заключается в том, что с каждой задачей можно связать некоторую другую задачу - продолжение первой задачи. Эта вторая задача будет выполнена после завершения работы первой - основной, задачи.
Библиотека PFX версии June 2008 Future<T> и класса CountDownEvent:
static Task ContinueWhenAll(
Action<Task[]> continuation, params Task[] tasks)
{
var starter = Future<bool>.Create();
var task = starter.ContinueWith(o => continuation(tasks));
CountdownEvent ce = new CountdownEvent(tasks.Length);
Action<Task> whenComplete = delegate {
if (ce.Decrement()) starter.Value = true;
};
foreach (var t in tasks)
t.ContinueWith(whenComplete, TaskContinuationKind.OnAny,
TaskCreationOptions.None, true);
return task;
}
Суть решения состоит в использовании объекта starter класса Future<T>, который становится готовым (т.е., его свойство Value приобретает значение), когда оканчивают работу все задачи заданного множества задач tasks. Отслеживание окончания работы множества задач происходит с помощью счетчика ce класса CountdownEvent. Уменьшение счетчика на единицу производится каждой задачей из исходного множества задач, а точнее, с помощью задачи-продолжения whenComplete, которая зарегистрирована как задача-продолжение для каждой задачи исходного множества.
Вызывая метод ContinueWhenAll, программист в качестве одного из его параметров указывает список задач, по завершении которых необходимо запустить continuation. В механизме продолжений, реализованном в ContinueWhenAll, в отличие от стандартного механизма, передается целый массив задач, работу которых этот
Ниже приведен пример использования метода ContinueWhenAll:
Task t1 = ..., t2 = ...;
Task t3 = ContinueWhenAll(
delegate { Console.WriteLine("t1 and t2 finished"); }, t1, t2);
Реализация метода ContinueWhenAny, который позволяет запустить исполнение продолжения, когда хотя бы одна задача из множества задач завершилась, похожа на реализацию метода ContinueWhenAll и показана ниже:
static Task ContinueWhenAny(
Action<Task> continuation, params Task[] tasks)
{
WriteOnce<Task> theCompletedTask = new WriteOnce<Task>();
var starter = Future<bool>.Create();
var task = starter.ContinueWith(o =>
continuation(theCompletedTask.Value));
Action<Task> whenComplete = t => {
if (theCompletedTask.TrySetValue(t)) starter.Value = true;
};
foreach (var t in tasks)
t.ContinueWith(whenComplete, TaskContinuationKind.OnAny,
TaskCreationOptions.None, true);
return task;
}
Отметим особенности реализации метода ContinueWhenAll. Для того чтобы отследить момент завершения исполнения какой-либо задачи из множества задач tasks используется переменная с однократным присваиванием WriteOnce< >. Завершившаяся задача пробует присвоить этой переменной ссылку на себя и, если данная задача завершила свое исполнение первой, то метод TrySetValue вернет True и будет запущено исполнение продолжения, иначе метод TrySetValue вернет False, что будет означать, что данная задача завершилась не первой. Остальной код метода ContinueWhenAny аналогичен реализации метода ContinueWhenAll.
Иногда, в некоторых приложениях требуется выполнить друг за другом последовательность задач одновременно (асинхронно) с основным потоком. Другими словами, мы хотели бы создать некоторую перечислимую коллекцию (массив или список) задач и отдать их на выполнение другому потоку. Такой механизм можно реализовать, в частности, с помощью конструкций продолжение ContinueWith:
static void RunAsync(IEnumerable<Task> iterator)
{
var enumerator = iterator.GetEnumerator();
Action a = null;
a = delegate
{
if (enumerator.MoveNext())
{
enumerator.Current.ContinueWith(delegate { a(); });
}
else enumerator.Dispose();
};
a();
}
Метод RunAsync работает следующим образом: методу RunAsync в качестве аргумента передается перечисляемая коллекция задач; затем метод получает
В лекции 5 было описано несколько способов ожидания завершения исполнения одной или нескольких задач с помощью методов Task. и Task.WaitAll.
Предположим, что нам необходимо параллельно обработать некоторую совокупность данных, запуская для обработки одного элемента совокупности отдельную задачу, и затем дождаться завершения работы всех задач. В этом случае, соответствующий фрагмент кода может выглядеть так:
IEnumerable<Data> data = ...;
List<Task> tasks = new List<Task>();
foreach(var item in data) tasks.Add(Task.Create(delegate { Process(item); });
Task.WaitAll(tasks.ToArray());
Основным недостатком этого решения является сохранение ссылок на все порожденные задачи с вписке tasks. При большом размере исходной совокупности данных, и, следовательно, большом количестве запущенных задач, сборщик мусора не сможет освободить ресурсы, связанные с уже завершившимися задачами, поскольку ссылки на них собраны в массив tasks, который используется в операторе ожидания WaitAll в качестве аргумента.
На самом деле, в этой ситуации достаточно хранить только значение счетчика запущенных задач. В этом случае, условие "дождаться завершения всех созданных задач" будет эквивалентно условию обнуления этого счетчика. Именно эта идея реализована в классе TaskWaiter, код которого приведен ниже:
public class TaskWaiter
{
public CountdownEvent _ce = new CountdownEvent(1);
private bool _doneAdding = false;
public void Add(Task t)
{
if (t == null) throw new ArgumentNullException("t");
_ce.Increment();
t.ContinueWith(ct => _ce.Decrement());
}
public void Wait()
{
if (!_doneAdding) { _doneAdding = true; _ce.Decrement(); }
_ce.Wait();
}
}
При вызове метода TaskWaiter.Add происходит увеличение счетчика запущенных задач и установка уменьшения счетчика при завершении добавляемой задачи. Первое реализовано на базе потокобезопасного счетчика CountdownEvent, второе - на базе механизма продолжения задач (см. 9.1 и 9.2).
При вызове метода происходит блокирование вызывающего потока. Поток будет заблокирован до тех пор, пока счетчик _ce не примет нулевого значения или, что эквивалентно, до тех пор, пока все запущенные задачи не завершатся.
Пример использования данного класса приведен ниже:
IEnumerable<Data> data = ...;
TaskWaiter tasks = new TaskWaiter();
foreach(var item in data) tasks.Add(Task.Create(delegate { Process(item); });
tasks.Wait();
Обратите внимание, что ссылки на запущенные задачи более не сохраняются, что позволяет упростить код и повысить эффективность работы сборщика мусора.
Аналогично другим наборам инструментов, библиотека PFX доставляет программисту возможность на основе базовых конструкций, включенных в PFX, строить собственные специальные конструкции (шаблоны) для решения тех или иных задач. Ниже будут показаны способы построения одной из таких конструкций - ParallelWhileNotEmpty. С помощью этой конструкции можно обрабатывать элементы данных из некоторого множества, причем
Например, используя эту конструкцию можно запрограммировать обработку дерева, где при обработке отдельного узла в множество необработанных узлов добавляются узлы-потомки данного узла:
ParallelWhileNotEmpty(treeRoots, (node, adder) =>
{
foreach(var child in node.Children) adder(child);
Process(node);
});
Естественно, что существует несколько способов реализации такого шаблона. Рассмотрим некоторые из них.
Во-первых, одно из решений может быть построено на основе отношения "родитель-TaskCreationOptions. ). (Вспомнить эти механизмы можно еще раз обратившись к лекции 5 и разделу 3 данной лекции):
public static void ParallelWhileNotEmpty<T>(
IEnumerable<T> initialValues, Action<T, Action<T>> body)
{
Action<T> addMethod = null;
addMethod = v => Task.Create(delegate { body(v, addMethod); });
Parallel.ForEach(initialValues, addMethod);
}
Ключевым элементом данной реализации является addMethod, который будет создавать (body. Аргументами функции body является сам элемент, который нужно обработать (параметр v ), а также v ). Этим addMethod, т.е., мы воспользовались рекурсивным вызовом , передавая ему в качестве аргументов множество элементов, которое нужно обработать, и соответствующий не завершит свою работу пока все задачи (и родительская, и дочерние), созданные при вызове addMethod, не будут завершены. Данное решение имеет тот существенный недостаток, что любая задача, у которой существуют потомки, не будет завершена и убрана сборщиком мусора до тех пор, пока не завершат работу ее потомки. Это может привести к большому перерасходу памяти, особенно при работе со структурами с большой степенью вложенности. Устранить этот недостаток можно, воспользовавшись классом TaskWaiter из раздела 3 данной лекции:
public static void ParallelWhileNotEmpty<T>(
IEnumerable<T> initialValues, Action<T, Action<T>> body)
{
TaskWaiter tasks = new TaskWaiter();
Action<T> addMethod = null;
addMethod = v =>
tasks.Add(Task.Create(delegate { body(v, addMethod); },
TaskCreationOptions.Detached));
Parallel.ForEach(initialValues, addMethod);
tasks.Wait();
}
Из реализации класса TaskWaiter видно (см. раздел 3 данной лекции), что он позволяет разорвать связь "родитель-
Еще одна реализация шаблона ParallelWhileNotEmpty основана на использовании двух списков: в одном из них хранятся элементы, которые обрабатываются на текущей стадии, а в другом накапливаются элементы, которые будут обработаны на следующем шаге. В
public static void ParallelWhileNotEmpty<T>(
IEnumerable<T> initialValues, Action<T, Action<T>> body)
{
var lists = new [] {
new ConcurrentStack<T>(initialValues), new ConcurrentStack<T>() };
for(int i=0; ; i++)
{
int fromIndex = i % 2;
var from = lists[fromIndex];
var to = lists[fromIndex ^ 1];
if (from.IsEmpty) break;
Action<T> addMethod = v => to.Push(v);
Parallel.ForEach(from.ToArray(), v => body(v, addMethod));
from.Clear();
}
}
Одной из целей при проектировании библиотеки для .NET Framework было облегчить реализацию рекурсивных параллельных операций и обеспечить для них максимально возможную эффективность.
В частности, используя средства PLINQ, обход и обработка вершин дерева записывается в виде нескольких строк кода. Используя , который был реализован в примере семинара 8, метод Process запишется следующим образом:
public static void Process<T> (Tree<T> tree, Action<T> action)
{
if (tree == null) return;
GetNodes(tree).AsParallel().ForAll(action);
}
Внутренние действия, производимые библиотекой PFX при реализации данного PLINQ-фрагмента, очень похожи на действия из примера семинара, в котором создаются несколько потоков, которые выбирают вершины для обработки с помощью конструкции перечисления, используя механизм запирания ( lock ). Однако, в PLINQ такой подход реализован более эффективно путем увеличения количества вершин, извлекаемых из перечисления за один раз, что минимизирует использование конструкции lock (см. Задачу 7).
На самом деле, как и в ранее приведенных реализациях, в осуществляется последовательный проход по дереву с запуском действий обработки ( actions ) в асинхронном режиме. Чтобы реализовать параллельный проход по дереву, когда различные Task ):
public static void Process<T> (Tree<T> tree, Action<T> action)
{
if (tree == null) return;
Parallel.Invoke(
() => action(tree.Data),
() => Process(tree.Left, action),
() => Process(tree.Right, action));
}
Оператор потенциально запускает параллельное исполнение всех трех операторов из данного примера, заканчивая свою работу только по завершении их всех. Отличие от аналогичного примера, использующего класс ThreadPool, состоит в том, что здесь не может наступить состояние дедлока ввиду нехватки потоков для исполнения, поскольку реализация
Если необходимо более гибкое управление параллельным исполнением отдельных фрагментов кода, то они могут быть оформлены в виде задач:
public static void Process<T> (Tree<T> tree, Action<T> action)
{
if (tree == null) return;
var t1 = Task.Create(
delegate { Process(tree.Left, action); });
var t2 = Task.Create(
delegate { Process(tree.Right, action); });
action(tree.Data);
Task.WaitAll(new Task[] { t1, t2 });
}
Аналогичные методы могут быть применены к реализации алгоритма быстрой сортировки:
public static void Quicksort<T> (T[] arr, int left, int right)
where T : IComparable<T>
{
if (right > left)
{
int pivot = Partition(arr, left, right);
Parallel.Invoke(
() => Quicksort(arr, left, pivot - 1),
() => Quicksort(arr, pivot + 1, right));
}
}
а также к реализации алгоритма
public static void Mergesort<T> (
T[] arr, int left, int right) where T : IComparable<T>
{
if (right > left)
{
int mid = (right + left) / 2;
Parallel.Invoke(
() => Mergesort(arr, left, mid),
() => Mergesort(arr, mid + 1, right));
Merge(arr, left, mid + 1, right);
}
}
Хотя параллельные версии этих алгоритмов мы получили очень просто, однако, этот код, по-прежнему, имеет проблемы. Реализация конструкции ) для отдельных операций и ожидании завершения работы этих задач. Если вычислительная сложность отдельных операций мала, то накладные расходы на создание задач и ожидание завершения их работы будут велики по сравнению с основными вычислительными операциями. Это может приводить к тому, что параллельная реализация будет работать медленнее, чем последовательная.
Для того, чтобы решить эту проблему, необходимо немного усложнить код, применив широко известный механизм на базе использования порога ( ). Идея применения порога состоит в том, что мы вводим Process может выглядеть следующим образом:
private static void Process<T> (
Tree<T> tree, Action<T> action, int depth)
{
if (tree == null) return;
if (depth > 5)
{
action(tree.Data);
Process(tree.Left, action, depth + 1);
Process(tree.Right, action, depth + 1);
}
else
{
Parallel.Invoke(
() => action(tree.Data),
() => Process(tree.Left, action, depth + 1),
() => Process(tree.Right, action, depth + 1));
}
}
В присутствуют одновременно последовательная и параллельная реализации
Задача 10.
Реализуйте алгоритмы и Mergesort с помощью конструкции , используя механизм порогов.
Для получения официальных документов о завершении программы дополнительного профессионального образования (удостоверения о повышении квалификации, дипломов о профессиональной переподготовке и MBA) необходимо предоставить:
Внимание! Вы можете не заказывать доставку бумажной версии официального документы, а скачать его в электронном виде и распечатать самостоятельно. Информация о выданном документе в течение 1 месяца загружается в Федеральную информационную систему «Федеральный реестр сведений о документах об образовании и (или) о квалификации, документах об обучении» - ФИС ФРДО.
Доступ на новый сайт осуществляется с использованием адреса электронной почты, который был указан вами при регистрации на "старом". Мы постарались перенести все ваши данные с прежнего ресурса, однако не исключена вероятность потери части информации.
При возникновении проблемы со входом, воспользуйтесь функцией сброса пароля
Если вы обнаружите несоответствия, пожалуйста, сообщите нам.