Spec-Zone.ru › D

std.parallelism

std.parallelism реализует высокоуровневые примитивы для SMP-параллелизма. Они включают в себя параллельный foreach, параллельное reduce, параллельное eager map, конвейеризацию и будущее/обещание параллелизма. std.parallelism рекомендуется, когда одно и то же действие должно выполняться параллельно на разных данных или когда функция должна выполняться в фоновом потоке, а ее результат возвращается в определенный основной поток. Для взаимодействия между произвольными потоками см. std.concurrency.

std.parallelism основано на концепции Task. Task — это объект, который представляет собой основную единицу работы в этой библиотеке и может выполняться параллельно с любой другой Task. Использование Task напрямую позволяет программировать с парадигмой будущего/обещания. Все остальные поддерживаемые парадигмы параллелизма (параллельный foreach, map, reduce, конвейеризация) представляют дополнительный уровень абстракции над Task. Они автоматически создают один или несколько Task объектов или тесно связанные типы, которые концептуально идентичны, но не являются частью публичного API.

После создания Task может быть выполнено в новом потоке или отправлено в TaskPool для выполнения. TaskPool инкапсулирует очередь задач и ее рабочие потоки. Его цель — эффективно отобразить большое количество Task на меньшее количество потоков. Очередь задач — это очередь FIFO объектов Task, которые были отправлены в TaskPool и ожидают выполнения. Рабочий поток — это поток, который связан ровно с одной очередью задач. Он выполняет Task в начале своей очереди, когда в очереди есть работа, или приостанавливается, когда работы нет. Каждая очередь задач связана с нулем или более рабочими потоками. Если результат Task необходим до начала выполнения рабочим потоком, Task может быть удален из очереди задач и выполнен немедленно в потоке, где нужен результат.

Предупреждение
Если не отмечено как @trusted или @safe, артефакты в этом модуле допускают неявное совместное использование данных между потоками и не могут гарантировать, что клиентский код свободен от низкоуровневых гонок данных.
Источник
std/parallelism.d
Автор
David Simcha
Лицензия:
Лицензия Boost 1.0
struct Task(alias fun, Args...);

Task представляет собой основную единицу работы. Task можно выполнить параллельно с любым другим Task. Использование этой структуры напрямую позволяет использовать параллелизм задач/обещаний. В этой парадигме функция (или делегат или другой вызываемый объект) выполняется в потоке, отличном от того, из которого она была вызвана. Вызывающий поток не блокируется во время выполнения функции. Вызов workForce, yieldForce, или spinForce используется для того, чтобы гарантировать, что Task завершила выполнение, и получить возвращаемое значение, если таковое имеется. Эти функции и done также действуют как полные барьеры памяти, что означает, что все записи в память, сделанные в потоке, в котором выполнялся Task, гарантированно будут видны в вызывающем потоке после возвращения одной из этих функций.

Функции std.parallelism.task и std.parallelism.scopedTask могут быть использованы для создания экземпляра этой структуры. См. task для примеров использования.

Результаты функций возвращаются из yieldForce, spinForce и workForce по ссылке. Если fun возвращает по ссылке, эта ссылка будет указывать на возвращенную ссылку fun. В противном случае она будет указывать на поле в этой структуре.

Копирование этой структуры отключено, так как это не имеет смысла. Если вы хотите передать эту структуру, вы должны передавать ее по ссылке или указателю.

Ошибки:
Изменения в аргументах ref и out не распространяются на место вызова, а только на args в этой структуре.
alias args = _args[1 .. __dollar];

Аргументы, с которыми была вызвана функция. Изменения в аргументах out и ref будут видны здесь.

alias ReturnType = typeof(fun(_args));

Тип возвращаемого значения функции, вызванной этим Task. Это может быть void.

@property ref @trusted ReturnType spinForce();

Если Task еще не запущена, выполните ее в текущем потоке. Если она завершена, верните ее возвращаемое значение, если таковое имеется. Если она выполняется, ждите, пока она не завершится, а затем верните возвращаемое значение. Если возникло исключение, перебросьте это исключение.

Эта функция должна использоваться, когда вы ожидаете, что результат Task будет доступен в течение времени, меньшего, чем переключение контекста ОС.

@property ref @trusted ReturnType yieldForce();

Если Task еще не запущена, выполните ее в текущем потоке. Если она завершена, верните ее возвращаемое значение, если таковое имеется. Если она выполняется, дождитесь события. Если возникло исключение, перебросьте это исключение.

Эта функция должна использоваться для дорогостоящих функций, так как ожидание события вносит задержку, но избегает ненужных циклов CPU.

@property ref @trusted ReturnType workForce();

Если эта Task еще не запущена, выполните её в текущем потоке. Если она завершена, верните её результат. Если она выполняется, выполните все остальные Task из экземпляра TaskPool, которому была передана эта Task, пока она не завершится. Если возникло исключение, перебросьте это исключение. Если другие задачи недоступны или эта Task была выполнена с помощью executeInNewThread, дождитесь события.

