Многопроцессорная обработка и распределенные вычисления
Реализация параллельных вычислений с распределенной памятью предоставляется модулем Distributed в составе стандартной библиотеки Julia.
Большинство современных компьютеров имеют более одного процессора, а несколько компьютеров могут быть объединены в кластер. Использование мощности этих нескольких процессоров позволяет быстрее выполнить многие вычисления. На производительность влияют два основных фактора: скорость самих процессоров и скорость их доступа к памяти. В кластере очевидно, что данный процессор будет иметь самый быстрый доступ к ОЗУ на том же компьютере (узле). Возможно, более неожиданно, что аналогичные проблемы актуальны на типичных многоядерных ноутбуках из-за разницы в скорости оперативной памяти и кэша. Следовательно, хорошая многопроцессорная среда должна позволять управлять «владением» блоком памяти конкретным процессором. Julia предоставляет многопроцессорную среду, основанную на передаче сообщений, позволяющую программам выполняться на нескольких процессах в разных областях памяти одновременно.
Реализация передачи сообщений в Julia отличается от других сред, таких как MPI[1]. Обмен сообщениями в Julia обычно «односторонний», что означает, что программист должен явно управлять только одним процессом в операции с двумя процессами. Кроме того, эти операции обычно не выглядят как «отправка сообщения» и «прием сообщения», а скорее напоминают операции более высокого уровня, такие как вызовы пользовательских функций.
Распределенное программирование в Julia основано на двух примитивах: удаленных ссылках и удаленных вызовах. Удаленная ссылка — это объект, который может использоваться любым процессом для ссылки на объект, хранящийся на конкретном процессе. Удаленный вызов — это запрос одного процесса для вызова определенной функции с определенными аргументами на другом (возможно, на том же) процессе.
Удаленные ссылки бывают двух типов: Future и RemoteChannel.
Удаленный вызов возвращает Future в качестве результата. Удаленные вызовы возвращаются немедленно; процесс, который сделал вызов, переходит к следующей операции, в то время как удаленный вызов происходит где-то в другом месте. Вы можете дождаться завершения удаленного вызова, вызвав wait на возвращенном Future, и вы можете получить полное значение результата, используя fetch.
С другой стороны, RemoteChannel могут быть перезаписаны. Например, несколько процессов могут координировать свою обработку, ссылаясь на один и тот же удаленный Channel.
Каждый процесс имеет ассоциированный идентификатор. Процесс, предоставляющий интерактивный приглашение Julia, всегда имеет id, равный 1. Процессы, используемые по умолчанию для параллельных операций, называются «рабочими». Когда есть только один процесс, процесс 1 считается рабочим. В противном случае рабочие считаются всеми процессами, кроме процесса 1. В результате, для получения преимуществ от методов параллельной обработки, таких как pmap, требуется добавление 2 или более процессов. Добавление одного процесса полезно, если вы хотите выполнять другие задачи в основном процессе, пока длительное вычисление выполняется на рабочем процессе.
Попробуем это. Начав с julia -p n предоставляется n рабочих процессов на локальном компьютере. Обычно имеет смысл, чтобы n равнялось числу потоков процессора (логических ядер) на компьютере. Обратите внимание, что аргумент -p неявно загружает модуль Distributed.
$ ./julia -p 2
julia> r = remotecall(rand, 2, 2, 2)
Future(2, 1, 4, nothing)
julia> s = @spawnat 2 1 .+ fetch(r)
Future(2, 1, 5, nothing)
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.
Для упрощения, символ :any может передаваться в @spawnat, который выбирает, где выполнить операцию за вас:
julia> r = @spawnat :any rand(2,2)
Future(2, 1, 4, nothing)
julia> s = @spawnat :any 1 .+ fetch(r)
Future(3, 1, 5, nothing)
julia> fetch(s)
2×2 Array{Float64,2}:
1.38854 1.9098
1.20939 1.57158
Обратите внимание, что мы использовали 1 .+ fetch(r) вместо 1 .+ r. Это связано с тем, что мы не знаем, где будет выполняться код, поэтому, как правило, может потребоваться fetch для перемещения r на процесс, выполняющий сложение. В этом случае @spawnat достаточно умна, чтобы выполнить вычисление на процессе, который владеет r, поэтому вызов fetch будет бездействующим (никакой работы не выполняется).
(Стоит отметить, что @spawnat не является встроенной, а определяется в Julia как макрос. Можно определить собственные такие конструкции.)
Важно помнить, что после получения Future кэширует свое значение локально. Дальнейшие вызовы fetch не подразумевают сетевой обмен. После того, как все ссылающиеся на Future получат данные, удалённое хранимое значение удаляется.
@async похож на @spawnat, но выполняет задачи только на локальном процессе. Мы используем его для создания задачи «подачи» для каждого процесса. Каждая задача выбирает следующий индекс, который нужно вычислить, затем ожидает завершения своего процесса, а затем повторяет до тех пор, пока не закончатся индексы. Обратите внимание, что задачи «подачи» не начинаются до тех пор, пока основная задача не достигнет конца блока @sync, в этот момент она передает управление и ожидает завершения всех локальных задач перед возвратом из функции. Что касается версий 0.7 и выше, задачи «подачи» могут совместно использовать состояние через nextidx, так как они все выполняются на одном процессе. Даже если Tasks планируются кооперативно, блокировка может по-прежнему потребоваться в некоторых контекстах, как в асинхронном вводе-выводе. Это означает, что переключения контекста происходят только в определенные моменты: в этом случае, когда вызывается remotecall_fetch. Это текущее состояние реализации, и оно может измениться в будущих версиях Julia, поскольку это призвано сделать возможным запуск до N Tasks на M Process, известный как многопоточность M:N. Затем потребуется модель приобретения\освобождения блокировки для nextidx , так как небезопасно, чтобы несколько процессов одновременно читают\записывали ресурс.
Доступность кода и загрузка пакетов
Ваш код должен быть доступен на любом процессе, который его выполняет. Например, введите следующее в приглашение 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(@spawnat :any rand2(2,2))
ERROR: RemoteException(2, CapturedException(UndefVarError(Symbol("#rand2"))
Stacktrace:
[...]
Процесс 1 знал о функции rand2, но процесс 2 — нет.
Чаще всего вы будете загружать код из файлов или пакетов, и у вас есть значительная гибкость в управлении тем, какие процессы загружают код. Рассмотрим файл DummyModule.jl, содержащий следующий код:
module DummyModule
export MyType, f
mutable struct MyType
a::Int
end
f(x) = x^2+1
println("loaded")
end
Для ссылки на MyType во всех процессах, DummyModule.jl необходимо загрузить на каждый процесс. Вызов include("DummyModule.jl") загружает его только на один процесс. Чтобы загрузить его на каждый процесс, используйте макрос @everywhere (начало Julia с julia -p 2):
julia> @everywhere include("DummyModule.jl")
loaded
From worker 3: loaded
From worker 2: loaded
Как обычно, это не включает DummyModule в область видимости ни на одном из процессов, что требует using или import. Кроме того, когда DummyModule включено в область видимости одного процесса, оно не находится в области видимости других:
julia> using .DummyModule julia> MyType(7) MyType(7) julia> fetch(@spawnat 2 MyType(7)) ERROR: On worker 2: UndefVarError: MyType not defined ⋮ julia> fetch(@spawnat 2 DummyModule.MyType(7)) MyType(7)
Однако по-прежнему возможно, например, отправить MyType на процесс, который загрузил DummyModule, даже если он не находится в области видимости:
julia> put!(RemoteChannel(2), MyType(7))
RemoteChannel{Channel{Any}}(2, 1, 13)
Файл также может быть предварительно загружен на несколько процессов при запуске с флагом -L, а сценарий драйвера может использоваться для управления вычислением:
julia -p <n> -L file1.jl -L file2.jl driver.jl
Процесс Julia, выполняющий сценарий драйвера в примере выше, имеет id равный 1, как и процесс, предоставляющий интерактивное приглашение.
Наконец, если DummyModule.jl не является отдельным файлом, а пакетом, то using DummyModule загрузит DummyModule.jl на всех процессах, но включит его в область видимости только на том процессе, где using был вызван.
Запуск и управление рабочими процессами
Базовая установка Julia имеет встроенную поддержку двух типов кластеров:
- Локальный кластер, указанный с помощью опции
-p, как показано выше. - Кластер, охватывающий машины, использующий опцию
--machine-file. Это использует бесклеточнуюsshавторизацию для запуска процессов Julia worker (из той же директории, что и текущий хост) на указанных машинах. Описание каждой машины имеет вид[count*][user@]host[:port] [bind_addr[:port]].userпо умолчанию соответствует текущему пользователю,port— стандартному порту ssh.count— количество рабочих процессов, запускаемых на узле, по умолчанию равно 1. Необязательная опцияbind-to bind_addr[:port]указывает IP-адрес и порт, которые другие рабочие процессы должны использовать для подключения к этому рабочему процессу.
Функции addprocs, rmprocs, workers и другие доступны как программистский способ добавления, удаления и запроса процессов в кластере.
julia> using Distributed
julia> addprocs(2)
2-element Array{Int64,1}:
2
3
Модуль Distributed должен быть явно загружен на главном процессе перед вызовом addprocs. Он автоматически становится доступным на рабочих процессах.
Обратите внимание, что рабочие процессы не выполняют скрипт запуска ~/.julia/config/startup.jl, а также не синхронизируют свое глобальное состояние (такое как глобальные переменные, новые определения методов и загруженные модули) ни с одним из других работающих процессов. Вы можете использовать addprocs(exeflags="--project") для инициализации рабочего процесса с определенной средой, а затем @everywhere using <modulename> или @everywhere include("file.jl").
Другие типы кластеров могут быть поддерживаемы написанием собственных пользовательских ClusterManager, как описано ниже в разделе ClusterManagers.
Перемещение данных
Отправка сообщений и перемещение данных составляют большую часть накладных расходов в распределенной программе. Снижение количества сообщений и объема передаваемых данных имеет решающее значение для достижения производительности и масштабируемости. Для этого важно понимать перемещение данных, выполняемое различными конструкциями распределенного программирования Julia.
fetch можно рассматривать как явную операцию перемещения данных, поскольку она напрямую запрашивает перемещение объекта на локальную машину. @spawnat (и несколько связанных конструкций) также перемещает данные, но это не так очевидно, поэтому его можно назвать неявной операцией перемещения данных. Рассмотрим эти два подхода к построению и возведению в квадрат случайной матрицы:
Способ 1:
julia> A = rand(1000,1000); julia> Bref = @spawnat :any A^2; [...] julia> fetch(Bref);
Способ 2:
julia> Bref = @spawnat :any rand(1000,1000)^2; [...] julia> fetch(Bref);
Разница кажется незначительной, но на самом деле она довольно существенна из-за поведения @spawnat. В первом методе случайная матрица строится локально, затем отправляется на другой процесс, где она возводится в квадрат. Во втором методе случайная матрица как строится, так и возводится в квадрат на другом процессе. Поэтому во втором методе отправляется намного меньше данных, чем в первом.
В этом примере с игрушечной задачей два метода легко различить и выбрать. Однако в реальной программе проектирование перемещения данных может потребовать большего внимания и, вероятно, некоторых измерений. Например, если первому процессу нужна матрица A, то первый метод может быть лучше. Или, если вычисление A дорогостоящее и только текущий процесс его имеет, то перемещение его на другой процесс может быть неизбежным. Или, если у текущего процесса очень мало работы между @spawnat и fetch(Bref), то может быть лучше полностью исключить распараллеливание. Или представьте, что rand(1000,1000) заменяется более дорогостоящей операцией. Тогда может иметь смысл добавить еще одну инструкцию @spawnat только для этого шага.
Глобальные переменные
Выражения, выполняемые удаленно через @spawnat, или замыкания, указанные для удаленного выполнения с помощью remotecall, могут ссылаться на глобальные переменные. Глобальные связи в модуле Main обрабатываются немного по-другому по сравнению с глобальными связями в других модулях. Рассмотрим следующий фрагмент кода:
A = rand(10,10) remotecall_fetch(()->sum(A), 2)
В этом случае sum ДОЛЖЕН быть определен в удаленном процессе. Обратите внимание, что A — это глобальная переменная, определенная в локальном рабочем пространстве. У рабочего процесса 2 нет переменной с именем A в Main. Акт отправки замыкания ()->sum(A) рабочему процессу 2 приводит к определению Main.A на 2. Main.A продолжает существовать на рабочем процессе 2 даже после того, как вызов remotecall_fetch возвращается. Удаленные вызовы с вложенными ссылками на глобальные переменные (только в модуле Main) управляют глобальными переменными следующим образом:
Новые глобальные связи создаются на целевых рабочих процессах, если на них есть ссылки в рамках удаленного вызова.
Глобальные константы объявляются как константы и на удаленных узлах.
-
Глобальные переменные пересылаются на целевой рабочий процесс только в контексте удаленного вызова и только в том случае, если их значение изменилось. Кроме того, кластер не синхронизирует глобальные связи между узлами. Например:
A = rand(10,10) remotecall_fetch(()->sum(A), 2) # worker 2 A = rand(10,10) remotecall_fetch(()->sum(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> @fetchfrom 2 InteractiveUtils.varinfo()
name size summary
––––––––– ––––––––– ––––––––––––––––––––––
A 800 bytes 10×10 Array{Float64,2}
Base Module
Core Module
Main Module
Как видно, глобальная переменная A определена на рабочем процессе 2, но B захватывается как локальная переменная, и поэтому привязка для B на рабочем процессе 2 отсутствует.
Параллельные отображения и циклы
К счастью, многие полезные параллельные вычисления не требуют перемещения данных. Типичным примером является моделирование Монте-Карло, где несколько процессов могут обрабатывать независимые испытания симуляции одновременно. Мы можем использовать @spawnat для подбрасывания монет на двух процессах. Сначала напишите следующую функцию в 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_string(Main, $(read("count_heads.jl", String)), "count_heads.jl")
julia> a = @spawnat :any count_heads(100000000)
Future(2, 1, 6, nothing)
julia> b = @spawnat :any count_heads(100000000)
Future(3, 1, 7, nothing)
julia> fetch(a)+fetch(b)
100001564
Этот пример демонстрирует мощный и часто используемый шаблон параллельного программирования. Много итераций выполняются независимо на нескольких процессах, а затем их результаты объединяются с помощью некоторой функции. Процесс объединения называется сведением, так как он обычно сводит тензорный ранг: вектор чисел сводится к одному числу, матрица — к одной строке или столбцу и т. д. В коде это обычно выглядит как шаблон x = f(x,v[i]), где x — накопитель, f — функция сведения, а v[i] — элементы, которые сводятся. Желательно, чтобы f была ассоциативной, чтобы порядок выполнения операций не влиял на результат.
Заметьте, что наше использование этого шаблона с count_heads может быть обобщено. Мы использовали две явные инструкции @spawnat, что ограничивает распараллеливание двумя процессами. Чтобы запустить на любом количестве процессов, мы можем использовать параллельный цикл for, работающий в распределенной памяти, который может быть написан в Julia с помощью @distributed следующим образом:
nheads = @distributed (+) for i = 1:200000000
Int(rand(Bool))
end
Эта конструкция реализует шаблон назначения итераций множеству процессов и объединения их с заданной операцией сведения (в данном случае (+)). Результат каждой итерации принимается как значение последнего выражения в цикле. Само выражение параллельного цикла оценивается в окончательный результат.
Обратите внимание, что хотя параллельные циклы for выглядят как последовательные циклы for, их поведение радикально отличается. В частности, порядок итераций не задан, и записи в переменные или массивы не будут глобально видимыми, так как итерации выполняются на разных процессах. Любые переменные, используемые внутри параллельного цикла, будут скопированы и распространены на каждый процесс.
Например, следующий код не будет работать как ожидается:
a = zeros(100000)
@distributed for i = 1:100000
a[i] = i
end
Этот код не инициализирует все a, так как каждый процесс будет иметь свою отдельную копию. Параллельные циклы for подобного вида следует избегать. К счастью, Общие массивы можно использовать для преодоления этого ограничения:
using SharedArrays
a = SharedArray{Float64}(10)
@distributed for i = 1:10
a[i] = i
end
Использование «внешних» переменных в параллельных циклах вполне допустимо, если эти переменные являются только для чтения:
a = randn(1000)
@distributed (+) for i = 1:100000
f(a[rand(1:end)])
end
Здесь каждая итерация применяет f к случайному образцу из вектора a, общий для всех процессов.
Как вы могли видеть, оператор сведения может быть опущен, если он не нужен. В этом случае цикл выполняется асинхронно, т. е. он запускает независимые задачи на всех доступных рабочих процессах и возвращает массив Future немедленно, не дожидаясь завершения. Вызывающая программа может дождаться завершения Future в более поздний момент, вызвав fetch на них или дождаться завершения в конце цикла, добавив перед ним @sync, как в @sync @distributed for.
В некоторых случаях оператор сокращения не требуется, и нам нужно просто применить функцию ко всем целым числам в определенном диапазоне (или, более обще, ко всем элементам в некотором наборе). Это еще одна полезная операция, называемая параллельным отображением, реализованная в Julia как функция pmap. Например, мы можем вычислить сингулярные значения нескольких больших случайных матриц параллельно следующим образом:
julia> M = Matrix{Float64}[rand(1000,1000) for i = 1:10];
julia> pmap(svdvals, M);
Функция Julia pmap предназначена для случая, когда каждый вызов функции выполняет большой объем работы. В отличие от этого, @distributed for может обрабатывать ситуации, когда каждый итерационный шаг очень мал, например, простое суммирование двух чисел. Только рабочие процессы используются как для pmap, так и для @distributed for для параллельного вычисления. В случае @distributed for, окончательное сокращение выполняется на вызывающем процессе.
Удаленные ссылки и абстрактные каналы
Удаленные ссылки всегда ссылаются на реализацию 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, реализованные пользователем. Простой пример этого представлен в dictchannel.jl в репозитории Примеры, который использует словарь в качестве удаленного хранилища.
Каналы и RemoteChannels
- Канал
Channelлокален для процесса. Рабочий процесс 2 не может напрямую ссылаться наChannelна рабочем процессе 3 и наоборот.RemoteChannel, однако, может помещать и извлекать значения между рабочими процессами. RemoteChannelможно рассматривать как обработчик дляChannel.- Идентификатор процесса,
pid, связанный сRemoteChannel, идентифицирует процесс, где хранилище, т.е. базовыйChannel, существует. - Любой процесс со ссылкой на
RemoteChannelможет помещать и извлекать элементы из канала. Данные автоматически отправляются (или извлекаются) в процесс, с которым связанRemoteChannel. - Сериализация
Channelтакже сериализует все данные, присутствующие в канале. Десериализация, следовательно, фактически создает копию исходного объекта. - С другой стороны, сериализация
RemoteChannelвключает только сериализацию идентификатора, который идентифицирует местоположение и экземплярChannel, на который указывает обработчик. ОбъектRemoteChannel, десериализованный (на любом рабочем процессе), следовательно, также указывает на то же хранилище, что и исходный.
Пример каналов из вышеприведенного текста может быть изменен для межпроцессорного взаимодействия, как показано ниже.
Мы запускаем 4 рабочих процесса для обработки одного jobs удаленного канала. Задания, идентифицированные по ID (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> @async 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
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; digits=2)) seconds on worker $where")
global 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 собирается мусором локально, узел, владеющий значением, снова уведомляется. Это реализовано в внутреннем кластерно-ориентированном сериализаторе. Удаленные ссылки действительны только в контексте работающего кластера. Сериализация и десериализация ссылок на обычные IO объекты не поддерживаются.
Уведомления выполняются путем отправки сообщений "отслеживания" — сообщение "добавить ссылку" при сериализации ссылки на другой процесс и сообщение "удалить ссылку" при локальной сборке мусора ссылки.
Поскольку Future имеют свойство записи один раз и кэшируются локально, действие fetch для Future также обновляет информацию о слежении за ссылками на узле, владеющем значением.
Узел, который владеет значением, освобождает его, когда все ссылки на него будут удалены.
В случае с Future, сериализация уже извлеченного Future на другой узел также отправляет значение, так как исходное удаленное хранилище может собрать значение к этому моменту.
Важно отметить, что когда объект собирается мусором локально, зависит от размера объекта и текущей потребности памяти в системе.
В случае с удаленными ссылками размер локального объекта ссылки довольно мал, в то время как значение, хранящееся на удаленном узле, может быть довольно большим. Поскольку локальный объект может не быть собран немедленно, рекомендуется явно вызывать finalize для локальных экземпляров RemoteChannel или не извлеченных Future. Поскольку вызов fetch для Future также удаляет его ссылку из удаленного хранилища, это не требуется для извлеченных Future. Явное вызов finalize приводит к немедленному отправлению сообщения на удаленный узел, чтобы тот удалил свою ссылку на значение.
После финализации ссылка становится недействительной и не может использоваться в дальнейших вызовах.
Локальные вызовы
Данные неизбежно копируются на удаленный узел для выполнения. Это относится как к удаленным вызовам, так и к случаям, когда данные хранятся в RemoteChannel / Future на другом узле. Как ожидалось, это приводит к созданию копии сериализованных объектов на удаленном узле. Однако, когда целевой узел является локальным узлом, т.е. идентификатор вызывающего процесса совпадает с идентификатором удаленного узла, он выполняется как локальный вызов. Обычно (но не всегда) он выполняется в другой задаче — но нет сериализации/десериализации данных. Следовательно, вызов ссылается на те же экземпляры объектов, что и переданные — копии не создаются. Это поведение показано ниже:
julia> using Distributed;
julia> rc = RemoteChannel(()->Channel(3)); # RemoteChannel created on local node
julia> v = [0];
julia> for i in 1:3
v[1] = i # Reusing `v`
put!(rc, v)
end;
julia> result = [take!(rc) for _ in 1:3];
julia> println(result);
Array{Int64,1}[[3], [3], [3]]
julia> println("Num Unique objects : ", length(unique(map(objectid, result))));
Num Unique objects : 1
julia> addprocs(1);
julia> rc = RemoteChannel(()->Channel(3), workers()[1]); # RemoteChannel created on remote node
julia> v = [0];
julia> for i in 1:3
v[1] = i
put!(rc, v)
end;
julia> result = [take!(rc) for _ in 1:3];
julia> println(result);
Array{Int64,1}[[1], [2], [3]]
julia> println("Num Unique objects : ", length(unique(map(objectid, result))));
Num Unique objects : 3
Как видно, put! для локального RemoteChannel с тем же объектом v модифицированным между вызовами приводит к тому же единственному экземпляру объекта, хранящегося. В противоположность этому, при создании копий v при вызове на узел, владеющем rc, на другом узле.
Следует отметить, что это, как правило, не проблема. Это нужно учитывать только в том случае, если объект хранится локально и изменяется после вызова. В таких случаях может быть целесообразно хранить deepcopy объекта.
Это также верно для удаленных вызовов на локальном узле, как показано в следующем примере:
julia> using Distributed; addprocs(1);
julia> v = [0];
julia> v2 = remotecall_fetch(x->(x[1] = 1; x), myid(), v); # Executed on local node
julia> println("v=$v, v2=$v2, ", v === v2);
v=[1], v2=[1], true
julia> v = [0];
julia> v2 = remotecall_fetch(x->(x[1] = 1; x), workers()[1], v); # Executed on remote node
julia> println("v=$v, v2=$v2, ", v === v2);
v=[0], v2=[1], false
Как можно видеть еще раз, удаленный вызов на локальный узел ведет себя так же, как прямой вызов. Вызов изменяет локальные объекты, переданные в качестве аргументов. В удаленном вызове он работает с копией аргументов.
Повторим, что это, как правило, не проблема. Если локальный узел также используется в качестве вычислительного узла, и аргументы используются после вызова, это поведение необходимо учитывать, и при необходимости должны передаваться глубокие копии аргументов для вызова на локальном узле. Вызовы удаленных узлов всегда будут работать с копиями аргументов.
Общие массивы
Общие массивы используют общую системную память для отображения одного и того же массива во многих процессах. Хотя есть некоторые сходства с DArray, поведение SharedArray отличается. В DArray каждый процесс имеет локальный доступ только к части данных, и никакие два процесса не делят одну и ту же часть; в отличие от этого, в SharedArray каждый «участвующий» процесс имеет доступ ко всему массиву. SharedArray — хороший выбор, когда вам нужно, чтобы большой объем данных был совместно доступен двум или более процессам на одном компьютере.
Поддержка общих массивов доступна через модуль SharedArrays, который должен быть явно загружен на всех участвующих рабочих узлах.
SharedArray индексирование (присвоение и доступ к значениям) работает так же, как и с обычными массивами, и является эффективным, потому что базовая память доступна локальному процессу. Поэтому большинство алгоритмов естественным образом работают с SharedArray, хотя и в режиме работы с одним процессом. В тех случаях, когда алгоритм настаивает на вводе Array, базовый массив можно извлечь из SharedArray путем вызова sdata. Для других AbstractArray типов, sdata просто возвращает сам объект, поэтому использовать sdata на любом Array-типе объекта безопасно.
Конструктор для общего массива имеет вид:
SharedArray{T,N}(dims::NTuple; init=false, pids=Int[])
что создает N-мерный общий массив типа bits T и размера dims по всем процессам, указанным в pids. В отличие от распределенных массивов, общий массив доступен только из тех участвующих рабочих узлов, которые указаны параметром pids, (и создающий процесс тоже, если он находится на том же хосте). Обратите внимание, что в SharedArray поддерживаются только элементы, являющиеся isbits.
Если функция init, с сигнатурой initfn(S::SharedArray), указана, то она вызывается на всех участвующих рабочих узлах. Вы можете указать, чтобы каждый рабочий узел выполнял функцию init на отдельной части массива, тем самым обеспечивая параллельную инициализацию.
Вот простой пример:
julia> using Distributed
julia> addprocs(3)
3-element Array{Int64,1}:
2
3
4
julia> @everywhere using SharedArrays
julia> S = SharedArray{Int,2}((3,4), init = S -> S[localindices(S)] = repeat([myid()], length(localindices(S))))
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
SharedArrays.localindices предоставляет несвязные одномерные диапазоны индексов и иногда удобен для распределения задач между процессами. Конечно, вы можете разделить работу любым желаемым способом:
julia> S = SharedArray{Int,2}((3,4), init = S -> S[indexpids(S):length(procs(S)):length(S)] = repeat([myid()], length( indexpids(S):length(procs(S)):length(S))))
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 range(0, stop=size(q,2), length=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);
одна, использующая @distributed:
julia> function advection_parallel!(q, u)
for t = 1:size(q,3)-1
@sync @distributed 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! заключается в том, что оно минимизирует трафик между рабочими узлами, позволяя каждому из них вычислять более длительное время на назначенном фрагменте.
Общие массивы и распределенное управление сборкой мусора
Как и удаленные ссылки, общие массивы также зависят от управления сборкой мусора на узле создания, чтобы освободить ссылки от всех участвующих рабочих узлов. Код, который создает много кратковременных объектов общих массивов, извлечет выгоду из явного завершения этих объектов как можно скорее. Это приводит к более быстрому освобождению и памяти, и файловых дескрипторов, отображающих общий сегмент.
Управляющие программы кластеров
Запуск, управление и сетевое взаимодействие процессов 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-сокеты.
Все рабочие процессы в кластере используют тот же куки, что и мастер. Когда куки не указан, т. е. с параметром --worker, рабочий процесс пытается прочитать его из стандартного ввода. LocalManager и SSHManager оба передают куки новым запущенным рабочим процессам через их стандартные вводы.
По умолчанию рабочий процесс будет прослушивать свободный порт по адресу, возвращаемому вызовом getipaddr(). Конкретный адрес для прослушивания можно указать с помощью необязательного аргумента --bind-to bind_addr[:port]. Это полезно для многоадресных хостов.
В качестве примера транспортного средства, отличного от TCP/IP, реализация может выбрать использование MPI, в этом случае --worker не должно быть указано. Вместо этого новые запущенные рабочие процессы должны вызвать init_worker(cookie) перед использованием каких-либо параллельных конструкций.
Для каждого запущенного рабочего процесса метод launch должен добавить объект WorkerConfig (с соответствующими инициализированными полями) в launched
mutable struct WorkerConfig
# Common fields relevant to all cluster managers
io::Union{IO, Nothing}
host::Union{AbstractString, Nothing}
port::Union{Integer, Nothing}
# Used when launching additional workers at a host
count::Union{Int, Symbol, Nothing}
exename::Union{AbstractString, Cmd, Nothing}
exeflags::Union{Cmd, Nothing}
# 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::Any
# SSHManager / SSH tunnel connections to workers
tunnel::Union{Bool, Nothing}
bind_addr::Union{AbstractString, Nothing}
sshflags::Union{Cmd, Nothing}
max_parallel::Union{Integer, Nothing}
# Used by Local/SSH managers
connect_at::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 — это поток, который может обрабатываться асинхронно.
Папка 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() для указанного удалённого рабочего процесса.
Папка примеров clustermanager/simple — это пример, демонстрирующий простую реализацию с использованием сокетов UNIX-домена для настройки кластера.
Требования к сети для LocalManager и SSHManager
Кластеры Julia предназначены для выполнения в уже защищённых средах на такой инфраструктуре, как локальные ноутбуки, ведомственные кластеры или даже облака. В этом разделе рассматриваются требования к сетевой безопасности для встроенных LocalManager и SSHManager:
Главный процесс не прослушивает какой-либо порт. Он только подключается к рабочим процессам.
Каждый рабочий процесс привязывается только к одному из локальных интерфейсов и прослушивает порт с временным номером, назначенным ОС.
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=`-i <keyfile>`.В топологии «все со всеми» (по умолчанию) все рабочие процессы подключаются друг к другу через обычные TCP-сокеты. Таким образом, политика безопасности на узлах кластера должна гарантировать свободное подключение между рабочими процессами для диапазона портов с временным номером (зависит от ОС).
Защита и шифрование всего трафика между рабочими процессами (через SSH) или шифрование отдельных сообщений можно выполнить с помощью пользовательской
ClusterManager. Если вы укажете
multiplex=trueкак опцию дляaddprocs, используется SSH-мьюльтиплексирование для создания туннеля между мастером и рабочими процессами. Если вы настроили SSH-мьюльтиплексирование самостоятельно и соединение уже установлено, SSH-мьюльтиплексирование используется независимо от опцииmultiplex. Если мьюльтиплексирование включено, переадресация настраивается с использованием существующего соединения (опция-O forwardв ssh). Это полезно, если вашим серверам требуется аутентификация по паролю; вы можете избежать аутентификации в Julia, войдя на сервер доaddprocs. Сокет управления будет расположен по адресу~/.ssh/julia-%r@%h:%pво время сессии, если не используется существующее мьюльтиплексированное соединение. Обратите внимание, что пропускная способность может быть ограничена, если вы создадите несколько процессов на узле и включите мьюльтиплексирование, потому что в этом случае процессы используют одно TCP-соединение мьюльтиплексирования.
Куки кластера
Все процессы в кластере используют один и тот же куки, который по умолчанию является случайной строкой в главном процессе:
-
cluster_cookie()возвращает куки, в то время какcluster_cookie(cookie)()устанавливает его и возвращает новый куки. - Все подключения аутентифицируются с обеих сторон, чтобы гарантировать, что только рабочие процессы, запущенные мастером, могут подключаться друг к другу.
- Куки можно передать рабочим процессам при запуске через аргумент
--worker=<cookie>. Если аргумент--workerуказан без куки, рабочий процесс пытается прочитать куки из стандартного ввода (stdin).stdinзакрывается сразу после извлечения куки. -
Управляющие программы
ClusterManagerмогут извлечь куки из мастера, вызвавcluster_cookie(). Управляющие программы кластеров, не использующие стандартный транспорт TCP/IP (и, следовательно, не указывающие--worker), должны вызватьinit_worker(cookie, manager)с тем же куки, что и в главном процессе.
Обратите внимание, что среды, требующие более высокого уровня безопасности, могут реализовать это с помощью пользовательской ClusterManager. Например, куки могут быть предварительно объявлены и, следовательно, не указываться в качестве аргумента запуска.
Указание топологии сети (экспериментальная функция)
Ключевой аргумент topology, переданный в addprocs используется для указания того, как рабочие процессы должны быть связаны друг с другом:
-
:all_to_all, по умолчанию: все рабочие процессы подключены друг к другу. -
:master_worker: только процесс драйвера, т.е.pid1, имеет подключения к рабочим процессам. -
:custom: методlaunchуправляющей программы кластера задает топологию подключения через поляidentиconnect_identsвWorkerConfig. Рабочий процесс с предоставленным управляющей программой кластера идентификаторомidentподключится ко всем рабочим процессам, указанным вconnect_idents.
Ключевой аргумент lazy=true|false влияет только на опцию topology :all_to_all. Если true, кластер начинает работу с мастером, подключенным ко всем рабочим процессам. Установление конкретных подключений между рабочими процессами происходит при первом удаленном вызове между двумя рабочими процессами. Это помогает уменьшить начальные ресурсы, выделенные для внутрикластерной коммуникации. Подключения устанавливаются в зависимости от временных требований параллельной программы. Значение по умолчанию для lazy — true.
В настоящее время отправка сообщения между неподключенными рабочими процессами приводит к ошибке. Это поведение, как и функциональность и интерфейс, следует считать экспериментальным и может измениться в будущих версиях.
Заслуживающие внимания внешние пакеты
Вне параллелизма Julia существует множество внешних пакетов, которые стоит упомянуть. Например, MPI.jl — это обёртка для Julia для протокола MPI или DistributedArrays.jl, как показано в Общие массивы. Следует упомянуть и экосистему Julia для программирования на GPU, которая включает:
Операции на низком уровне (ядра на C) OpenCL.jl и CUDAdrv.jl, представляющие соответственно интерфейс OpenCL и обёртку CUDA.
Интерфейсы на низком уровне (ядро Julia) такие как CUDAnative.jl, являющийся реализацией CUDA на основе Julia.
Абстракции высокого уровня, специфичные для поставщиков, такие как CuArrays.jl и CLArrays.jl
Библиотеки высокого уровня, такие как ArrayFire.jl и GPUArrays.jl
В следующем примере мы будем использовать как DistributedArrays.jl, так и CuArrays.jl, чтобы распределить массив по нескольким процессам, сначала преобразовав его через distribute() и CuArray().
Помните, при импорте DistributedArrays.jl импортируйте его во всех процессах, используя @everywhere
$ ./julia -p 4
julia> addprocs()
julia> @everywhere using DistributedArrays
julia> using CuArrays
julia> B = ones(10_000) ./ 2;
julia> A = ones(10_000) .* π;
julia> C = 2 .* A ./ B;
julia> all(C .≈ 4*π)
true
julia> typeof(C)
Array{Float64,1}
julia> dB = distribute(B);
julia> dA = distribute(A);
julia> dC = 2 .* dA ./ dB;
julia> all(dC .≈ 4*π)
true
julia> typeof(dC)
DistributedArrays.DArray{Float64,1,Array{Float64,1}}
julia> cuB = CuArray(B);
julia> cuA = CuArray(A);
julia> cuC = 2 .* cuA ./ cuB;
julia> all(cuC .≈ 4*π);
true
julia> typeof(cuC)
CuArray{Float64,1}
Помните, что некоторые функции Julia в настоящее время не поддерживаются CUDAnative.jl[2], особенно некоторые функции, такие как sin, будут нуждаться в замене на CUDAnative.sin(cc: @maleadt).
В следующем примере мы будем использовать как DistributedArrays.jl, так и CuArrays.jl, чтобы распределить массив по нескольким процессам и вызвать на нём общую функцию.
function power_method(M, v)
for i in 1:100
v = M*v
v /= norm(v)
end
return v, norm(M*v) / norm(v) # or (M*v) ./ v
end
power_method многократно создаёт новый вектор и нормирует его. Мы не указали сигнатуру типа в объявлении функции, давайте посмотрим, будет ли она работать с вышеупомянутыми типами данных:
julia> M = [2. 1; 1 1];
julia> v = rand(2)
2-element Array{Float64,1}:
0.40395
0.445877
julia> power_method(M,v)
([0.850651, 0.525731], 2.618033988749895)
julia> cuM = CuArray(M);
julia> cuv = CuArray(v);
julia> curesult = power_method(cuM, cuv);
julia> typeof(curesult)
CuArray{Float64,1}
julia> dM = distribute(M);
julia> dv = distribute(v);
julia> dC = power_method(dM, dv);
julia> typeof(dC)
Tuple{DistributedArrays.DArray{Float64,1,Array{Float64,1}},Float64}
Для завершения краткого обзора внешних пакетов, рассмотрим MPI.jl, обёртку Julia для протокола MPI. Поскольку рассмотрение каждой внутренней функции займёт слишком много времени, лучше просто оценить подход, используемый для реализации протокола.
Рассмотрим этот тестовый скрипт, который просто вызывает каждый подпроцесс, инициализирует его ранг, а когда достигает главного процесса, суммирует ранги.
import MPI
MPI.Init()
comm = MPI.COMM_WORLD
MPI.Barrier(comm)
root = 0
r = MPI.Comm_rank(comm)
sr = MPI.Reduce(r, MPI.SUM, root, comm)
if(MPI.Comm_rank(comm) == root)
@printf("sum of ranks: %s\n", sr)
end
MPI.Finalize()
mpirun -np 4 ./julia example.jl
- 1В этом контексте MPI относится к стандарту MPI-1. Начиная с MPI-2, комитет стандартов MPI ввёл новый набор механизмов коммуникации, которые коллективно известны как удалённый доступ к памяти (RMA). Целью добавления RMA в стандарт MPI было облегчение односторонних схем коммуникации. Для получения дополнительной информации об последнем стандарте MPI, смотрите https://mpi-forum.org/docs.
- 2Документация по GPU Julia
© 2009–2021 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.6.0/manual/distributed-computing/