Многопроцессорность-и-распределённые-вычисления
Реализация параллельных вычислений с распределённой памятью предоставляется модулем 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 создать 2×2 случайную матрицу, а во второй строке — добавить к ней 1. Результат обоих вычислений доступен в двух фьючерсах, r и s. Макрос @spawnat вычисляет выражение во втором аргументе на процессе, указанном в первом аргументе.
Иногда вам может потребоваться немедленно получить удалённо вычисленное значение. Это обычно происходит, когда вы считываете удалённый объект, чтобы получить данные, необходимые для следующей локальной операции. Для этой цели существует функция remotecall_fetch. Она эквивалентна fetch(remotecall(...)) но более эффективна.
julia> remotecall_fetch(r-> fetch(r)[1, 1], 2, r) 0.18526337335308085
Это извлекает массив на рабочем процессе 2 и возвращает первое значение. Обратите внимание, что fetch в этом случае не перемещает данные, так как выполняется на рабочем процессе, который владеет массивом. Также можно записать:
julia> remotecall_fetch(getindex, 2, r, 1, 1) 0.10824216411304866
Помните, что 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— количество процессов worker, которые нужно запустить на узле, по умолчанию равно 1. Необязательный параметрbind-to bind_addr[:port]указывает IP-адрес и порт, который должны использовать другие worker для подключения к этому worker.
Функции addprocs, rmprocs, workers и другие доступны в качестве программного способа добавления, удаления и запроса процессов в кластере.
julia> using Distributed
julia> addprocs(2)
2-element Array{Int64,1}:
2
3
Модуль Distributed должен быть явно загружен на главном процессе перед вызовом addprocs. Он автоматически становится доступным на worker-процессах.
Обратите внимание, что worker не выполняют скрипт запуска ~/.julia/config/startup.jl, и не синхронизируют своё глобальное состояние (такое как глобальные переменные, определения новых методов и загруженные модули) с какими-либо другими работающими процессами. Вы можете использовать addprocs(exeflags="--project") для инициализации worker с определённой средой, а затем @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 — это глобальная переменная, определённая в локальном рабочем пространстве. Worker 2 не имеет переменной, называемой A в Main. Действие отправки замыкания ()->sum(A) worker 2 приводит к определению Main.A на 2. Main.A продолжает существовать на worker 2 даже после того, как вызов remotecall_fetch возвращается. Удалённые вызовы с вложенными ссылками на глобальные переменные (только в модуле Main ) управляют глобальными переменными следующим образом:
Новые глобальные связи создаются на целевых worker, если они упоминаются в рамках удалённого вызова.
Глобальные константы также объявляются как константы на удалённых узлах.
-
Глобальные переменные пересылаются на целевой worker только в контексте удалённого вызова, и только если их значение изменилось. Кроме того, кластер не синхронизирует глобальные связи между узлами. Например:
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на worker 2 имеет другое значение, чемMain.Aна worker 3, в то время как значениеMain.Aна узле 1 установлено вnothing.
Как вы могли заметить, в то время как память, связанная с глобальными переменными, может быть освобождена, когда они переопределяются на главном узле, такого действия не происходит на worker, поскольку связи остаются действительными. 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 определена на worker 2, но B захвачена как локальная переменная, и поэтому связи для B не существует на worker 2.
Параллельные Map и циклы
К счастью, многие полезные параллельные вычисления не требуют перемещения данных. Распространённый пример — симуляция Монте-Карло, где несколько процессов могут одновременно обрабатывать независимые испытания моделирования. Мы можем использовать @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, подобные этому, следует избегать. К счастью, Shared Arrays позволяют обойти это ограничение:
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, общий для всех процессов.
Как вы могли видеть, оператор редукции можно опустить, если он не нужен. В этом случае цикл выполняется асинхронно, т. е. он порождает независимые задачи на всех доступных worker и возвращает массив 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 в репозитории Примеры, который использует словарь в качестве удалённого хранилища.
Каналы и удаленные каналы
- Канал
Channelлокален для процесса. Рабочий процесс 2 не может напрямую обратиться кChannelна рабочем процессе 3 и наоборот. ОднакоRemoteChannelможет помещать и извлекать значения между рабочими процессами. - Можно рассматривать
RemoteChannelкак указатель наChannel. - Идентификатор процесса,
pid, связанный сRemoteChannel, идентифицирует процесс, где существует подчинённое хранилище, то есть подчинённыйChannel. - Любой процесс с ссылкой на
RemoteChannelможет помещать и извлекать элементы из канала. Данные автоматически отправляются (или извлекаются) в (или из) процесс, с которым связанRemoteChannel. - Сериализация
Channelтакже сериализует любые данные, присутствующие в канале. Десериализация, следовательно, создаёт копию исходного объекта. - С другой стороны, сериализация
RemoteChannelвключает только сериализацию идентификатора, который идентифицирует местоположение и экземплярChannel, на который указывает указатель. Десериализованный объектRemoteChannel(на любом рабочем процессе) поэтому также ссылается на то же подчинённое хранилище, что и оригинал.
Пример каналов из вышеприведённого описания можно модифицировать для межпроцессорного взаимодействия, как показано ниже.
Мы запускаем 4 рабочих процесса для обработки одного jobs удалённого канала. Задачи, идентифицированные по идентификатору (job_id), записываются в канал. Каждая выполняющаяся на удалённом процессе задача в этой симуляции считывает job_id, ждёт случайное количество времени и записывает в канал результатов кортеж job_id, времени выполнения и собственного pid в канал результатов. В конце концов, все results выводятся на главном процессе.
julia> addprocs(4); # add worker processes
julia> const jobs = RemoteChannel(()->Channel{Int}(32));
julia> const results = RemoteChannel(()->Channel{Tuple}(32));
julia> @everywhere function do_work(jobs, results) # define work function everywhere
while true
job_id = take!(jobs)
exec_time = rand()
sleep(exec_time) # simulates elapsed time doing actual work
put!(results, (job_id, exec_time, myid()))
end
end
julia> function make_jobs(n)
for i in 1:n
put!(jobs, i)
end
end;
julia> n = 12;
julia> errormonitor(@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 используются встроенными менеджерами. Пользовательские менеджеры кластеров обычно указывают только host или port:
Если
ioуказан, он используется для чтения информации об хосте/порте. Рабочий процесс Julia выводит адрес и порт привязки при запуске. Это позволяет рабочим процессам Julia прослушивать любой свободный порт, а не требовать ручного конфигурирования портов рабочих процессов.Если
ioне указан,hostиportиспользуются для подключения.-
count,exenameиexeflagsактуальны для запуска дополнительных рабочих процессов из рабочего процесса. Например, менеджер кластера может запустить один рабочий процесс на узел и использовать его для запуска дополнительных рабочих процессов.-
countсо значением целого числаnзапустит всегоnрабочих процессов. -
countсо значением:autoзапустит столько рабочих процессов, сколько потоков CPU (логических ядер) на данной машине. -
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 , Dagger.jl предоставляет функциональность, аналогичную Python's Dask, а DistributedArrays.jl предоставляет операции над массивами, распределенными по рабочим процессам, как показано в Общих массивах.
Следует упомянуть экосистему Julia для программирования на GPU, которая включает:
CUDA.jl оборачивает различные библиотеки CUDA и поддерживает компиляцию ядер Julia для графических процессоров Nvidia.
oneAPI.jl оборачивает унифицированную модель программирования oneAPI и поддерживает выполнение ядер Julia на поддерживаемых ускорителях. В настоящее время поддерживается только Linux.
AMDGPU.jl оборачивает библиотеки AMD ROCm и поддерживает компиляцию ядер Julia для графических процессоров AMD. В настоящее время поддерживается только Linux.
Высокоуровневые библиотеки, такие как KernelAbstractions.jl, Tullio.jl и ArrayFire.jl.
В приведенном ниже примере мы будем использовать как DistributedArrays.jl , так и CUDA.jl для распределения массива по нескольким процессам, сначала преобразовав его через distribute() и CuArray().
При импорте DistributedArrays.jl помните об импорте его во все процессы, используя @everywhere
$ ./julia -p 4
julia> addprocs()
julia> @everywhere using DistributedArrays
julia> using CUDA
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}
В приведенном ниже примере мы будем использовать как DistributedArrays.jl , так и CUDA.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.
© 2009–2022 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.8/manual/distributed-computing/