Spec-Zone.ru › Julia 0.5

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

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

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

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

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

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

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

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

$ ./julia -p 2

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

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

julia> fetch(s)
2×2 Array{Float64,2}:
 1.60401  1.50111
 1.17457  1.15741

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

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

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

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

Помните, что 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)
1.10824216411304866 1.13798233877923116
1.12376292706355074 1.18750497916607167

Обратите внимание, что мы использовали 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: On worker 2:
function rand2 not defined on process 2

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

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

module DummyModule

export MyType, f

type 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, выполняющий скрипт драйвера в примере выше. Процессы, используемые по умолчанию для параллельных операций, называются «рабочими». Когда есть только один процесс, процесс 1 считается рабочим. В противном случае рабочие процессы — это все процессы, кроме процесса 1.

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

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

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

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

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

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

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

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

# method 1
A = rand(1000,1000)
Bref = @spawn A^2
...
fetch(Bref)

# method 2
Bref = @spawn rand(1000,1000)^2
...
fetch(Bref)

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

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

Параллельные 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 случайных битов. Вот как мы можем выполнить несколько испытаний на двух машинах и сложить результаты:

@everywhere include("count_heads.jl")

a = @spawn count_heads(100000000)
b = @spawn count_heads(100000000)
fetch(a)+fetch(b)

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

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

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(). Например, мы могли бы вычислить сингулярные значения нескольких больших случайных матриц параллельно следующим образом:

M = Matrix{Float64}[rand(1000,1000) for i=1:10]
pmap(svd, M)

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

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

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

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

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

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

M = Matrix{Float64}[rand(800,800), rand(600,600), rand(800,800), rand(600,600)]
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().

Каналы

Каналы обеспечивают быстрый способ межзадачной коммуникации. Channel{T}(n::Int) — это общая очередь максимальной длины n, хранящая объекты типа T. Несколько читателей могут читать из Channel через fetch() и take!(). Несколько писателей могут добавлять в Channel через put!(). isready() проверяет наличие каких-либо объектов в канале, в то время как wait() ожидает, пока объект станет доступным. close() закрывает Channel. В закрытом Channel, put!() завершится с ошибкой, в то время как take!() и fetch() успешно возвращают любые существующие значения до момента опустошения.

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

Удаленные ссылки и AbstractChannels

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

Для реализации AbstractChannel (например, Channel) необходимо реализовать put!(), take!(), fetch(), isready() и wait(). Удаленный объект, на который указывает 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, который использует словарь в качестве удалённого хранилища.

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

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

Узел, где хранится значение, отслеживает, какие из рабочих процессов имеют ссылку на него. Каждый раз, когда 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::Type, dims::NTuple; init=false, pids=Int[])

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

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

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

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

julia> S = SharedArray(Int, (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, (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]. В таких случаях лучше разбить массив вручную. Разделим по второму измерению:

# This function retuns the (irange,jrange) indexes assigned to this worker
@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

# Here's the kernel
@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

# Here's a convenience wrapper for a SharedArray implementation
@everywhere advection_shared_chunk!(q, u) = advection_chunk!(q, u, myrange(q)..., 1:size(q,3)-1)

Теперь сравним три различных версии: одна, выполняющаяся в одном процессе:

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

одна, использующая @parallel:

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

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

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

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

q = SharedArray(Float64, (500,500,500))
u = SharedArray(Float64, (500,500,500))

# Run once to JIT-compile
advection_serial!(q, u)
advection_parallel!(q, u)
advection_shared!(q,u)

# Now the real results:
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, ответственный за запуск рабочих процессов на одном хосте:

immutable 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

type 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 all-to-all на пользовательский транспортный уровень немного сложнее. Каждый процесс Julia имеет столько же задач связи, сколько рабочих, с которыми он соединён. Например, рассмотрим кластер Julia из 32 процессов в сети all-to-all:

  • Таким образом, каждый процесс 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 всё ещё логически соединены друг с другом — любой рабочий может напрямую отправлять сообщения любому другому рабочему, не зная о том, что ZeroMQ используется в качестве транспортного уровня.

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

  • Рабочие 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-подключения и для мастер-рабочего. Типичный сценарий для этого — локальный ноутбук, на котором запущен Julia REPL (т.е. мастер), а остальной кластер находится в облаке, например, на Amazon EC2. В этом случае необходимо открыть только порт 22 в удалённом кластере в сочетании с аутентификацией SSH-клиента через инфраструктуру открытых ключей (PKI). Учётные данные для аутентификации можно указать через sshflags, например, sshflags=`-e <keyfile>`.

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

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

Куки кластера

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

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

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

Указание топологии сети (Экспериментально)

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

  • :all_to_all : это значение по умолчанию, когда все рабочие подключены друг к другу.
  • :master_slave : только процесс-драйвер, т.е. pid 1, имеет подключения к рабочим.
  • :custom : метод launch управляющего кластером определяет топологию подключения. Поля ident и connect_idents в WorkerConfig используются для указания того же. connect_idents — это список идентификаторов ClusterManager, предоставленных рабочим, к которым рабочий, идентифицированный по ident должен подключиться.

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

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

В дополнение к задачам, удалённым вызовам и удалённым ссылкам, 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 на этих платформах, вы должны использовать ключевое слово 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, что этот цикл является многопоточным регионом.

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() (аналог yield() в C).

Обратите внимание, что, хотя код 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://www.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.5/manual/parallel-computing/

Spec-Zone.ru

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