Spec-Zone.ru › Julia 0.6

Параллельное вычисление

Большинство современных компьютеров обладают более чем одним процессором, а несколько компьютеров могут быть объединены в кластер. Использование мощности этих нескольких процессоров позволяет быстрее выполнять множество вычислений. Два основных фактора влияют на производительность: скорость самих процессоров и скорость их доступа к памяти. В кластере очевидно, что данный процессор будет иметь самый быстрый доступ к ОЗУ внутри того же компьютера (узла). Возможно, более неожиданно, что подобные проблемы актуальны и на типичном многоядерном ноутбуке из-за различий в скорости оперативной памяти и кэша. Следовательно, хорошая многопроцессорная среда должна позволять управлять «владением» частью памяти конкретным процессором. Julia предоставляет многопроцессорную среду на основе обмена сообщениями, позволяющую программам работать на нескольких процессах в отдельных доменах памяти одновременно.

Реализация обмена сообщениями в Julia отличается от других сред, таких как MPI [1]. Связь в Julia обычно «односторонняя», что означает, что программист должен явно управлять только одним процессором в двухпроцессорной операции. Кроме того, эти операции обычно не выглядят как «отправка сообщения» и «прием сообщения», а скорее напоминают операции более высокого уровня, такие как вызовы пользовательских функций.

Параллельное программирование в Julia основано на двух примитивах: удаленных ссылках и удаленных вызовах. Удаленная ссылка — это объект, который может быть использован из любого процесса для ссылки на объект, хранящийся на конкретном процессе. Удаленный вызов — это запрос одного процесса на вызов определенной функции с определенными аргументами на другом (возможно, том же) процессе.

Удаленные ссылки бывают двух типов: Future и RemoteChannel.

Удаленный вызов возвращает Future своего результата. Удаленные вызовы возвращаются немедленно; процесс, который сделал вызов, переходит к следующей операции, в то время как удаленный вызов происходит где-то в другом месте. Вы можете дождаться завершения удаленного вызова, вызвав wait() на возвращенном Future, и вы можете получить полное значение результата, используя fetch().

С другой стороны, RemoteChannel перезаписываемые. Например, несколько процессов могут координировать свою обработку, ссылаясь на один и тот же удаленный Channel.

Каждый процесс имеет связанный с ним идентификатор. Процесс, предоставляющий интерактивный приглашение Julia, всегда имеет id равный 1. Процессы, используемые по умолчанию для параллельных операций, называются «рабочими». Когда существует только один процесс, процесс 1 считается рабочим. В противном случае рабочие считаются всеми процессами, кроме процесса 1.

Давайте попробуем это. Начиная с julia -p n предоставляет n рабочих процессов на локальной машине. Как правило, имеет смысл, чтобы n равнялось количеству ядер процессора на машине.

$ ./julia -p 2

julia> r = remotecall(rand, 2, 2, 2)
Future(2, 1, 4, Nullable{Any}())

julia> s = @spawnat 2 1 .+ fetch(r)
Future(2, 1, 5, Nullable{Any}())

julia> fetch(s)
2×2 Array{Float64,2}:
 1.18526  1.50912
 1.16296  1.60607

Первый аргумент к remotecall() — функция для вызова. Большинство параллельных программ в Julia не ссылаются на конкретные процессы или количество доступных процессов, но remotecall() считается интерфейсом низкого уровня, предоставляющим более тонкий контроль. Второй аргумент к remotecall() — id процесса, который выполнит работу, а оставшиеся аргументы будут переданы в вызываемую функцию.

Как вы можете видеть, в первой строке мы попросили процесс 2 создать случайную матрицу 2x2, а во второй строке мы попросили его добавить к ней 1. Результат обоих вычислений доступен в двух будущем, r и s. Макрос @spawnat вычисляет выражение во втором аргументе в процессе, указанном первым аргументом.

Иногда вам может потребоваться значение, вычисленное удалённо, немедленно. Это обычно происходит, когда вы считываете данные из удаленного объекта, необходимые для следующей локальной операции. Функция remotecall_fetch() предназначена для этой цели. Она эквивалентна fetch(remotecall(...)) , но более эффективна.

julia> remotecall_fetch(getindex, 2, r, 1, 1)
0.18526337335308085

Помните, что getindex(r,1,1) эквивалентно r[1,1], поэтому этот вызов извлекает первый элемент будущего r.

Синтаксис remotecall() не очень удобен. Макрос @spawn упрощает задачу. Он работает с выражением, а не с функцией и сам выбирает, где выполнить операцию:

julia> r = @spawn rand(2,2)
Future(2, 1, 4, Nullable{Any}())

julia> s = @spawn 1 .+ fetch(r)
Future(3, 1, 5, Nullable{Any}())

julia> fetch(s)
2×2 Array{Float64,2}:
 1.38854  1.9098
 1.20939  1.57158

Обратите внимание, что мы использовали 1 .+ fetch(r) вместо 1 .+ r. Это потому, что нам неизвестно, где будет выполняться код, поэтому в общем случае может потребоваться fetch(), чтобы переместить r в процесс, выполняющий добавление. В данном случае @spawn достаточно умён, чтобы выполнить вычисление в процессе, которому принадлежит r, поэтому fetch() станет бесполезной операцией (никакой работы не будет выполнено).

(Стоит отметить, что @spawn не встроен, а определен в Julia как макрос. Можно определить собственные конструкции такого типа.)

Важно помнить, что после извлечения Future кеширует своё значение локально. Дальнейшие fetch() вызовы не предполагают сетевого перехода. После извлечения всех ссылок на Future удалённое сохранённое значение удаляется.

Доступность кода и загрузка пакетов

Ваш код должен быть доступен на любом процессе, который его выполняет. Например, введите следующее в приглашение Julia:

julia> function rand2(dims...)
           return 2*rand(dims...)
       end

julia> rand2(2,2)
2×2 Array{Float64,2}:
 0.153756  0.368514
 1.15119   0.918912