@property @trusted bool done();

Возвращает true , если Task завершила выполнение.

Исключения:
Перебрасывает любое исключение, сгенерированное во время выполнения Task.
@trusted void executeInNewThread();

@trusted void executeInNewThread(int priority);

Создайте новый поток для выполнения этой Task, выполните его в новом потоке, а затем завершите поток. Это может быть использовано для параллелизма задач/обещаний. Задаче может быть задан явный приоритет. Если он предоставлен, его значение передается core.thread.Thread.priority. См. std.parallelism.task для примера использования.

auto task(alias fun, Args...)(Args args);

Создаёт Task в куче GC, который вызывает алиас. Это может быть выполнено с помощью Task.executeInNewThread или путём передачи в std.parallelism.TaskPool. Глобально доступный экземпляр TaskPool предоставляется std.parallelism.taskPool.

Возвращает:
Указатель на Task.
Пример
// Read two files into memory at the same time.
import std.file;

void main()
{
    // Create and execute a Task for reading
    // foo.txt.
    auto file1Task = task!read("foo.txt");
    file1Task.executeInNewThread();

    // Read bar.txt in parallel.
    auto file2Data = read("bar.txt");

    // Get the results of reading foo.txt.
    auto file1Data = file1Task.yieldForce;
}
// Sorts an array using a parallel quick sort algorithm.
// The first partition is done serially.  Both recursion
// branches are then executed in parallel.
//
// Timings for sorting an array of 1,000,000 doubles on
// an Athlon 64 X2 dual core machine:
//
// This implementation:               176 milliseconds.
// Equivalent serial implementation:  280 milliseconds
void parallelSort(T)(T[] data)
{
    // Sort small subarrays serially.
    if (data.length < 100)
    {
         std.algorithm.sort(data);
         return;
    }

    // Partition the array.
    swap(data[$ / 2], data[$ - 1]);
    auto pivot = data[$ - 1];
    bool lessThanPivot(T elem) { return elem < pivot; }

    auto greaterEqual = partition!lessThanPivot(data[0..$ - 1]);
    swap(data[$ - greaterEqual.length - 1], data[$ - 1]);

    auto less = data[0..$ - greaterEqual.length - 1];
    greaterEqual = data[$ - greaterEqual.length..$];

    // Execute both recursion branches in parallel.
    auto recurseTask = task!parallelSort(greaterEqual);
    taskPool.put(recurseTask);
    parallelSort(less);
    recurseTask.yieldForce;
}
auto task(F, Args...)(F delegateOrFp, Args args)
Constraints: if (is(typeof(delegateOrFp(args))) && !isSafeTask!F);

Создаёт Task в куче GC, который вызывает указатель на функцию, делегат или класс/структуру с перегруженным оператором вызова.

Пример
// Read two files in at the same time again,
// but this time use a function pointer instead
// of an alias to represent std.file.read.
import std.file;

void main()
{
    // Create and execute a Task for reading
    // foo.txt.
    auto file1Task = task(&read!string, "foo.txt", size_t.max);
    file1Task.executeInNewThread();

    // Read bar.txt in parallel.
    auto file2Data = read("bar.txt");

    // Get the results of reading foo.txt.
    auto file1Data = file1Task.yieldForce;
}
Примечания
Эта функция принимает делегат без области видимости, что означает, что её можно использовать с замыканиями. Если вы не можете выделить замыкание из-за объектов на стеке, у которых есть уничтожение области видимости, см. scopedTask, которая принимает делегат области видимости.
@trusted auto task(F, Args...)(F fun, Args args)
Constraints: if (is(typeof(fun(args))) && isSafeTask!F);

Версия task , пригодная для использования из @safe кода. Механизм использования идентичен случаю без @safe, но безопасность вводит некоторые ограничения:

  1. fun должен быть @safe или @trusted.
  2. F не должен иметь никаких алиасов без общего использования, как определено в std.traits.hasUnsharedAliasing. Это означает, что это не может быть делегат без общего использования или не-общий класс или структура с перегруженным оператором opCall. Это также исключает прием параметров алиасов шаблонов.
  3. Args не должен иметь алиасов без общего использования.
  4. fun не должен возвращать по ссылке.
  5. Тип возвращаемого значения не должен иметь алиасов без общего использования, если только fun не является pure или Task выполняется с помощью executeInNewThread вместо использования TaskPool.

auto scopedTask(alias fun, Args...)(Args args);

auto scopedTask(F, Args...)(scope F delegateOrFp, Args args)
Constraints: if (is(typeof(delegateOrFp(args))) && !isSafeTask!F);

@trusted auto scopedTask(F, Args...)(F fun, Args args)
Constraints: if (is(typeof(fun(args))) && isSafeTask!F);

Эти функции позволяют создавать объекты Task на стеке, а не в куче GC. Жизненный цикл Task созданного с помощью scopedTask не может превышать жизненного цикла области видимости, в которой он был создан.

