Многопроцессорность-и-распределённые-вычисления
Реализация параллельных вычислений с распределённой памятью предоставляется модулем 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 в репозитории Примеры, который использует словарь в качестве удаленного хранилища.
Каналы и 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> 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 , или DistributedArrays.jl, как показано в Раздел по общим массивам. Следует упомянуть экосистему Julia для программирования на GPU, которая включает:
Операции на низком уровне (на основе C-ядра) OpenCL.jl и CUDAdrv.jl, которые представляют собой соответственно интерфейс OpenCL и обертку CUDA.
Интерфейсы на низком уровне (ядро Julia), такие как CUDAnative.jl, который представляет собой собственное нативное для Julia CUDA-решение.
Абстракции высокого уровня, специфичные для поставщика, такие как 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Страницы руководства Julia GPU
© 2009–2021 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.7.0/manual/distributed-computing/