julia> fetch(@spawn rand2(2,2))
ERROR: RemoteException(2, CapturedException(UndefVarError(Symbol("#rand2"))
[...]

Процесс 1 знал о функции rand2, но процесс 2 — нет.

Чаще всего вы будете загружать код из файлов или пакетов, и у вас есть значительная гибкость в управлении тем, какие процессы загружают код. Рассмотрим файл DummyModule.jl, содержащий следующий код:

module DummyModule

export MyType, f

mutable struct MyType
    a::Int
end

f(x) = x^2+1

println("loaded")

end

Запустив Julia с julia -p 2, вы можете использовать это для проверки следующего:

  • include("DummyModule.jl") загружает файл только на один процесс (на том, который выполняет операцию).

  • using DummyModule заставляет модуль загрузиться на все процессы; однако модуль становится доступным только в том процессе, который выполняет операцию.

  • До тех пор, пока DummyModule загружен на процессе 2, команды, такие как

    rr = RemoteChannel(2)
    put!(rr, MyType(7))

    позволяют хранить объект типа MyType на процессе 2, даже если DummyModule не доступен в области видимости процесса 2.

Вы можете принудительно заставить команду выполняться на всех процессах, используя макрос @everywhere. Например, @everywhere можно также использовать для прямого определения функции на всех процессах:

julia> @everywhere id = myid()

julia> remotecall_fetch(()->id, 2)
2

Файл также может быть предварительно загружен на несколько процессов при запуске, а скрипт драйвера может использоваться для управления вычислениями:

julia -p <n> -L file1.jl -L file2.jl driver.jl

Процесс Julia, выполняющий скрипт драйвера в примере выше, имеет id равный 1, как и процесс, предоставляющий интерактивное приглашение.

Базовая установка Julia имеет встроенную поддержку двух типов кластеров:

  • Локальный кластер, указанный с помощью параметра -p , как показано выше.

  • Кластер, охватывающий несколько машин, с помощью параметра --machinefile . Это использует бесклеточную ssh авторизацию для запуска процессов Julia-рабочих (из той же директории, что и текущий хост) на указанных машинах.

Функции addprocs(), rmprocs(), workers() и другие доступны в качестве программного средства для добавления, удаления и запроса процессов в кластере.

Обратите внимание, что рабочие не выполняют скрипт запуска .juliarc.jl, а также не синхронизируют своё глобальное состояние (такое как глобальные переменные, новые определения методов и загруженные модули) ни с одним из других запущенных процессов.

Другие типы кластеров могут поддерживаться за счёт написания собственного пользовательского ClusterManager, как описано ниже в разделе ClusterManagers.

Перемещение данных

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

fetch() можно рассматривать как явную операцию перемещения данных, поскольку она напрямую запрашивает перемещение объекта на локальную машину. @spawn (и несколько связанных конструкций) также перемещает данные, но это не так очевидно, поэтому можно назвать неявной операцией перемещения данных. Рассмотрим эти два подхода к построению и возведению в квадрат случайной матрицы:

Метод 1:

julia> A = rand(1000,1000);

julia> Bref = @spawn A^2;

[...]

julia> fetch(Bref);

Метод 2:

julia> Bref = @spawn rand(1000,1000)^2;

[...]

julia> fetch(Bref);

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

В этом примере с игрушечными данными два метода легко различимы и выбираемы. Однако при разработке реальной программы перемещение данных может потребовать больше размышлений и, вероятно, некоторых измерений. Например, если первому процессу нужна матрица A, то первый метод может быть лучше. Или, если вычисление A является дорогостоящим и только текущий процесс его имеет, то перемещение его в другой процесс может быть неизбежным. Или, если текущему процессу почти ничего не нужно делать между @spawn и fetch(Bref), то, возможно, лучше вообще исключить параллелизм. Или представьте, что rand(1000,1000) заменено более дорогой операцией. Тогда может иметь смысл добавить еще одну @spawn инструкцию только для этого шага.

Глобальные переменные

Выражения, выполняемые удалённо через @spawn, или замыкания, определённые для удалённого выполнения с помощью remotecall, могут ссылаться на глобальные переменные. Глобальные привязки в модуле Main обрабатываются немного иначе по сравнению с глобальными привязками в других модулях. Рассмотрим следующий фрагмент кода:

A = rand(10,10)
remotecall_fetch(()->foo(A), 2)

Обратите внимание, что A — это глобальная переменная, определённая в локальном рабочем пространстве. У Рабочего узла 2 нет переменной, называемой A в Main. Действие отправки замыкания ()->foo(A) Рабочему узлу 2 приводит к определению Main.A на 2. Main.A продолжает существовать на Рабочем узле 2 даже после возвращения вызова remotecall_fetch. Удалённые вызовы с встроенными глобальными ссылками (только в модуле Main) управляют глобальными переменными следующим образом:

  • Новые глобальные привязки создаются на целевых рабочих узлах, если они ссылаются в рамках удалённого вызова.

  • Глобальные константы объявляются как константы и на удалённых узлах.

  • Глобальные переменные отправляются целевому рабочему узлу только в контексте удалённого вызова и только в том случае, если их значение изменилось. Кроме того, кластер не синхронизирует глобальные привязки между узлами. Например:

    A = rand(10,10)
    remotecall_fetch(()->foo(A), 2) # worker 2
    A = rand(10,10)
    remotecall_fetch(()->foo(A), 3) # worker 3
    A = nothing

    Выполнение приведенного выше фрагмента кода приводит к тому, что Main.A на Рабочем узле 2 имеет другое значение, чем Main.A на Рабочем узле 3, в то время как значение Main.A на узле 1 установлено в nothing.

Как вы, возможно, поняли, хотя память, связанная с глобальными переменными, может быть освобождена, когда они переопределяются на главном узле, никаких подобных действий не предпринимается на рабочих узлах, так как привязки остаются действительными. clear! может быть использован для ручного переопределения конкретных глобальных переменных на удалённых узлах на nothing после того, как они больше не требуются. Это позволит освободить связанную с ними память в рамках обычного цикла сборки мусора.

Таким образом, программы должны быть осторожны при использовании ссылок на глобальные переменные в удалённых вызовах. На самом деле, если это возможно, предпочтительнее вообще их избегать. Если вам необходимо ссылаться на глобальные переменные, рассмотрите возможность использования блоков let для локализации глобальных переменных.

Например:

julia> A = rand(10,10);

julia> remotecall_fetch(()->A, 2);

julia> B = rand(10,10);

julia> let B = B
           remotecall_fetch(()->B, 2)
       end;

julia> @spawnat 2 whos();

julia>  From worker 2:                               A    800 bytes  10×10 Array{Float64,2}
        From worker 2:                            Base               Module
        From worker 2:                            Core               Module
        From worker 2:                            Main               Module

Как видно, глобальная переменная A определена на Рабочем узле 2, но B захвачена как локальная переменная, и поэтому привязка для B не существует на Рабочем узле 2.

Параллельные map и циклы

К счастью, многие полезные параллельные вычисления не требуют перемещения данных. Типичный пример — моделирование Монте-Карло, где несколько процессов могут одновременно обрабатывать независимые испытания моделирования. Мы можем использовать @spawn для подбрасывания монет на двух процессах. Сначала напишите следующую функцию в count_heads.jl:

function count_heads(n)
    c::Int = 0
    for i = 1:n
        c += rand(Bool)
    end
    c
end

Функция count_heads просто складывает вместе n случайных битов. Вот как мы можем выполнить несколько испытаний на двух машинах и сложить результаты:

julia> @everywhere include("count_heads.jl")

julia> a = @spawn count_heads(100000000)
Future(2, 1, 6, Nullable{Any}())

julia> b = @spawn count_heads(100000000)
Future(3, 1, 7, Nullable{Any}())

julia> fetch(a)+fetch(b)
100001564

Этот пример демонстрирует мощный и часто используемый шаблон параллельного программирования. Много итераций выполняются независимо на нескольких процессах, а затем их результаты комбинируются с помощью некоторой функции. Процесс объединения называется редукцией, так как он, как правило, уменьшает тензорный ранг: вектор чисел сводится к одному числу, или матрица — к одной строке или столбцу и т. д. В коде это обычно выглядит как шаблон x = f(x,v[i]), где x — это накопитель, f — это функция редукции, а v[i] — это элементы, подлежащие редукции. Желательно, чтобы f была ассоциативной, чтобы порядок выполнения операций не влиял на результат.

Обратите внимание, что наше использование этого шаблона с count_heads может быть обобщено. Мы использовали две явные инструкции @spawn, что ограничивает параллелизм двумя процессами. Для выполнения на любом количестве процессов мы можем использовать параллельный цикл for, который можно написать в Julia с помощью @parallel вот так:

nheads = @parallel (+) for i = 1:200000000
    Int(rand(Bool))
end

Эта конструкция реализует шаблон, при котором итерации назначаются нескольким процессам и объединяются с заданной редукцией (в данном случае (+)). Результат каждой итерации принимается как значение последнего выражения в цикле. Само выражение параллельного цикла вычисляется в окончательный ответ.

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

Например, следующий код не будет работать так, как ожидается:

a = zeros(100000)
@parallel for i = 1:100000
    a[i] = i
end

Этот код не инициализирует все a, так как каждый процесс будет иметь свою копию. Такие параллельные циклы for следует избегать. К счастью, для решения этой проблемы можно использовать Общие массивы:

a = SharedArray{Float64}(10)
@parallel for i = 1:10
    a[i] = i
end

Использование «внешних» переменных в параллельных циклах вполне приемлемо, если переменные являются только для чтения:

a = randn(1000)
@parallel (+) for i = 1:100000
    f(a[rand(1:end)])
end

Здесь каждая итерация применяет f к случайной выборке из вектора a, доступного всем процессам.

Как вы могли видеть, оператор редукции можно опустить, если он не нужен. В этом случае цикл выполняется асинхронно, т. е. он порождает независимые задачи на всех доступных рабочих узлах и возвращает массив Future немедленно, не дожидаясь завершения. Звонящий может дождаться завершения Future впоследствии, вызвав fetch() на них или дождаться завершения в конце цикла, предварив его @sync, как в @sync @parallel for.

В некоторых случаях оператор редукции не нужен, и мы просто хотим применить функцию ко всем целым числам в некотором диапазоне (или, более обобщённо, ко всем элементам в некотором наборе). Это ещё одна полезная операция, называемая параллельным map, реализованная в Julia как функция pmap(). Например, мы можем вычислить сингулярные значения нескольких больших случайных матриц параллельно следующим образом:

julia> M = Matrix{Float64}[rand(1000,1000) for i = 1:10];

julia> pmap(svd, M);

Функция pmap() Julia предназначена для случая, когда каждый вызов функции выполняет большой объём работы. В отличие от этого, @parallel for может обрабатывать ситуации, где каждая итерация очень мала, возможно, просто суммирует два числа. Для параллельного вычисления обе функции pmap() и @parallel for используют только рабочие узлы. В случае @parallel for, окончательная редукция выполняется на вызывающем процессе.

Синхронизация с удалёнными ссылками

Планирование

Параллельная платформа Julia использует Задачи (также известные как сопрограммы) для переключения между несколькими вычислениями. Всякий раз, когда код выполняет операцию связи, такую как fetch() или wait(), текущая задача приостанавливается, а планировщик выбирает другую задачу для выполнения. Задача перезапускается, когда ожидаемое событие завершается.

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

В качестве примера рассмотрим вычисление сингулярных значений матриц различного размера:

julia> M = Matrix{Float64}[rand(800,800), rand(600,600), rand(800,800), rand(600,600)];

julia> pmap(svd, M);

Если один процесс обрабатывает матрицы 800×800 и другой — 600×600, мы не получим такой масштабируемости, как могли бы. Решением является создание локальной задачи для «подачи» работы каждому процессу по завершении текущей задачи. Например, рассмотрим простую реализацию pmap():

function pmap(f, lst)
    np = nprocs()  # determine the number of processes available
    n = length(lst)
    results = Vector{Any}(n)
    i = 1
    # function to produce the next work item from the queue.
    # in this case it's just an index.
    nextidx() = (idx=i; i+=1; idx)
    @sync begin
        for p=1:np
            if p != myid() || np == 1
                @async begin
                    while true
                        idx = nextidx()
                        if idx > n
                            break
                        end
                        results[idx] = remotecall_fetch(f, p, lst[idx])
                    end
                end
            end
        end
    end
    results
end

@async похожа на @spawn, но выполняет задачи только на локальном процессе. Мы используем её для создания задачи-«поставщика» для каждого процесса. Каждая задача выбирает следующий индекс, который нужно вычислить, затем ждёт завершения своего процесса, затем повторяет это, пока не исчерпаются индексы. Обратите внимание, что задачи-«поставщики» не начинают выполняться до тех пор, пока основная задача не достигнет конца блока @sync, в этот момент она передаёт управление и ждёт завершения всех локальных задач перед возвратом из функции. Задачи-«поставщики» могут совместно использовать состояние через nextidx(), поскольку они все выполняются на одном процессе. Блокировка не требуется, так как потоки планируются кооперативно, а не прерываются. Это означает, что переключение контекста происходит только в определённых точках: в данном случае, когда вызывается remotecall_fetch().

Каналы

В разделе о Task в Потоке управления обсуждалось выполнение нескольких функций в кооперативном режиме. Channel могут быть очень полезны для передачи данных между выполняемыми задачами, особенно при работе с операциями ввода-вывода.

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

Канал можно визуализировать как трубу, то есть у него есть входной и выходной концы.

  • Несколько писателей в разных задачах могут одновременно записывать в один и тот же канал через вызовы put!().

  • Несколько читателей в разных задачах могут одновременно читать данные через вызовы take!().

  • В качестве примера:

    # Given Channels c1 and c2,
    c1 = Channel(32)
    c2 = Channel(32)
    
    # and a function `foo()` which reads items from from c1, processes the item read
    # and writes a result to c2,
    function foo()
        while true
            data = take!(c1)
            [...]               # process data
            put!(c2, result)    # write out result
        end
    end
    
    # we can schedule `n` instances of `foo()` to be active concurrently.
    for _ in 1:n
        @schedule foo()
    end
  • Каналы создаются с помощью конструктора Channel{T}(sz). Канал будет содержать только объекты типа T. Если тип не указан, канал может содержать объекты любого типа. sz относится к максимальному количеству элементов, которые могут храниться в канале в любой момент. Например, Channel(32) создает канал, который может хранить максимум 32 объекта любого типа. Channel{MyType}(64) может хранить до 64 объектов MyType в любой момент.

  • Если Channel пуст, читатели (при вызове take!()) будут блокироваться до тех пор, пока данные не станут доступны.

  • Если Channel заполнен, писатели (при вызове put!()) будут блокироваться до тех пор, пока не освободится место.

  • isready() проверяет наличие каких-либо объектов в канале, а wait() ожидает, пока объект не станет доступным.

  • Channel изначально находится в открытом состоянии. Это означает, что к нему можно свободно читать и писать с помощью вызовов take!() и put!(). close() закрывает Channel. В закрытом Channel, put!() завершится с ошибкой. Например:

julia> c = Channel(2);

julia> put!(c, 1) # `put!` on an open channel succeeds
1

julia> close(c);

julia> put!(c, 2) # `put!` on a closed channel throws an exception.
ERROR: InvalidStateException("Channel is closed.",:closed)
[...]
  • take!() и fetch() (которая извлекает, но не удаляет значение) в закрытом канале успешно возвращают любые существующие значения до тех пор, пока он не опустеет. Продолжая приведенный выше пример:

julia> fetch(c) # Any number of `fetch` calls succeed.
1

julia> fetch(c)
1

julia> take!(c) # The first `take!` removes the value.
1

julia> take!(c) # No more data available on a closed channel.
ERROR: InvalidStateException("Channel is closed.",:closed)
[...]

Канал может использоваться как итерируемый объект в цикле for, в этом случае цикл работает до тех пор, пока у канала есть данные или он открыт. Переменная цикла принимает все значения, добавленные в Channel. Цикл for завершается после того, как Channel закрыт и опустел.

Например, следующее приведет к тому, что цикл for будет ожидать новых данных:

julia> c = Channel{Int}(10);

julia> foreach(i->put!(c, i), 1:3) # add a few entries

julia> data = [i for i in c]

в то время как это вернётся после чтения всех данных:

julia> c = Channel{Int}(10);

julia> foreach(i->put!(c, i), 1:3); # add a few entries

julia> close(c);                    # `for` loops can exit

julia> data = [i for i in c]
3-element Array{Int64,1}:
 1
 2
 3

Рассмотрим простой пример использования каналов для межзадачной коммуникации. Мы запускаем 4 задачи для обработки данных из одного канала jobs. Задачи, идентифицируемые по идентификатору (job_id), записываются в канал. Каждая задача в этом моделировании считывает job_id, ждёт случайное время и записывает обратно кортеж из job_id и моделируемого времени в канал результатов. Наконец, все results выводятся на экран.

julia> const jobs = Channel{Int}(32);

julia> const results = Channel{Tuple}(32);

julia> function do_work()
           for job_id in jobs
               exec_time = rand()
               sleep(exec_time)                # simulates elapsed time doing actual work
                                               # typically performed externally.
               put!(results, (job_id, exec_time))
           end
       end;

julia> function make_jobs(n)
           for i in 1:n
               put!(jobs, i)
           end
       end;

julia> n = 12;

julia> @schedule make_jobs(n); # feed the jobs channel with "n" jobs

julia> for i in 1:4 # start 4 tasks to process requests in parallel
           @schedule do_work()
       end

julia> @elapsed while n > 0 # print out results
           job_id, exec_time = take!(results)
           println("$job_id finished in $(round(exec_time,2)) seconds")
           n = n - 1
       end
4 finished in 0.22 seconds
3 finished in 0.45 seconds
1 finished in 0.5 seconds
7 finished in 0.14 seconds
2 finished in 0.78 seconds
5 finished in 0.9 seconds
9 finished in 0.36 seconds
6 finished in 0.87 seconds
8 finished in 0.79 seconds
10 finished in 0.64 seconds
12 finished in 0.5 seconds
11 finished in 0.97 seconds
0.029772311

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

Удаленные ссылки и абстрактные каналы

Удаленные ссылки всегда ссылаются на реализацию AbstractChannel.

Для реализации put!(), take!(), fetch(), isready() и wait() требуется конкретная реализация AbstractChannel (например, Channel). Удаленный объект, на который ссылается Future, хранится в Channel{Any}(1), т. е. в Channel размером 1, способном хранить объекты типа Any.

RemoteChannel, которая может перезаписываться, может указывать на каналы любого типа и размера, или на любую другую реализацию AbstractChannel.

Конструктор RemoteChannel(f::Function, pid)() позволяет нам создавать ссылки на каналы, содержащие более одного значения определенного типа. f() — это функция, выполняемая на pid, и она должна возвращать AbstractChannel.

Например, RemoteChannel(()->Channel{Int}(10), pid), вернёт ссылку на канал типа Int и размером 10. Канал существует на узле pid.

Методы put!(), take!(), fetch(), isready() и wait() для RemoteChannel делегируются на хранилище в удаленном процессе.

RemoteChannel можно использовать для ссылки на реализованные пользователем объекты AbstractChannel . Простой пример этого показан в examples/dictchannel.jl, который использует словарь в качестве удалённого хранилища.

Каналы и удалённые каналы

  • Channel локален для процесса. Узел 2 не может напрямую ссылаться на Channel на узле 3 и наоборот. RemoteChannel, однако, может помещать и извлекать значения между узлами.

  • RemoteChannel можно рассматривать как дескриптор для Channel.

  • Идентификатор процесса, pid, связанный с RemoteChannel, определяет процесс, где существует базовая структура данных, то есть базовый Channel.

  • Любой процесс со ссылкой на RemoteChannel может помещать и извлекать элементы из канала. Данные автоматически отправляются (или извлекаются) в процесс, к которому относится RemoteChannel.

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

  • С другой стороны, сериализация RemoteChannel включает только сериализацию идентификатора, который идентифицирует местоположение и экземпляр Channel, на который ссылается дескриптор. Десериализованный объект RemoteChannel (на любом узле) поэтому также ссылается на то же хранилище, что и исходный.

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

Мы запускаем 4 узла для обработки одного канала jobs удалённого канала. Задачи, идентифицируемые по идентификатору (job_id), записываются в канал. Каждая удалённо выполняемая задача в этом моделировании считывает job_id, ждёт случайное время и возвращает кортеж из job_id, затраченного времени и собственного pid в канал результатов. Наконец, все results выводятся на мастер-процессе.

julia> addprocs(4); # add worker processes

julia> const jobs = RemoteChannel(()->Channel{Int}(32));

julia> const results = RemoteChannel(()->Channel{Tuple}(32));

julia> @everywhere function do_work(jobs, results) # define work function everywhere
           while true
               job_id = take!(jobs)
               exec_time = rand()
               sleep(exec_time) # simulates elapsed time doing actual work
               put!(results, (job_id, exec_time, myid()))
           end
       end

julia> function make_jobs(n)
           for i in 1:n
               put!(jobs, i)
           end
       end;

julia> n = 12;

julia> @schedule make_jobs(n); # feed the jobs channel with "n" jobs

julia> for p in workers() # start tasks on the workers to process requests in parallel
           @async remote_do(do_work, p, jobs, results)
       end

julia> @elapsed while n > 0 # print out results
           job_id, exec_time, where = take!(results)
           println("$job_id finished in $(round(exec_time,2)) seconds on worker $where")
           n = n - 1
       end
1 finished in 0.18 seconds on worker 4
2 finished in 0.26 seconds on worker 5
6 finished in 0.12 seconds on worker 4
7 finished in 0.18 seconds on worker 4
5 finished in 0.35 seconds on worker 5
4 finished in 0.68 seconds on worker 2
3 finished in 0.73 seconds on worker 3
11 finished in 0.01 seconds on worker 3
12 finished in 0.02 seconds on worker 3
9 finished in 0.26 seconds on worker 5
8 finished in 0.57 seconds on worker 4
10 finished in 0.58 seconds on worker 2
0.055971741

Удаленные ссылки и распределённое удаление мусора

Объекты, на которые ссылаются удалённые ссылки, могут быть освобождены только тогда, когда все ссылки в кластере удалены.

Узел, на котором хранится значение, отслеживает, какие из узлов имеют на него ссылку. Каждый раз, когда RemoteChannel или (неизвлечённый) Future сериализуется на узел, узел, на который указывает ссылка, получает уведомление. И каждый раз, когда RemoteChannel или (неизвлечённый) Future удаляется локально, узел, владеющий значением, получает уведомление.

Уведомления осуществляются путём отправки сообщений отслеживания — сообщение «добавить ссылку» при сериализации ссылки на другой процесс и сообщение «удалить ссылку» при локальном удалении ссылки.

Поскольку Future записываются один раз и кешируются локально, действие fetch() Future также обновляет информацию об отслеживании ссылок на узле, владеющем значением.

Узел, владеющий значением, освобождает его после того, как все ссылки на него будут удалены.

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

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

В случае удаленных ссылок размер локального объекта ссылки довольно небольшой, в то время как значение, хранящееся на удаленном узле, может быть достаточно большим. Поскольку локальный объект может не быть собран немедленно, хорошей практикой является явное вызов finalize() на локальных экземплярах RemoteChannel или на невыгруженных Future. Поскольку вызов fetch() на Future также удаляет его ссылку из удаленного хранилища, это не требуется для выгруженных Future. Явное вызов finalize() приводит к немедленному сообщению, отправленному на удаленный узел, чтобы он удалил свою ссылку на значение.

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

Общие массивы

Общие массивы используют системную общую память для отображения одного и того же массива в нескольких процессах. Несмотря на некоторые сходства с DArray, поведение SharedArray существенно отличается. В DArray каждый процесс имеет локальный доступ только к фрагменту данных, и никакие два процесса не делят один и тот же фрагмент; в отличие от SharedArray, каждый «участвующий» процесс имеет доступ ко всему массиву. SharedArray — хороший выбор, когда вам нужно, чтобы большой объем данных был совместно доступен двум или более процессам на одной машине.

SharedArray индексирование (присваивание и доступ к значениям) работает так же, как и с обычными массивами, и является эффективным, потому что основная память доступна локальному процессу. Поэтому большинство алгоритмов естественным образом работают с SharedArray, хотя бы в однопроцессном режиме. В тех случаях, когда алгоритм настаивает на вводе Array, базовый массив можно извлечь из SharedArray вызвав sdata(). Для других типов AbstractArray, sdata() просто возвращает сам объект, поэтому использовать sdata() на любом объекте типа Array безопасно.

Конструктор для общего массива имеет вид:

SharedArray{T,N}(dims::NTuple; init=false, pids=Int[])

что создает N-мерный общий массив типа бит T и размера dims среди процессов, указанных в pids. В отличие от распределенных массивов, общий массив доступен только для тех участвующих рабочих процессов, которые указаны в аргументе pids (и для создающего процесса тоже, если он находится на том же хосте).

Если указана функция init с сигнатурой initfn(S::SharedArray), она вызывается на всех участвующих рабочих процессах. Вы можете указать, чтобы каждый рабочий процесс выполнял функцию init на отдельной части массива, тем самым обеспечивая распараллеливание инициализации.

Вот краткий пример:

julia> addprocs(3)
3-element Array{Int64,1}:
 2
 3
 4

julia> S = SharedArray{Int,2}((3,4), init = S -> S[Base.localindexes(S)] = myid())
3×4 SharedArray{Int64,2}:
 2  2  3  4
 2  3  3  4
 2  3  4  4

julia> S[3,2] = 7
7

julia> S
3×4 SharedArray{Int64,2}:
 2  2  3  4
 2  3  3  4
 2  7  4  4

Base.localindexes() предоставляет несвязные одномерные диапазоны индексов и иногда удобен для разделения задач между процессами. Конечно, вы можете разделить работу любым способом:

julia> S = SharedArray{Int,2}((3,4), init = S -> S[indexpids(S):length(procs(S)):length(S)] = myid())
3×4 SharedArray{Int64,2}:
 2  2  2  2
 3  3  3  3
 4  4  4  4

Поскольку все процессы имеют доступ к основным данным, вам нужно быть осторожным, чтобы не создать конфликтов. Например:

@sync begin
    for p in procs(S)
        @async begin
            remotecall_wait(fill!, p, S, p)
        end
    end
end

приведет к неопределенному поведению. Поскольку каждый процесс заполняет весь массив своим собственным pid, тот процесс, который последним выполнится (для любого конкретного элемента S), сохранит свой pid.

В качестве более расширенного и сложного примера рассмотрим выполнение следующего «ядра» параллельно:

q[i,j,t+1] = q[i,j,t] + u[i,j,t]

В этом случае, если мы попытаемся разделить работу с использованием одномерного индекса, мы, вероятно, столкнемся с проблемами: если q[i,j,t] находится в конце блока, назначенного одному рабочему процессу, а q[i,j,t+1] — в начале блока, назначенного другому, очень вероятно, что q[i,j,t] не будет готов к моменту, когда он понадобится для вычисления q[i,j,t+1]. В таких случаях лучше вручную разбить массив на куски. Давайте разделим его по второму измерению. Определим функцию, которая возвращает (irange, jrange) индексы, назначенные этому рабочему процессу:

julia> @everywhere function myrange(q::SharedArray)
           idx = indexpids(q)
           if idx == 0 # This worker is not assigned a piece
               return 1:0, 1:0
           end
           nchunks = length(procs(q))
           splits = [round(Int, s) for s in linspace(0,size(q,2),nchunks+1)]
           1:size(q,1), splits[idx]+1:splits[idx+1]
       end

Далее, определим ядро:

julia> @everywhere function advection_chunk!(q, u, irange, jrange, trange)
           @show (irange, jrange, trange)  # display so we can see what's happening
           for t in trange, j in jrange, i in irange
               q[i,j,t+1] = q[i,j,t] + u[i,j,t]
           end
           q
       end

Мы также определим удобную обёртку для реализации SharedArray

julia> @everywhere advection_shared_chunk!(q, u) =
           advection_chunk!(q, u, myrange(q)..., 1:size(q,3)-1)

Теперь давайте сравним три разных версии, одна из которых работает в одном процессе:

julia> advection_serial!(q, u) = advection_chunk!(q, u, 1:size(q,1), 1:size(q,2), 1:size(q,3)-1);

другая использует @parallel:

julia> function advection_parallel!(q, u)
           for t = 1:size(q,3)-1
               @sync @parallel for j = 1:size(q,2)
                   for i = 1:size(q,1)
                       q[i,j,t+1]= q[i,j,t] + u[i,j,t]
                   end
               end
           end
           q
       end;

и одна, которая делегирует по частям:

julia> function advection_shared!(q, u)
           @sync begin
               for p in procs(q)
                   @async remotecall_wait(advection_shared_chunk!, p, q, u)
               end
           end
           q
       end;

Если мы создаем SharedArray и замеряем время выполнения этих функций, мы получим следующие результаты (с julia -p 4):

julia> q = SharedArray{Float64,3}((500,500,500));

julia> u = SharedArray{Float64,3}((500,500,500));

Выполните функции один раз, чтобы выполнить JIT-компиляцию, и замерьте время их выполнения во второй раз с помощью @time:

julia> @time advection_serial!(q, u);
(irange,jrange,trange) = (1:500,1:500,1:499)
 830.220 milliseconds (216 allocations: 13820 bytes)

julia> @time advection_parallel!(q, u);
   2.495 seconds      (3999 k allocations: 289 MB, 2.09% gc time)

julia> @time advection_shared!(q,u);
        From worker 2:       (irange,jrange,trange) = (1:500,1:125,1:499)
        From worker 4:       (irange,jrange,trange) = (1:500,251:375,1:499)
        From worker 3:       (irange,jrange,trange) = (1:500,126:250,1:499)
        From worker 5:       (irange,jrange,trange) = (1:500,376:500,1:499)
 238.119 milliseconds (2264 allocations: 169 KB)

Главное преимущество advection_shared! заключается в минимизации трафика между рабочими процессами, позволяя каждому из них вычислять более длительное время на назначенной части.

Общие массивы и распределенное управление сборкой мусора

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

ClusterManagers

Запуск, управление и сетевое взаимодействие процессов Julia в логическом кластере выполняется с помощью менеджеров кластеров. ClusterManager отвечает за

  • запуск рабочих процессов в среде кластера

  • управление событиями в течение срока службы каждого рабочего процесса

  • необязательно, предоставление транспорта данных

Кластер Julia имеет следующие характеристики:

  • Изначальный процесс Julia, также называемый master, является особым и имеет id равный 1.

  • Только процесс master может добавлять или удалять рабочие процессы.

  • Все процессы могут напрямую общаться друг с другом.

Соединения между рабочими процессами (используя встроенный транспорт TCP/IP) устанавливаются следующим образом:

  • addprocs() вызывается в главном процессе с объектом ClusterManager.

  • addprocs() вызывает соответствующий метод launch(), который запускает необходимое количество рабочих процессов на соответствующих машинах.

  • Каждый рабочий процесс начинает прослушивание на свободном порту и записывает информацию о своем хосте и порте в STDOUT.

  • Менеджер кластера перехватывает STDOUT каждого рабочего процесса и делает его доступным для главного процесса.

  • Главный процесс анализирует эту информацию и устанавливает подключения TCP/IP к каждому рабочему процессу.

  • Каждый рабочий процесс также уведомляется об остальных рабочих процессах в кластере.

  • Каждый рабочий процесс подключается ко всем рабочим процессам, чей id меньше, чем у самого рабочего процесса id.

  • Таким образом, создается сеть «сетка», в которой каждый рабочий процесс напрямую подключен к каждому другому.

Хотя по умолчанию транспортный уровень использует стандартный TCPSocket, возможно, что кластер Julia предоставляет свой собственный транспорт.

Julia предоставляет два встроенных менеджера кластеров:

  • LocalManager, используемый при вызове addprocs() или addprocs(np::Integer)

  • SSHManager, используемый при вызове addprocs(hostnames::Array) со списком имен хостов

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

Таким образом, минимальному менеджеру кластера потребуется:

  • быть подтипом абстрактного ClusterManager

  • реализовать launch(), метод, ответственный за запуск новых рабочих процессов

  • реализовать manage(), который вызывается при различных событиях в течение срока службы рабочего процесса (например, отправка сигнала прерывания)

addprocs(manager::FooManager) требует, чтобы FooManager реализовал:

function launch(manager::FooManager, params::Dict, launched::Array, c::Condition)
    [...]
end

function manage(manager::FooManager, id::Integer, config::WorkerConfig, op::Symbol)
    [...]
end

В качестве примера давайте посмотрим, как реализован LocalManager, менеджер, ответственный за запуск рабочих процессов на одном хосте:

struct LocalManager <: ClusterManager
    np::Integer
end

function launch(manager::LocalManager, params::Dict, launched::Array, c::Condition)
    [...]
end

function manage(manager::LocalManager, id::Integer, config::WorkerConfig, op::Symbol)
    [...]
end

Метод launch() принимает следующие аргументы:

  • manager::ClusterManager: менеджер кластера, с которым вызывается addprocs()

  • params::Dict: все ключевые аргументы, переданные в addprocs()

  • launched::Array: массив, в который необходимо добавить один или несколько объектов WorkerConfig

  • c::Condition: переменная условия, которая должна быть уведомлена по мере запуска рабочих процессов

Метод launch() вызывается асинхронно в отдельной задаче. Завершение этой задачи сигнализирует о запуске всех запрошенных рабочих процессов. Следовательно, функция launch() ДОЛЖНА завершиться, как только все запрошенные рабочие процессы будут запущены.

Новоинициализированные рабочие процессы подключаются друг к другу и к мастер-процессу способом «все-со-всеми». Указание аргумента команды --worker <cookie> приводит к тому, что запущенные процессы инициализируют себя как рабочие процессы, а подключения устанавливаются через сокеты TCP/IP. По желанию, также может быть указано --bind-to bind_addr[:port], чтобы позволить другим рабочим процессам подключиться к нему по указанному bind_addr и port. Это полезно для хостов с несколькими сетевыми интерфейсами.

В качестве примера транспортного средства, отличного от TCP/IP, реализация может выбрать использование MPI, в этом случае --worker указывать не нужно. Вместо этого, новоинициализированные рабочие процессы должны вызвать init_worker(cookie) перед использованием любых параллельных конструкций.

Для каждого запущенного рабочего процесса метод launch() должен добавить объект WorkerConfig (с соответствующими инициализированными полями) в launched

mutable struct WorkerConfig
    # Common fields relevant to all cluster managers
    io::Nullable{IO}
    host::Nullable{AbstractString}
    port::Nullable{Integer}

    # Used when launching additional workers at a host
    count::Nullable{Union{Int, Symbol}}
    exename::Nullable{AbstractString}
    exeflags::Nullable{Cmd}

    # External cluster managers can use this to store information at a per-worker level
    # Can be a dict if multiple fields need to be stored.
    userdata::Nullable{Any}

    # SSHManager / SSH tunnel connections to workers
    tunnel::Nullable{Bool}
    bind_addr::Nullable{AbstractString}
    sshflags::Nullable{Cmd}
    max_parallel::Nullable{Integer}

    connect_at::Nullable{Any}

    [...]
end

Большинство полей в WorkerConfig используются встроенными менеджерами. Пользовательские менеджеры кластеров обычно указывают только io или host / port:

  • Если указано io, оно используется для чтения информации о хосте/порту. Рабочий процесс Julia выводит свой адрес привязки и порт при запуске. Это позволяет рабочим процессам Julia прослушивать любой свободный порт, вместо необходимости ручной настройки портов рабочих процессов.

  • Если io не указано, используются host и port для подключения.

  • count, exename и exeflags важны для запуска дополнительных рабочих процессов из рабочего процесса. Например, менеджер кластера может запустить по одному рабочему процессу на узел и использовать его для запуска дополнительных рабочих процессов.

    • count со значением целого числа n запустит всего n рабочих процессов.

    • count со значением :auto запустит столько рабочих процессов, сколько ядер на данном компьютере.

    • exename — имя исполняемого файла julia включая полный путь.

    • exeflags должно быть установлено в требуемые аргументы командной строки для новых рабочих процессов.

  • tunnel, bind_addr, sshflags и max_parallel используются, когда для подключения к рабочим процессам от мастер-процесса требуется SSH-туннель.

  • userdata предоставляется для пользовательских менеджеров кластеров для хранения собственной информации, специфичной для рабочих процессов.

manage(manager::FooManager, id::Integer, config::WorkerConfig, op::Symbol) вызывается в разное время в течение жизненного цикла рабочего процесса с соответствующими значениями op:

  • со значениями :register/:deregister при добавлении/удалении рабочего процесса из пула рабочих процессов Julia.

  • со значением :interrupt при вызове interrupt(workers). ClusterManager должен сигнализировать соответствующему рабочему процессу с помощью сигнала прерывания.

  • со значением :finalize для целей очистки.

Менеджеры кластеров с пользовательскими транспортными средствами

Замена стандартных подключений TCP/IP «все-со-всеми» на пользовательский транспортный слой немного сложнее. Каждый процесс Julia имеет столько же задач связи, сколько рабочих процессов, к которым он подключен. Например, рассмотрим кластер Julia из 32 процессов в сети «все-со-всеми»:

  • Каждый процесс Julia имеет 31 задачу связи.

  • Каждая задача обрабатывает все входящие сообщения от одного удаленного рабочего процесса в цикле обработки сообщений.

  • Цикл обработки сообщений ожидает на объекте IO (например, TCPSocket в стандартной реализации), считывает все сообщение, обрабатывает его и ожидает следующего.

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

Для замены стандартного транспорта новой реализации необходимо настроить подключения к удаленным рабочим процессам и предоставить соответствующие объекты IO , на которых циклы обработки сообщений могут ожидать. Реализуемые обратные вызовы, специфичные для менеджера:

connect(manager::FooManager, pid::Integer, config::WorkerConfig)
kill(manager::FooManager, pid::Int, config::WorkerConfig)

Стандартная реализация (использующая сокеты TCP/IP) реализована как connect(manager::ClusterManager, pid::Integer, config::WorkerConfig).

connect должен вернуть пару объектов IO , один для чтения данных, отправленных рабочим процессом pid, а другой — для записи данных, которые необходимо отправить рабочему процессу pid. Пользовательские менеджеры кластеров могут использовать внутренний BufferStream в качестве канала для передачи данных между пользовательским, возможно не-IO транспортным средством и встроенной параллельной инфраструктурой Julia.

BufferStream — это внутренний IOBuffer, который ведет себя как IO — это поток, который может обрабатываться асинхронно.

Папка examples/clustermanager/0mq содержит пример использования ZeroMQ для подключения рабочих процессов Julia в топологии звезда с брокером 0MQ посередине. Примечание: процессы Julia все еще логически подключены друг к другу — любой рабочий процесс может напрямую отправлять сообщение любому другому рабочему процессу, не зная, что используется 0MQ в качестве транспортного слоя.

При использовании пользовательских транспортных средств:

  • Рабочие процессы Julia не должны запускаться с --worker. Запуск с --worker приведет к тому, что новозапущенные рабочие процессы по умолчанию будут использовать реализацию транспортного средства TCP/IP.

  • Для каждого входящего логического подключения к рабочему процессу необходимо вызвать Base.process_messages(rd::IO, wr::IO)(). Это запускает новую задачу, которая обрабатывает чтение и запись сообщений от/в рабочий процесс, представленный объектами IO.

  • init_worker(cookie, manager::FooManager) ОБЯЗАТЕЛЬНО должен вызываться в рамках инициализации процесса рабочего процесса.

  • Поле connect_at::Any в WorkerConfig может быть установлено менеджером кластера при вызове launch(). Значение этого поля передается во все connect() обратные вызовы. Как правило, оно содержит информацию о том, как подключиться к рабочему процессу. Например, транспортный слой TCP/IP использует это поле для указания кортежа (host, port) для подключения к рабочему процессу.

kill(manager, pid, config) вызывается для удаления рабочего процесса из кластера. В мастер-процессе соответствующие объекты IO должны быть закрыты реализацией для обеспечения правильной очистки. Стандартная реализация просто выполняет вызов exit() для указанного удаленного рабочего процесса.

examples/clustermanager/simple — это пример простой реализации с использованием сокетов UNIX-домена для настройки кластера.

Требования к сети для LocalManager и SSHManager

Кластеры Julia предназначены для выполнения в уже защищенных средах на инфраструктуре, такой как локальные ноутбуки, отдельные кластеры или облако. В этом разделе рассматриваются требования к безопасности сети для встроенных LocalManager и SSHManager:

  • Мастер-процесс не прослушивает ни один порт. Он только подключается к рабочим процессам.

  • Каждый рабочий процесс привязывается только к одному локальному интерфейсу и прослушивает первый свободный порт, начиная с 9009.

  • LocalManager, используемый addprocs(N), по умолчанию привязывается только к интерфейсу обратной связи. Это означает, что рабочие процессы, запущенные позже на удаленных хостах (или кем-либо с дурными намерениями), не могут подключиться к кластеру. Обращение addprocs(4) за которым следует addprocs(["remote_host"]) не удастся. Некоторые пользователям может потребоваться создать кластер, включающий их локальную систему и несколько удаленных систем. Это можно сделать, явно запросив LocalManager привязаться к внешнему сетевому интерфейсу с помощью ключевого аргумента restrict: addprocs(4; restrict=false).

  • SSHManager, используемый addprocs(list_of_remote_hosts), запускает рабочие процессы на удаленных хостах через SSH. По умолчанию SSH используется только для запуска рабочих процессов Julia. Последующие подключения мастер-рабочий и рабочий-рабочий используют обычные незашифрованные сокеты TCP/IP. На удаленных хостах должна быть включена функция безпарольного входа. Дополнительные флаги SSH или учетные данные могут быть указаны с помощью ключевого аргумента sshflags.

  • addprocs(list_of_remote_hosts; tunnel=true, sshflags=<ssh keys and other flags>) полезно, когда мы хотим использовать SSH-соединения и для мастер-рабочего процесса. Типичный сценарий для этого — локальный ноутбук, на котором выполняется REPL Julia (то есть мастер), а остальной кластер находится в облаке, скажем, на Amazon EC2. В этом случае необходимо открыть только порт 22 в удаленном кластере в сочетании с SSH-клиентом, аутентифицированным по инфраструктуре открытых ключей (PKI). Учетные данные для аутентификации можно предоставить через sshflags, например sshflags=`-e <keyfile>`.

    Обратите внимание, что соединения рабочий-рабочий по-прежнему являются обычным TCP, и локальная политика безопасности на удаленном кластере должна допускать свободные соединения между узлами рабочих процессов, по крайней мере, для портов 9009 и выше.

    Защита и шифрование всего трафика рабочий-рабочий (через SSH) или шифрование отдельных сообщений может быть реализовано с помощью пользовательского ClusterManager.

Файл cookie кластера

Все процессы в кластере используют один и тот же файл cookie, который по умолчанию является случайно сгенерированной строкой в мастер-процессе:

  • Base.cluster_cookie() возвращает файл cookie, а Base.cluster_cookie(cookie)() устанавливает и возвращает новый.

  • Все подключения аутентифицируются с обеих сторон для обеспечения того, что к друг другу могут подключаться только рабочие процессы, запущенные мастером.

  • Файл cookie должен быть передан рабочим процессам при запуске в качестве аргумента --worker <cookie>. Пользовательские ClusterManager могут получить файл cookie в мастер-процессе, вызвав Base.cluster_cookie(). Менеджеры кластеров, не использующие стандартный транспорт TCP/IP (и, следовательно, не указывающие --worker), должны вызвать init_worker(cookie, manager) с тем же файлом cookie, что и в мастер-процессе.

Обратите внимание, что среды, требующие более высокого уровня безопасности, могут реализовать это с помощью пользовательского ClusterManager . Например, файлы cookie могут быть предварительно распределены и поэтому не указываться как аргумент запуска.

Указание топологии сети (экспериментальная функция)

Ключевой аргумент topology, переданный addprocs , используется для указания того, как рабочие процессы должны быть подключены друг к другу:

  • :all_to_all, по умолчанию: все рабочие процессы соединены друг с другом.

  • :master_slave: только процесс драйвера, т.е. pid 1, имеет подключения к рабочим процессам.

  • :custom: метод launch менеджера кластера определяет топологию соединения через поля ident и connect_idents в WorkerConfig. Рабочий процесс с предоставленным менеджером кластера идентификатором ident подключится ко всем рабочим процессам, указанным в connect_idents.

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

Многопоточность (Экспериментальная)

В дополнение к задачам, удалённым вызовам и удалённым ссылкам, Julia из v0.5 будет встраивать поддержку многопоточности. Обратите внимание, что этот раздел экспериментальный, и интерфейсы могут измениться в будущем.

Настройка

По умолчанию Julia запускается с одним потоком выполнения. Это можно проверить, используя команду Threads.nthreads():

julia> Threads.nthreads()
1

Количество потоков, с которыми запускается Julia, контролируется переменной среды, называемой JULIA_NUM_THREADS. Теперь запустим Julia с 4 потоками:

export JULIA_NUM_THREADS=4

(Вышеуказанная команда работает в оболочках Bourne на Linux и OSX. Обратите внимание, что если вы используете оболочку C Shell на этих платформах, вы должны использовать ключевое слово set вместо export. Если вы работаете в Windows, запустите командную строку в расположении julia.exe и используйте set вместо export.)

Давайте проверим, что у нас есть 4 потока.

julia> Threads.nthreads()
4

Но мы сейчас в главном потоке. Чтобы проверить, используем команду Threads.threadid()

julia> Threads.threadid()
1

Макрос @threads

Давайте рассмотрим простой пример с нашими встроенными потоками. Создадим массив нулей:

julia> a = zeros(10)
10-element Array{Float64,1}:
 0.0
 0.0
 0.0
 0.0
 0.0
 0.0
 0.0
 0.0
 0.0
 0.0

Давайте одновременно обработаем этот массив, используя 4 потока. Каждый поток запишет свой идентификатор потока в каждое место.

Julia поддерживает параллельные циклы, используя макрос Threads.@threads. Этот макрос прикрепляется перед циклом for для указания Julia, что этот цикл является многопоточной областью:

julia> Threads.@threads for i = 1:10
           a[i] = Threads.threadid()
       end

Пространство итераций разделяется между потоками, после чего каждый поток записывает свой идентификатор потока в свои назначенные места:

julia> a
10-element Array{Float64,1}:
 1.0
 1.0
 1.0
 2.0
 2.0
 2.0
 3.0
 3.0
 4.0
 4.0

Обратите внимание, что Threads.@threads не имеет необязательного параметра редукции, как @parallel.

@threadcall (Экспериментальная)

Все задачи ввода-вывода, таймеры, команды REPL и т.д. мультиплексируются на один поток ОС через цикл событий. Модифицированная версия libuv (http://docs.libuv.org/en/v1.x/) предоставляет эту функциональность. Точки возврата обеспечивают кооперативное планирование нескольких задач в одном потоке ОС. Задачи ввода-вывода и таймеры возвращаются неявно, ожидая возникновения события. Явный вызов yield() позволяет планировать другие задачи.

Таким образом, задача, выполняющая ccall, фактически предотвращает планировщик Julia от выполнения других задач до возвращения вызова. Это верно для всех вызовов во внешние библиотеки. Исключениями являются вызовы в пользовательский C-код, который вызывает обратно в Julia (который может затем возвращаться) или C-код, который вызывает jl_yield() (C-эквивалент yield()).

Обратите внимание, что, хотя код Julia выполняется в одном потоке (по умолчанию), библиотеки, используемые Julia, могут запускать свои внутренние потоки. Например, библиотека BLAS может запустить столько потоков, сколько ядер на машине.

Макрос @threadcall решает сценарии, в которых мы не хотим, чтобы ccall блокировал основной цикл событий Julia. Он планирует выполнение C-функции в отдельном потоке. Для этого используется пул потоков с размером по умолчанию 4. Размер пула потоков контролируется переменной среды UV_THREADPOOL_SIZE. В ожидании свободного потока и во время выполнения функции, как только поток доступен, запрашиваемая задача (в главном цикле событий Julia) уступает место другим задачам. Обратите внимание, что @threadcall не возвращается до завершения выполнения. С точки зрения пользователя это, следовательно, блокирующий вызов, как и другие API Julia.

Очень важно, чтобы вызываемая функция не вызывала обратно в Julia.

@threadcall может быть удалена/изменена в будущих версиях Julia.

[1]

В этом контексте MPI относится к стандарту MPI-1. Начиная с MPI-2, комитет по стандартам MPI ввёл новую группу механизмов связи, объединённых под названием Доступ к удалённой памяти (RMA). Мотивация для добавления RMA в стандарт MPI заключалась в облегчении односторонних коммуникационных шаблонов. Для дополнительной информации о последнем стандарте MPI см. http://mpi-forum.org/docs.

© 2009–2016 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/release-0.6/manual/parallel-computing/

Spec-Zone.ru

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