scopedTask может быть предпочтительнее task:

  1. Когда создаётся Task , вызывающий делегат, и замыкание не может быть выделено из-за объектов на стеке, у которых есть уничтожение области видимости. Перегрузка делегата scopedTask принимает делегат scope.
  2. В качестве микрооптимизации, для избежания выделения памяти в куче, связанного с task или с созданием замыкания.
Использование в остальном идентично task.

Примечания
Объекты Task , созданные с помощью scopedTask , автоматически вызовут Task.yieldForce в своём деструкторе, если это необходимо, чтобы гарантировать завершение Task до того, как будет уничтожена кадр стека, на котором они находятся.
alias totalCPUs = __lazilyInitializedConstant!(immutable(uint), 4294967295u, totalCPUsImpl).__lazilyInitializedConstant;

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

class TaskPool;

Этот класс encapsulates очередь задач и набор потоков-рабочих. Его назначение — эффективно отобразить большое количество Task на меньшее количество потоков. Очередь задач — это очередь FIFO объектов Task , которые были переданы TaskPool и ожидают выполнения. Поток-рабочий — это поток, который выполняет Task в начале очереди, когда она доступна, и останавливается, когда очередь пуста.

Этот класс обычно используется через глобальный экземпляр, доступный через свойство std.parallelism.taskPool. Иногда бывает полезно явно создать TaskPool:

  1. Когда вы хотите экземпляры TaskPool с несколькими приоритетами, например, очередь с низким приоритетом и очередь с высоким приоритетом.
  2. Когда потоки в глобальном пуле задач ожидают синхронизирующего примитива (например, мьютекса), и вы хотите распараллелить код, который должен быть выполнен, прежде чем эти потоки могут быть возобновлены.

Примечание
Потоки-рабочие в этом пуле не остановятся, пока не будет вызвано stop или finish, даже если основной поток уже завершен. Это может привести к программам, которые никогда не завершаются. Если вы не хотите этого поведения, вы можете установить isDaemon в true.
@trusted this();

Конструктор по умолчанию, который инициализирует TaskPool с totalCPUs - 1 рабочими потоками. Минус 1 включён, потому что основной поток также будет доступен для выполнения работы.

Примечание
На одноядерных машинах примитивы, предоставляемые TaskPool , работают прозрачно в однопоточном режиме.
@trusted this(size_t nWorkers);

Позволяет настроить пользовательское количество рабочих потоков.

ParallelForeach!R parallel(R)(R range, size_t workUnitSize);

ParallelForeach!R parallel(R)(R range);

Реализует цикл foreach в параллельном режиме над диапазоном. Это работает, неявно создавая и отправляя по одному Task в TaskPool для каждого рабочего потока. Рабочая единица — это набор последовательных элементов range , которые должны быть обработаны рабочим потоком между коммуникациями с другими потоками. Количество элементов, обрабатываемых за одну рабочую единицу, регулируется параметром workUnitSize. Меньшие рабочие единицы обеспечивают лучшую балансировку нагрузки, но более крупные рабочие единицы избегают накладных расходов на частые коммуникации с другими потоками для получения следующей рабочей единицы. Крупные рабочие единицы также избегают ложного совместного использования в тех случаях, когда диапазон изменяется. Чем меньше время занимает одна итерация цикла, тем больше должно быть значение workUnitSize. Для очень ресурсоёмких циклов тела workUnitSize должно быть равно 1. Также доступен перегруз, который выбирает размер рабочей единицы по умолчанию.

Пример
// Find the logarithm of every number from 1 to
// 10_000_000 in parallel.
auto logs = new double[10_000_000];

// Parallel foreach works with or without an index
// variable.  It can be iterate by ref if range.front
// returns by ref.

// Iterate over logs using work units of size 100.
foreach (i, ref elem; taskPool.parallel(logs, 100))
{
    elem = log(i + 1.0);
}

// Same thing, but use the default work unit size.
//
// Timings on an Athlon 64 X2 dual core machine:
//
// Parallel foreach:  388 milliseconds
// Regular foreach:   619 milliseconds
foreach (i, ref elem; taskPool.parallel(logs))
{
    elem = log(i + 1.0);
}
Примечания
Потребление памяти в этой реализации гарантированно является постоянным в range.length.
Прерывание цикла parallel foreach с помощью break, labeled break, labeled continue, return или goto вызывает ParallelForeachError. В случае диапазонов без произвольного доступа parallel foreach лениво буферизует в массив размером workUnitSize перед выполнением параллельной части цикла. Исключением является тот случай, когда parallel foreach выполняется над диапазоном, возвращаемым asyncBuf или map, копирование опускается, и буферы просто меняются местами. В этом случае workUnitSize игнорируется, а размер рабочей единицы устанавливается в размер буфера range. Гарантируется, что при выходе из цикла будет выполнен барьер памяти, чтобы результаты, полученные всеми потоками, были видны в вызывающем потоке. **Обработка исключений**: Если хотя бы одно исключение возникает внутри цикла parallel foreach, отправка дополнительных Task объектов прекращается как можно скорее, непредсказуемым образом. Все выполняемые или ожидающие рабочие единицы допускаются до завершения. Затем все исключения, сгенерированные любой рабочей единицей, цепляются с использованием Throwable.next и перебрасываются повторно. Порядок цепочки исключений непредсказуем.
template amap(functions...)
auto amap(Args...)(Args args)
Constraints: if (isRandomAccessRange!(Args[0]));

Жёсткая параллельная функция map. Жесткость этой функции означает, что она имеет меньшие накладные расходы, чем лениво вычисляемая TaskPool.map, и её следует предпочитать, если допустимы требования памяти к жёсткости. functions — это функции, подлежащие вычислению, передаваемые как параметры шаблонов-псевдонимов в стиле, похожем на std.algorithm.iteration.map. Первый аргумент должен быть диапазоном произвольного доступа. По соображениям производительности, функция amap будет предполагать, что элементы диапазона ещё не инициализированы. Элементы будут перезаписаны без вызова деструктора и выполнения присваивания. Таким образом, диапазон не должен содержать значимых данных: либо неинициализированные объекты, либо объекты в состоянии .init.

auto numbers = iota(100_000_000.0);

// Find the square roots of numbers.
//
// Timings on an Athlon 64 X2 dual core machine:
//
// Parallel eager map:                   0.802 s
// Equivalent serial implementation:     1.768 s
auto squareRoots = taskPool.amap!sqrt(numbers);


Непосредственно после аргумента диапазона может быть указан необязательный аргумент размера рабочей единицы. Рабочие единицы, используемые amap, идентичны тем, которые определены для parallel foreach. Если размер рабочей единицы не указан, используется размер рабочей единицы по умолчанию.

// Same thing, but make work unit size 100.
auto squareRoots = taskPool.amap!sqrt(numbers, 100);


Диапазон вывода для возвращения результатов может быть указан в качестве последнего аргумента. Если он не указан, массив соответствующего типа будет выделен в куче с управляемой сборкой мусора. Если он указан, он должен быть диапазоном произвольного доступа с присваиваемыми элементами, должен иметь семантику ссылок относительно присваивания к его элементам и иметь такую же длину, что и входной диапазон. Запись в смежные элементы из разных потоков должна быть безопасной.

// Same thing, but explicitly allocate an array
// to return the results in.  The element type
// of the array may be either the exact type
// returned by functions or an implicit conversion
// target.
auto squareRoots = new float[numbers.length];
taskPool.amap!sqrt(numbers, squareRoots);

// Multiple functions, explicit output range, and
// explicit work unit size.
auto results = new Tuple!(float, real)[numbers.length];
taskPool.amap!(sqrt, log)(numbers, 100, results);

Примечание
Гарантируется, что после записи всех результатов, но перед возвратом, будет выполнен барьер памяти, чтобы результаты, полученные всеми потоками, были видны в вызывающем потоке.
Советы
Для выполнения операции отображения на месте укажите тот же диапазон для входного и выходного диапазона.
Для параллелизации копирования диапазона с дорогостоящими для вычисления элементами в массив передайте функцию тождества (функцию, которая просто возвращает переданный ей аргумент) в amap. **Обработка исключений**: Если хотя бы одно исключение возникает внутри функций map, отправка дополнительных Task объектов прекращается как можно скорее, непредсказуемым образом. Все выполняемые или ожидающие рабочие единицы допускаются до завершения. Затем все исключения, сгенерированные любой рабочей единицей, цепляются с использованием Throwable.next и перебрасываются повторно. Порядок цепочки исключений непредсказуем.
template map(functions...)
auto map(S)(S source, size_t bufSize = 100, size_t workUnitSize = size_t.max)
Constraints: if (isInputRange!S);

Полуленивая параллельная функция map, которую можно использовать для конвейеризации. Функции map вычисляются для первых bufSize элементов, сохраняются в буфере и делаются доступными для popFront. В то же время в фоновом режиме второй буфер той же размерности заполняется. Когда первый буфер исчерпан, он меняется местами со вторым буфером и заполняется, в то время как значения из того, что изначально было вторым буфером, считываются. Эта реализация позволяет записывать элементы в буфер без необходимости в атомарных операциях или синхронизации для каждой записи и позволяет эффективно вычислять функцию отображения в параллельном режиме.

map имеет большие накладные расходы, чем более простая процедура, используемая amap, но избегает необходимости держать все результаты в памяти одновременно и работает с диапазонами без произвольного доступа.

Параметры:
S source Входной диапазон, подлежащий отображению. Если source не имеет произвольного доступа, он будет лениво буферизован в массив размером bufSize перед вычислением функции map. (Исключение из этого правила см. в Примечаниях.)
size_t bufSize Размер буфера для хранения вычисленных элементов.
size_t workUnitSize Количество элементов, подлежащих вычислению за одну рабочую Task. Должно быть меньше или равно bufSize, и должно составлять дробь от bufSize таким образом, чтобы все рабочие потоки могли использоваться. Если используется значение по умолчанию size_t.max, workUnitSize будет установлено в значение по умолчанию для всего пула.
Возвращает:
Входной диапазон, представляющий результаты отображения. Этот диапазон имеет длину, если source имеет длину.
Примечания
Если диапазон, возвращаемый map или asyncBuf, используется в качестве входного для map, то в качестве оптимизации копирование из выходного буфера первого диапазона в входной буфер второго диапазона опускается, даже если диапазоны, возвращаемые map и asyncBuf, являются диапазонами без произвольного доступа. Это означает, что параметр bufSize , переданный в текущий вызов map, будет проигнорирован, и размер буфера будет равен размеру буфера source.
Пример
// Pipeline reading a file, converting each line
// to a number, taking the logarithms of the numbers,
// and performing the additions necessary to find
// the sum of the logarithms.

auto lineRange = File("numberList.txt").byLine();
auto dupedLines = std.algorithm.map!"a.idup"(lineRange);
auto nums = taskPool.map!(to!double)(dupedLines);
auto logs = taskPool.map!log10(nums);

double sum = 0;
foreach (elem; logs)
{
    sum += elem;
}
**Обработка исключений**: Любые исключения, возникающие при итерации над source или вычислении функции map, перебрасываются при вызове popFront или, если они возникают во время создания, просто допускаются до распространения вызывающей стороне. В случае исключений, возникших при вычислении функции map, исключения цепляются как в TaskPool.amap.
auto asyncBuf(S)(S source, size_t bufSize = 100)
Constraints: if (isInputRange!S);

При заданном source диапазоне, итерация по которому является ресурсоёмкой, возвращает диапазон, который асинхронно буферизует содержимое source в буфер размером bufSize элементов в рабочем потоке, делая ранее буферизованные элементы из второго буфера, также размером bufSize, доступными через интерфейс диапазона возвращаемого объекта. Возвращаемый диапазон имеет длину, если hasLength!S. asyncBuf полезен, например, при выполнении ресурсоёмких операций над элементами диапазонов, представляющих данные на диске или в сети.

Пример
import std.conv, std.stdio;

void main()
{
    // Fetch lines of a file in a background thread
    // while processing previously fetched lines,
    // dealing with byLine's buffer recycling by
    // eagerly duplicating every line.
    auto lines = File("foo.txt").byLine();
    auto duped = std.algorithm.map!"a.idup"(lines);

    // Fetch more lines in the background while we
    // process the lines already read into memory
    // into a matrix of doubles.
    double[][] matrix;
    auto asyncReader = taskPool.asyncBuf(duped);

    foreach (line; asyncReader)
    {
        auto ls = line.split("\t");
        matrix ~= to!(double[])(ls);
    }
}
**Обработка исключений**: Любые исключения, возникшие при итерации по source, перебрасываются при вызове popFront или, если они возникают во время создания, просто допускаются до распространения вызывающей стороне.
auto asyncBuf(C1, C2)(C1 next, C2 empty, size_t initialBufSize = 0, size_t nBuffers = 100)
Constraints: if (is(typeof(C2.init()) : bool) && (Parameters!C1.length == 1) && (Parameters!C2.length == 0) && isArray!(Parameters!C1[0]));

Принимая на вход вызываемый объект next, который записывает в предоставленный пользователем буфер, и второй вызываемый объект empty, который определяет, есть ли доступные данные для записи с помощью next, возвращает диапазон ввода, который асинхронно вызывает next с набором буферов размером nBuffers и делает результаты доступными в порядке их получения через интерфейс диапазона ввода возвращаемого объекта. Аналогично перегрузке диапазона ввода asyncBuf, первая половина буферов становится доступной через интерфейс диапазона, а вторая половина заполняется, и наоборот.

Параметры:
C1 next Вызываемый объект, принимающий единственный аргумент, который должен быть массивом с изменяемыми элементами. При вызове next записывает данные в массив, предоставленный вызывающей стороной.
C2 empty Вызываемый объект, не принимающий аргументов и возвращающий тип, неявно преобразуемый в bool. Это используется для обозначения того, что больше данных получить не удастся, вызвав next.
size_t initialBufSize Начальный размер каждого буфера. Если next принимает массив по ссылке, он может изменить размер буферов.
size_t nBuffers Количество буферов для циклического прохода при вызове next.
Пример
// Fetch lines of a file in a background
// thread while processing previously fetched
// lines, without duplicating any lines.
auto file = File("foo.txt");

void next(ref char[] buf)
{
    file.readln(buf);
}

// Fetch more lines in the background while we
// process the lines already read into memory
// into a matrix of doubles.
double[][] matrix;
auto asyncReader = taskPool.asyncBuf(&next, &file.eof);

foreach (line; asyncReader)
{
    auto ls = line.split("\t");
    matrix ~= to!(double[])(ls);
}
Обработка исключений: Любые исключения, возникшие во время итерации по range, повторно выбрасываются при вызове popFront.
Предупреждение
Использование диапазона, возвращаемого этой функцией, в параллельном цикле foreach не сработает, так как буферы могут перезаписываться, пока задача по их обработке находится в очереди. Это проверяется во время компиляции и приведет к ошибке статического утверждения.
template reduce(functions...)
auto reduce(Args...)(Args args);

Параллельное сокращение по диапазону с произвольным доступом. За исключением специально оговоренных случаев, использование аналогично std.algorithm.iteration.reduce. Также есть fold, которая делает то же самое, но с другим порядком параметров.

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

Поскольку сокращение выполняется параллельно, functions должна быть ассоциативной. Для простоты обозначений пусть # - инфиксный оператор, представляющий functions. Тогда (a # b) # c должно быть равно a # (b # c). Сложение чисел с плавающей точкой не является ассоциативным, хотя в точном арифметическом представлении сложение ассоциативно. Суммирование чисел с плавающей точкой с помощью этой функции может давать разные результаты, чем последовательное суммирование. Однако для многих практических целей сложение чисел с плавающей точкой можно считать ассоциативным.

Обратите внимание, что поскольку functions предполагаются ассоциативными, выполняются дополнительные оптимизации последовательной части алгоритма сокращения. Они используют уровень инструкций параллелизма современных процессоров, в дополнение к уровню потоков параллелизма, который использует остальная часть этого модуля. Это может привести к ускорениям, превышающим линейные, по сравнению с std.algorithm.iteration.reduce, особенно для тонкозернистых бенчмарков, таких как скалярные произведения.

Явное начальное значение может быть предоставлено в качестве первого аргумента. Если оно указано, оно используется в качестве начального значения для всех рабочих единиц и для окончательного сокращения результатов всех рабочих единиц. Поэтому, если это не тождественное значение для выполняемой операции, результаты могут отличаться от результатов, сгенерированных с помощью std.algorithm.iteration.reduce, или в зависимости от того, сколько рабочих единиц используется. Следующий аргумент должен быть диапазоном, подлежащим сокращению.

// Find the sum of squares of a range in parallel, using
// an explicit seed.
//
// Timings on an Athlon 64 X2 dual core machine:
//
// Parallel reduce:                     72 milliseconds
// Using std.algorithm.reduce instead:  181 milliseconds
auto nums = iota(10_000_000.0f);
auto sumSquares = taskPool.reduce!"a + b"(
    0.0, std.algorithm.map!"a * a"(nums)
);


Если явное начальное значение не указано, первый элемент каждой рабочей единицы используется в качестве начального значения. Для окончательного сокращения результат от первой рабочей единицы используется в качестве начального значения.
// Find the sum of a range in parallel, using the first
// element of each work unit as the seed.
auto sum = taskPool.reduce!"a + b"(nums);


Явный размер рабочей единицы может быть указан в качестве последнего аргумента. Указание слишком малого размера рабочей единицы фактически сериализует сокращение, так как окончательное сокращение результата каждой рабочей единицы будет доминировать во времени вычислений. Если TaskPool.size для этого экземпляра равно нулю, этот параметр игнорируется, и используется одна рабочая единица.
// Use a work unit size of 100.
auto sum2 = taskPool.reduce!"a + b"(nums, 100);

// Work unit size of 100 and explicit seed.
auto sum3 = taskPool.reduce!"a + b"(0.0, nums, 100);


Параллельное сокращение поддерживает несколько функций, например, std.algorithm.reduce.
// Find both the min and max of nums.
auto minMax = taskPool.reduce!(min, max)(nums);
assert(minMax[0] == reduce!min(nums));
assert(minMax[1] == reduce!max(nums));


Обработка исключений:

После завершения выполнения этой функции любые возникшие исключения объединяются с помощью Throwable.next и повторно выбрасываются. Порядок объединения не определен.

См. также:
fold функционально эквивалентна reduce, за исключением того, что параметр диапазона указывается первым, и нет необходимости использовать tuple для нескольких начальных значений.
template fold(functions...)
auto fold(Args...)(Args args);

Реализует одноимённую функцию (также известную как accumulate, compress, inject, или foldl), присутствующую в различных функциональных языках программирования.

fold функционально эквивалентна reduce, за исключением того, что параметр диапазона указывается первым, и нет необходимости использовать tuple для нескольких начальных значений.

Может быть один или несколько вызываемых объектов (аргумент functions).

Параметры:
Args args Только диапазон, по которому нужно выполнить свёртку; или диапазон и по одному начальному значению для каждой функции; или диапазон, по одному начальному значению для каждой функции и размер рабочей единицы
Возвращает:
Накопленный результат как одно значение для одной функции и как кортеж значений для нескольких функций
См. также:
Аналогично std.algorithm.iteration.fold, fold - это обёртка вокруг reduce.
Пример
static int adder(int a, int b)
{
    return a + b;
}
static int multiplier(int a, int b)
{
    return a * b;
}

// Just the range
auto x = taskPool.fold!adder([1, 2, 3, 4]);
assert(x == 10);

// The range and the seeds (0 and 1 below; also note multiple
// functions in this example)
auto y = taskPool.fold!(adder, multiplier)([1, 2, 3, 4], 0, 1);
assert(y[0] == 10);
assert(y[1] == 24);

// The range, the seed (0), and the work unit size (20)
auto z = taskPool.fold!adder([1, 2, 3, 4], 0, 20);
assert(z == 10);
const nothrow @property @safe size_t workerIndex();

Возвращает индекс текущего потока относительно этого TaskPool. Любой поток, не принадлежащий этому пулу, получит индекс 0. Рабочие потоки в этом пуле получают уникальные индексы от 1 до this.size.

Эта функция полезна для поддержания ресурсов, локальных для каждого рабочего потока.

Пример
// Execute a loop that computes the greatest common
// divisor of every number from 0 through 999 with
// 42 in parallel.  Write the results out to
// a set of files, one for each thread.  This allows
// results to be written out without any synchronization.

import std.conv, std.range, std.numeric, std.stdio;

void main()
{
    auto filesHandles = new File[taskPool.size + 1];
    scope(exit) {
        foreach (ref handle; fileHandles)
        {
            handle.close();
        }
    }

    foreach (i, ref handle; fileHandles)
    {
        handle = File("workerResults" ~ to!string(i) ~ ".txt");
    }

    foreach (num; parallel(iota(1_000)))
    {
        auto outHandle = fileHandles[taskPool.workerIndex];
        outHandle.writeln(num, '\t', gcd(num, 42));
    }
}
struct WorkerLocalStorage(T);

Структура для создания локального хранилища для рабочих потоков. Локальное хранилище для рабочих потоков — это хранилище, локальное для потока, существующее только для рабочих потоков в данном TaskPool плюс один поток вне пула. Оно выделяется в куче с сборкой мусора способом, предотвращающим ложное совместное использование, и не обязательно имеет глобальный охват в любом потоке. К нему можно получить доступ из любого рабочего потока в TaskPool , который его создал, и из одного потока вне этого TaskPool. Все потоки за пределами пула, создавшие данный экземпляр локального хранилища для рабочих потоков, разделяют один слот.

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

Основные случаи использования WorkerLocalStorageStorage:

  1. Выполнение параллельных сокращений с императивным, а не функциональным, стилем программирования. В этом случае WorkerLocalStorageStorage полезно рассматривать как локальное для каждого потока только для параллельной части алгоритма.
  2. Переиспользование временных буферов при итерациях параллельного цикла foreach.

Пример
// Calculate pi as in our synopsis example, but
// use an imperative instead of a functional style.
immutable n = 1_000_000_000;
immutable delta = 1.0L / n;

auto sums = taskPool.workerLocalStorage(0.0L);
foreach (i; parallel(iota(n)))
{
    immutable x = ( i - 0.5L ) * delta;
    immutable toAdd = delta / ( 1.0 + x * x );
    sums.get += toAdd;
}

// Add up the results from each worker thread.
real pi = 0;
foreach (threadResult; sums.toRange)
{
    pi += 4.0L * threadResult;
}
@property ref auto get(this Qualified)();

Получить текущий экземпляр потока. Возвращается по ссылке. Обратите внимание, что вызов get из любого потока вне TaskPool , который создал этот экземпляр, вернёт ту же ссылку, поэтому экземпляр локального хранилища рабочих потоков должен использоваться только одним потоком вне пула, который его создал. Если это правило нарушено, результатом будет неопределённое поведение.

Если включены проверки и toRange был вызван, то этот экземпляр WorkerLocalStorage больше не является локальным для рабочего потока, и при вызове этого метода возникнет ошибка утверждения. Это не проверяется при отключённых проверках по соображениям производительности.

@property void get(T val);

Присвоить значение текущему экземпляру потока. Эта функция имеет те же ограничения, что и её перегрузка.

@property WorkerLocalStorageRange!T toRange();

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

Вызов этой функции устанавливает флаг, указывающий, что эта структура больше не локальна для рабочих потоков, и попытка повторного использования метода get приведёт к ошибке утверждения, если проверки включены.

struct WorkerLocalStorageRange(T);

Базовые элементы диапазона для локального хранилища рабочих потоков. Цель этого — получить доступ к результатам, созданным каждым рабочим потоком, из одного потока, как только вы больше не используете локальное хранилище рабочих потоков из нескольких потоков. Не используйте эту структуру в параллельной части алгоритма.

Правильный способ создания этого объекта — вызов WorkerLocalStorage.toRange. После создания этот объект ведёт себя как конечный диапазон с произвольным доступом, со свойством присваивания элементов-ссылок и длиной, равной количеству рабочих потоков в TaskPool , который его создал, плюс 1.

WorkerLocalStorage!T workerLocalStorage(T)(lazy T initialVal = T.init);

Создаёт экземпляр локального хранилища для потока, инициализированный заданным значением. Значение является lazy, что позволяет, например, легко создавать по одному экземпляру класса для каждого потока. Пример использования см. в структуре WorkerLocalStorage.

@trusted void stop();

Посылает сигнал всем рабочим потокам завершиться как только они закончат текущую Task, или немедленно, если они не выполняют Task. Task в очереди не будут выполнены, пока не будет вызвано Task.workForce, Task.yieldForce или Task.spinForce, что приведёт к их выполнению.

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

@trusted void finish(bool blocking = false);

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

Если аргумент blocking равен true, то дожидается завершения всех рабочих потоков перед возвратом. Этот параметр может использоваться в приложениях, где результаты задач никогда не потребляются — например, когда TaskPool используется как примитивный планировщик задач, которые взаимодействуют другими способами, кроме возвращаемых значений.

Предупреждение
Вызов этой функции с blocking = true из рабочего потока, являющегося членом того же TaskPool, на котором вызывается finish, приведёт к тупику.
const pure nothrow @property @safe size_t size();

Возвращает количество рабочих потоков в пуле.

void put(alias fun, Args...)(ref Task!(fun, Args) task)
Constraints: if (!isSafeReturn!(typeof(task)));

void put(alias fun, Args...)(Task!(fun, Args)* task)
Constraints: if (!isSafeReturn!(typeof(*task)));

Помещает объект Task в конец очереди задач. Объект Task может быть передан по указателю или ссылке.

Пример
import std.file;

// Create a task.
auto t = task!read("foo.txt");

// Add it to the queue to be executed.
taskPool.put(t);
Примечания
@trusted перегрузки этой функции вызываются для Task если std.traits.hasUnsharedAliasing ложно для типа возвращаемого значения Task или функция, которую выполняет Task, является pure. Объекты Task, которые удовлетворяют всем остальным требованиям, указанным в @trusted перегрузках task и scopedTask, могут быть созданы и выполнены из кода @safe с помощью Task.executeInNewThread, но не с помощью TaskPool.
Хотя эта функция принимает адреса переменных, которые могут находиться в стеке, некоторые перегрузки помечены как @trusted. Task включает деструктор, который дожидается завершения задачи перед уничтожением фрейма стека, на котором он выделен. Следовательно, невозможно уничтожить фрейм стека до завершения задачи и пока он не перестал ссылаться на TaskPool.
@property @trusted bool isDaemon();

@property @trusted void isDaemon(bool newVal);

Эти свойства управляют тем, являются ли рабочие потоки демоновыми потоками. Демоновый поток автоматически завершается, когда завершаются все не-демоновые потоки. Не-демоновый поток будет препятствовать завершению программы, пока он не завершится.

Если какой-либо TaskPool с не-демоновыми потоками активен, то либо stop, либо finish должны быть вызваны на нём до завершения программы.

Рабочие потоки в экземпляре TaskPool, возвращённом свойством taskPool, по умолчанию являются демоновыми. Рабочие потоки вручную инициализированных пулов задач по умолчанию не-демоновые.

Примечание
Для пула размером ноль получатель произвольно возвращает true, а устанавливатель не имеет эффекта.
@property @trusted int priority();

@property @trusted void priority(int newPriority);

Эти функции позволяют получить и установить приоритет планирования операционной системы для рабочих потоков в этом TaskPool. Они передаются в core.thread.Thread.priority, поэтому заданное значение приоритета здесь означает то же самое, что и идентичное значение приоритета в core.thread.

Примечание
Для пула размером ноль получатель произвольно возвращает core.thread.Thread.PRIORITY_MIN, а устанавливатель не имеет эффекта.
@property @trusted TaskPool taskPool();

Возвращает лениво инициализированный глобальный экземпляр TaskPool. Этот метод можно безопасно вызывать одновременно из нескольких не-рабочих потоков. Рабочие потоки в этом пуле являются демоновыми потоками, что означает, что нет необходимости вызывать TaskPool.stop или TaskPool.finish перед завершением основного потока.

@property @trusted uint defaultPoolThreads();

@property @trusted void defaultPoolThreads(uint newVal);

Эти свойства получают и устанавливают количество рабочих потоков в экземпляре TaskPool возвращаемом taskPool. Значение по умолчанию равно totalCPUs - 1. Вызов устанавливателя после первого вызова taskPool не изменяет количество рабочих потоков в экземпляре, возвращаемом taskPool.

ParallelForeach!R parallel(R)(R range);

ParallelForeach!R parallel(R)(R range, size_t workUnitSize);

Функции-удобства, которые передаются в taskPool.parallel. Их цель — сделать параллельный foreach менее громоздким и более читабельным.

Пример
// Find the logarithm of every number from
// 1 to 1_000_000 in parallel, using the
// default TaskPool instance.
auto logs = new double[1_000_000];

foreach (i, ref elem; parallel(logs))
{
    elem = log(i + 1.0);
}

© 1999–2021 The D Language Foundation
Licensed under the Boost License 1.0.
https://dlang.org/phobos/std_parallelism.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API