Многопроцессорная и распределённая вычисления
Реализация распределённых вычислений с разделенной памятью предоставляется модулем Distributed в стандартной библиотеке Julia.
Большинство современных компьютеров имеют более одного процессора, а несколько компьютеров могут быть объединены в кластер. Использование мощности этих множественных процессоров позволяет ускорить многие вычисления. Два основных фактора влияют на производительность: скорость самих процессоров и скорость доступа к памяти. В кластере очевидно, что конкретный процессор имеет самый быстрый доступ к оперативной памяти в том же компьютере (узле). Возможно, более неожиданно, что аналогичные проблемы актуальны на типичном многоядерном ноутбуке из-за различий в скорости оперативной памяти и кэша. Поэтому хорошая многопроцессорная среда должна позволять управлять «владением» блоком памяти конкретным процессором. Julia предоставляет многопроцессорную среду на основе передачи сообщений, чтобы программы могли выполняться на нескольких процессах в отдельных областях памяти одновременно.
Реализация передачи сообщений в Julia отличается от других сред, таких как MPI[1]. Общение в Julia обычно является «односторонним», что означает, что программист должен явно управлять только одним процессом в операции с двумя процессами. Кроме того, эти операции обычно не выглядят как «отправка сообщения» и «прием сообщения», а скорее похожи на операции более высокого уровня, такие как вызовы пользовательских функций.
Распределённое программирование в Julia основано на двух примитивах: удалённых ссылках и удалённых вызовах. Удалённая ссылка — это объект, который может быть использован любым процессом для ссылки на объект, хранящийся на конкретном процессе. Удалённый вызов — это запрос одного процесса вызвать определённую функцию с определёнными аргументами на другом процессе (возможно, на том же).
Удалённые ссылки бывают двух типов: Future и RemoteChannel.
Удалённый вызов возвращает Future для своего результата. Удалённые вызовы возвращаются немедленно; процесс, который сделал вызов, переходит к следующей операции, в то время как удалённый вызов происходит где-то ещё. Вы можете дождаться завершения удалённого вызова, вызвав wait на возвращённой Future, и вы можете получить полное значение результата, используя fetch.
С другой стороны, RemoteChannel являются перезаписываемыми. Например, несколько процессов могут координировать свою обработку, ссылаясь на ту же удалённую Channel.
Каждый процесс имеет связанный идентификатор. Процесс, предоставляющий интерактивный приглашение Julia, всегда имеет id равный 1. Процессы, используемые по умолчанию для параллельных операций, называются «рабочими». Когда есть только один процесс, процесс 1 считается рабочим. В противном случае, рабочими считаются все процессы, кроме процесса 1. В результате, добавление 2 или более процессов необходимо для получения преимуществ от параллельных методов обработки, таких как pmap. Добавление одного процесса полезно, если вы просто хотите делать другие вещи в главном процессе, пока длительное вычисление выполняется на рабочем.
Давайте попробуем это. Начав с julia -p n предоставляются n рабочих процессов на локальном компьютере. Обычно имеет смысл, чтобы n было равно количеству потоков процессора (логических ядер) на компьютере. Обратите внимание, что аргумент -p неявно загружает модуль Distributed.
$ ./julia -p 2
julia> r = remotecall(rand, 2, 2, 2)
Future(2, 1, 4, nothing)
julia> s = @spawnat 2 1 .+ fetch(r)
Future(2, 1, 5, nothing)
julia> fetch(s)
2×2 Array{Float64,2}:
1.18526 1.50912
1.16296 1.60607
Первый аргумент для remotecall — функция для вызова. Большинство параллельных вычислений в Julia не ссылаются на конкретные процессы или количество доступных процессов, но remotecall считается интерфейсом низкого уровня, обеспечивающим более тонкий контроль. Второй аргумент remotecall — это id процесса, который выполнит работу, а оставшиеся аргументы будут переданы вызываемой функции.
Как видите, в первой строке мы попросили процесс 2 создать 2x2 случайную матрицу, а во второй строке — добавить к ней 1. Результат обоих вычислений доступен в двух фьючерсах, r и s. Макрос @spawnat оценивает выражение во втором аргументе на процессе, указанном в первом аргументе.
Иногда вам может потребоваться значение, вычисленное удалённо, немедленно. Это обычно происходит, когда вы считываете удалённый объект, чтобы получить данные, необходимые для следующей локальной операции. Функция remotecall_fetch предназначена для этой цели. Она эквивалентна fetch(remotecall(...)), но более эффективна.
julia> remotecall_fetch(getindex, 2, r, 1, 1) 0.18526337335308085
Помните, что getindex(r,1,1) эквивалентна r[1,1], поэтому этот вызов извлекает первый элемент будущего r.
Для упрощения, символ :any может быть передан в [@spawnat], который выбирает, где выполнить операцию за вас:
julia> r = @spawnat :any rand(2,2)
Future(2, 1, 4, nothing)
julia> s = @spawnat :any 1 .+ fetch(r)
Future(3, 1, 5, nothing)
julia> fetch(s)
2×2 Array{Float64,2}:
1.38854 1.9098
1.20939 1.57158
Обратите внимание, что мы использовали 1 .+ fetch(r) вместо 1 .+ r. Это связано с тем, что мы не знаем, где будет выполняться код, поэтому в общем случае может потребоваться fetch для перемещения r на процесс, выполняющий добавление. В этом случае @spawnat достаточно умён, чтобы выполнить вычисление на процессе, который владеет r, поэтому fetch будет пустой операцией (работа не выполняется).
(Стоит отметить, что @spawnat не является встроенным, а определён в Julia как макрос. Можно определить собственные конструкции такого рода.)
Важно помнить, что после извлечения Future кэширует своё значение локально. Дальнейшие fetch вызовы не требуют сетевого взаимодействия. После того, как все ссылающиеся на Future значения будут извлечены, удалённое хранящееся значение удаляется.
@async похож на @spawnat, но выполняет задачи только на локальном процессе. Мы используем его для создания задачи-подачи для каждого процесса. Каждая задача выбирает следующий индекс, который нужно вычислить, затем ждёт завершения своего процесса, а затем повторяет, пока не исчерпаются индексы. Обратите внимание, что задачи-подачи не начинают выполняться, пока главная задача не достигнет конца блока @sync, в этот момент она отдает управление и ждет завершения всех локальных задач, прежде чем вернуть управление из функции. Что касается версии v0.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 (с того же пути, что и у текущего хоста) на указанных машинах.
Функции addprocs, rmprocs, workers и другие доступны как программный способ добавления, удаления и запроса процессов в кластере.
julia> using Distributed
julia> addprocs(2)
2-element Array{Int64,1}:
2
3
Модуль Distributed должен быть явно загружен на главном процессе перед вызовом addprocs. Он автоматически становится доступным на рабочих процессах.
Обратите внимание, что рабочие процессы не выполняют скрипт запуска ~/.julia/config/startup.jl, и они не синхронизируют своё глобальное состояние (такое как глобальные переменные, новые определения методов и загруженные модули) ни с одним из других работающих процессов. Вы можете использовать addprocs(exeflags="--project") для инициализации рабочего процесса с определённой средой, а затем @everywhere using <modulename> или @everywhere include("file.jl").
Другие типы кластеров могут поддерживаться написанием собственного пользовательского ClusterManager, как описано ниже в разделе Управляющие кластерами.
Перемещение данных
Отправка сообщений и перемещение данных составляют большую часть накладных расходов в распределённой программе. Сокращение количества сообщений и объёма отправляемых данных имеет решающее значение для достижения производительности и масштабируемости. Для этого важно понимать перемещение данных, выполняемое различными конструкциями распределённого программирования Julia.
fetch можно рассматривать как явную операцию перемещения данных, поскольку она напрямую запрашивает перемещение объекта на локальную машину. @spawnat (и несколько связанных конструкций) также перемещает данные, но это не так очевидно, поэтому это можно назвать неявной операцией перемещения данных. Рассмотрим два подхода к построению и возведению в квадрат случайной матрицы:
Метод 1:
julia> A = rand(1000,1000); julia> Bref = @spawnat :any A^2; [...] julia> fetch(Bref);
Метод 2:
julia> Bref = @spawnat :any rand(1000,1000)^2; [...] julia> fetch(Bref);
Разница кажется незначительной, но на самом деле она весьма существенна из-за поведения @spawnat. В первом методе случайная матрица строится локально, затем отправляется на другой процесс, где она возводится в квадрат. Во втором методе случайная матрица строится и возводится в квадрат на другом процессе. Следовательно, второй метод отправляет гораздо меньше данных, чем первый.
В этом учебном примере два метода легко отличить и выбрать из них. Однако в реальной программе проектирование перемещения данных может потребовать больше размышлений и, вероятно, некоторых измерений. Например, если первому процессу нужна матрица A, то первый метод может быть лучше. Или, если вычисление A дорогостоящее и только текущий процесс его имеет, то перемещение его на другой процесс может быть неизбежным. Или, если текущий процесс имеет очень мало работы между @spawnat и fetch(Bref), то возможно лучше полностью отказаться от параллелизма. Или представьте, что rand(1000,1000) заменено на более дорогостоящую операцию. Тогда может иметь смысл добавить ещё одну инструкцию @spawnat только для этого шага.
Глобальные переменные
Выражения, выполняемые удалённо с помощью @spawnat, или лямбда-функции, указанные для удалённого выполнения с помощью remotecall, могут ссылаться на глобальные переменные. Глобальные привязки в модуле Main обрабатываются немного по-другому по сравнению с глобальными привязками в других модулях. Рассмотрим следующий фрагмент кода:
A = rand(10,10) remotecall_fetch(()->sum(A), 2)
В этом случае sum ДОЛЖЕН быть определён в удалённом процессе. Обратите внимание, что A — это глобальная переменная, определённая в локальном рабочем пространстве. Рабочий процесс 2 не имеет переменной, называемой A в Main. Действие отправки лямбда-функции ()->sum(A) на рабочий процесс 2 приводит к определению Main.A на процессе 2. Main.A продолжает существовать на рабочем процессе 2 даже после возвращения вызова remotecall_fetch. Удалённые вызовы с встроенными ссылками на глобальные переменные (только в модуле Main ) управляют глобальными переменными следующим образом:
Новые глобальные привязки создаются на целевых рабочих процессах, если они ссылаются в рамках удалённого вызова.
Глобальные константы объявляются как константы и на удалённых узлах.
-
Глобальные переменные пересылаются на целевой рабочий процесс только в контексте удалённого вызова, и только если их значение изменилось. Кроме того, кластер не синхронизирует глобальные привязки между узлами. Например:
A = rand(10,10) remotecall_fetch(()->sum(A), 2) # worker 2 A = rand(10,10) remotecall_fetch(()->sum(A), 3) # worker 3 A = nothing
Выполнение приведенного выше фрагмента кода приводит к тому, что
Main.Aна рабочем процессе 2 имеет другое значение, чемMain.Aна рабочем процессе 3, в то время как значениеMain.Aна узле 1 установлено вnothing.
Как вы могли заметить, хотя память, связанная с глобальными переменными, может быть освобождена при переприсваивании на главном процессе, такая же операция не выполняется на рабочих процессах, так как привязки продолжают быть действительными. clear! можно использовать для ручного переприсваивания определённых глобальных переменных на удалённых узлах в nothing после того, как они больше не требуются. Это позволит освободить связанную с ними память в рамках обычного цикла сбора мусора.
Таким образом, программы должны быть осторожны при ссылке на глобальные переменные в удалённых вызовах. На самом деле, если это возможно, их лучше вообще избегать. Если вам необходимо ссылаться на глобальные переменные, рассмотрите использование блоков let для локализации глобальных переменных.
Например:
julia> A = rand(10,10);
julia> remotecall_fetch(()->A, 2);
julia> B = rand(10,10);
julia> let B = B
remotecall_fetch(()->B, 2)
end;
julia> @fetchfrom 2 InteractiveUtils.varinfo()
name size summary
––––––––– ––––––––– ––––––––––––––––––––––
A 800 bytes 10×10 Array{Float64,2}
Base Module
Core Module
Main Module
Как видно, глобальная переменная A определена на рабочем процессе 2, но B захвачена как локальная переменная, и поэтому привязки для B не существует на рабочем процессе 2.
Параллельные 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 выглядят как последовательные, их поведение существенно отличается. В частности, итерации не происходят в заданном порядке, а записи в переменные или массивы не будут глобально видны, так как итерации выполняются на разных процессах. Любые переменные, используемые внутри параллельного цикла, будут скопированы и переданы каждому процессу.
Например, следующий код не будет работать как ожидается:
a = zeros(100000)
@distributed for i = 1:100000
a[i] = i
end
Этот код не инициализирует все a, так как у каждого процесса будет своя копия. Параллельные циклы for, подобные этому, следует избегать. К счастью, Общие массивы могут быть использованы для решения этой проблемы:
using SharedArrays
a = SharedArray{Float64}(10)
@distributed for i = 1:10
a[i] = i
end
Использование «внешних» переменных в параллельных циклах является вполне допустимым, если переменные только для чтения:
a = randn(1000)
@distributed (+) for i = 1:100000
f(a[rand(1:end)])
end
Здесь каждая итерация применяет f к случайно выбранному образцу из вектора a, общим для всех процессов.
Как вы могли заметить, оператор редукции можно опустить, если он не нужен. В этом случае цикл выполняется асинхронно, т. е. он порождает независимые задачи на всех доступных рабочих процессах и возвращает массив Future немедленно, не ожидая завершения. Вызывающий процесс может дождаться завершения Future в более поздний момент, вызвав fetch для них, или дождаться завершения в конце цикла, добавив @sync, например @sync @distributed for.
В некоторых случаях оператор редукции не нужен, и нам нужно лишь применить функцию ко всем целым числам в некотором диапазоне (или, более обще, ко всем элементам в некотором наборе). Это ещё одна полезная операция, называемая параллельным map, реализованная в 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. Задачи, идентифицируемые по идентификатору (job_id), записываются в канал. Каждая выполняемая удаленно задача в этом симуляторе считывает job_id, ждет случайное время и записывает в канал результатов кортеж из job_id, времени выполнения и собственного pid в канал результатов. Наконец, все results выводятся на главном процессе.
julia> addprocs(4); # add worker processes
julia> const jobs = RemoteChannel(()->Channel{Int}(32));
julia> const results = RemoteChannel(()->Channel{Tuple}(32));
julia> @everywhere function do_work(jobs, results) # define work function everywhere
while true
job_id = take!(jobs)
exec_time = rand()
sleep(exec_time) # simulates elapsed time doing actual work
put!(results, (job_id, exec_time, myid()))
end
end
julia> function make_jobs(n)
for i in 1:n
put!(jobs, i)
end
end;
julia> n = 12;
julia> @async make_jobs(n); # feed the jobs channel with "n" jobs
julia> for p in workers() # start tasks on the workers to process requests in parallel
remote_do(do_work, p, jobs, results)
end
julia> @elapsed while n > 0 # print out results
job_id, exec_time, where = take!(results)
println("$job_id finished in $(round(exec_time; digits=2)) seconds on worker $where")
global n = n - 1
end
1 finished in 0.18 seconds on worker 4
2 finished in 0.26 seconds on worker 5
6 finished in 0.12 seconds on worker 4
7 finished in 0.18 seconds on worker 4
5 finished in 0.35 seconds on worker 5
4 finished in 0.68 seconds on worker 2
3 finished in 0.73 seconds on worker 3
11 finished in 0.01 seconds on worker 3
12 finished in 0.02 seconds on worker 3
9 finished in 0.26 seconds on worker 5
8 finished in 0.57 seconds on worker 4
10 finished in 0.58 seconds on worker 2
0.055971741
Удаленные ссылки и распределённый сбор мусора
Объекты, на которые ссылаются удаленные ссылки, могут быть освобождены только тогда, когда все ссылки на них в кластере будут удалены.
Узел, где хранится значение, отслеживает, у каких рабочих узлов есть ссылка на него. Каждый раз, когда RemoteChannel или (неизвлеченное) Future сериализуется на рабочий узел, узел, на который указывает ссылка, уведомляется. И каждый раз, когда RemoteChannel или (неизвлеченное) Future собирается мусор локально, узел, владеющий значением, вновь уведомляется. Это реализовано в внутреннем кластерно-осознанном сериализаторе. Удаленные ссылки действительны только в контексте работающего кластера. Сериализация и десериализация ссылок на обычные IO объекты не поддерживается.
Уведомления осуществляются путём отправки сообщений «отслеживания» — сообщения «добавить ссылку» при сериализации ссылки на другой процесс и сообщения «удалить ссылку» при локальном сборе мусора ссылки.
Поскольку Future одноразовые и кешируются локально, действие fetch Future также обновляет информацию об отслеживании ссылок на узле, владеющем значением.
Узел, который владеет значением, освобождает его после того, как все ссылки на него будут удалены.
С Future, сериализация уже извлеченного Future на другой узел также отправляет значение, так как оригинальное удалённое хранилище может к этому времени собрать значение.
Важно отметить, что когда объект собирается мусором локально, зависит от размера объекта и текущего давления памяти в системе.
В случае удаленных ссылок размер локального объекта ссылки невелик, а значение, хранящееся на удалённом узле, может быть довольно большим. Поскольку локальный объект может не быть собран сразу, рекомендуется явно вызывать finalize на локальных экземплярах RemoteChannel или на неизвлеченных Future. Поскольку вызов fetch на Future также удаляет его ссылку из удалённого хранилища, это не требуется для извлечённых Future. Явное вызов finalize приводит к немедленной отправке сообщения на удалённый узел, чтобы он удалил ссылку на значение.
После завершения ссылка становится недействительной и не может быть использована в последующих вызовах.
Локальные вызовы
Данные обязательно копируются на удалённый узел для выполнения. Это справедливо как для удалённых вызовов, так и для случаев, когда данные хранятся в RemoteChannel / Future на другом узле. Как ожидалось, это приводит к созданию копии сериализованных объектов на удалённом узле. Однако, когда целевой узел — локальный узел, то есть идентификатор вызывающего процесса совпадает с идентификатором удалённого узла, он выполняется как локальный вызов. Обычно (но не всегда) выполняется в другой задаче — но не происходит сериализация/десериализация данных. В результате вызов ссылается на те же экземпляры объектов, что и переданные — не создаются копии. Это поведение показано ниже:
julia> using Distributed;
julia> rc = RemoteChannel(()->Channel(3)); # RemoteChannel created on local node
julia> v = [0];
julia> for i in 1:3
v[1] = i # Reusing `v`
put!(rc, v)
end;
julia> result = [take!(rc) for _ in 1:3];
julia> println(result);
Array{Int64,1}[[3], [3], [3]]
julia> println("Num Unique objects : ", length(unique(map(objectid, result))));
Num Unique objects : 1
julia> addprocs(1);
julia> rc = RemoteChannel(()->Channel(3), workers()[1]); # RemoteChannel created on remote node
julia> v = [0];
julia> for i in 1:3
v[1] = i
put!(rc, v)
end;
julia> result = [take!(rc) for _ in 1:3];
julia> println(result);
Array{Int64,1}[[1], [2], [3]]
julia> println("Num Unique objects : ", length(unique(map(objectid, result))));
Num Unique objects : 3
Как видно, put! на локальном RemoteChannel с тем же объектом v , изменённым между вызовами, приводит к тому же единственному экземпляру объекта, хранящемуся. В отличие от копий v, создаваемых, когда узел, владеющий rc, — это другой узел.
Следует отметить, что это обычно не проблема. Это нужно учитывать только в том случае, если объект хранится локально и изменяется после вызова. В таких случаях может быть уместно хранить deepcopy объекта.
Это также верно для удалённых вызовов на локальном узле, как показано в следующем примере:
julia> using Distributed; addprocs(1);
julia> v = [0];
julia> v2 = remotecall_fetch(x->(x[1] = 1; x), myid(), v); # Executed on local node
julia> println("v=$v, v2=$v2, ", v === v2);
v=[1], v2=[1], true
julia> v = [0];
julia> v2 = remotecall_fetch(x->(x[1] = 1; x), workers()[1], v); # Executed on remote node
julia> println("v=$v, v2=$v2, ", v === v2);
v=[0], v2=[1], false
Как видно ещё раз, удалённый вызов на локальный узел ведёт себя так же, как прямой вызов. Вызов изменяет локальные объекты, переданные в качестве аргументов. В удалённом вызове он работает с копией аргументов.
Повторим, что это обычно не проблема. Если локальный узел также используется как вычислительный узел и аргументы используются после вызова, это поведение нужно учитывать, и при необходимости следует передавать глубокие копии аргументов в вызываемый на локальном узле вызов. Вызовы на удалённые узлы всегда будут работать с копиями аргументов.
Общие массивы
Упорядоченные массивы используют общую системную память для отображения одного и того же массива во многих процессах. Хотя есть некоторое сходство с DArray, поведение SharedArray отличается. В DArray каждый процесс имеет локальный доступ только к части данных, и никакие два процесса не делят одну и ту же часть; напротив, в SharedArray каждый «участвующий» процесс имеет доступ ко всему массиву. SharedArray — хороший выбор, когда вам нужно, чтобы большой объем данных был совместно доступен двум или более процессам на одном компьютере.
Поддержка упорядоченных массивов доступна через модуль SharedArrays, который необходимо явно загрузить на всех участвующих рабочих процессах.
SharedArray индексация (присваивание и доступ к значениям) работает так же, как и для обычных массивов, и является эффективной, потому что базовая память доступна локальному процессу. Поэтому большинство алгоритмов естественным образом работают с SharedArray, хотя и в режиме одного процесса. В тех случаях, когда алгоритм требует входной Array, базовый массив можно извлечь из SharedArray, вызвав sdata. Для других типов AbstractArray, sdata просто возвращает сам объект, поэтому использовать sdata с любым объектом типа Array безопасно.
Конструктор для упорядоченного массива имеет вид:
SharedArray{T,N}(dims::NTuple; init=false, pids=Int[])
который создает N-мерный упорядоченный массив типа 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 ДОЛЖНА завершиться, как только все запрошенные рабочие процессы будут запущены.
Новoзапущенные рабочие процессы подключены друг к другу и к главному процессу по типу «все ко всем». Указание командной строки --worker[=<cookie>] приводит к тому, что запущенные процессы инициализируются как рабочие процессы, и подключения устанавливаются через TCP/IP-сокеты.
Все рабочие процессы в кластере используют один и тот же куки , что и мастер-процесс. Если куки не указан, то есть с опцией --worker, рабочий процесс пытается прочитать его из стандартного ввода. LocalManager и SSHManager оба передают куки новым запущенным рабочим процессам через их стандартный ввод.
По умолчанию рабочий процесс будет прослушивать свободный порт по адресу, возвращаемому вызовом getipaddr(). Указать конкретный адрес для прослушивания можно с помощью необязательного аргумента --bind-to bind_addr[:port]. Это полезно для многоадресных хостов.
В качестве примера транспортного уровня, отличного от TCP/IP, реализация может использовать MPI, в этом случае --worker не должен быть указан. Вместо этого, вновь запущенные рабочие процессы должны вызвать init_worker(cookie) перед использованием каких-либо параллельных конструкций.
Для каждого запущенного рабочего процесса метод launch должен добавить объект WorkerConfig (с инициализированными соответствующими полями) в launched
mutable struct WorkerConfig
# Common fields relevant to all cluster managers
io::Union{IO, Nothing}
host::Union{AbstractString, Nothing}
port::Union{Integer, Nothing}
# Used when launching additional workers at a host
count::Union{Int, Symbol, Nothing}
exename::Union{AbstractString, Cmd, Nothing}
exeflags::Union{Cmd, Nothing}
# External cluster managers can use this to store information at a per-worker level
# Can be a dict if multiple fields need to be stored.
userdata::Any
# SSHManager / SSH tunnel connections to workers
tunnel::Union{Bool, Nothing}
bind_addr::Union{AbstractString, Nothing}
sshflags::Union{Cmd, Nothing}
max_parallel::Union{Integer, Nothing}
# Used by Local/SSH managers
connect_at::Any
[...]
end
Большинство полей в WorkerConfig используются встроенными менеджерами. Пользовательские менеджеры кластеров, как правило, указывают только io или host / port:
Если
ioуказан, он используется для чтения информации о хосте/порту. Рабочий процесс Julia выводит адрес и порт привязки при запуске. Это позволяет рабочим процессам Julia прослушивать любой свободный порт, вместо необходимости ручной настройки портов рабочих процессов.Если
ioне указан, для подключения используютсяhostиport.-
count,exenameиexeflagsимеют отношение к запуску дополнительных рабочих процессов из рабочего процесса. Например, менеджер кластера может запускать по одному рабочему процессу на узел и использовать его для запуска дополнительных рабочих процессов.-
countсо значением в виде целого числаnзапустит общее количествоnрабочих процессов. -
countсо значением:autoзапустит столько рабочих процессов, сколько процессорных потоков (логических ядер) на этом компьютере. -
exename— имя исполняемого файлаjuliaвключая полный путь. -
exeflagsдолжен содержать требуемые аргументы командной строки для новых рабочих процессов.
-
tunnel,bind_addr,sshflagsиmax_parallelиспользуются, когда для подключения к рабочим процессам из мастер-процесса требуется SSH-туннель.userdataпредоставляется пользовательским менеджерам кластеров для хранения собственной информации, специфичной для рабочего процесса.
manage(manager::FooManager, id::Integer, config::WorkerConfig, op::Symbol) вызывается в разное время во время жизненного цикла рабочего процесса со значениями op:
- со значениями
:register/:deregisterпри добавлении/удалении рабочего процесса из пула рабочих процессов Julia. - со значением
:interruptпри вызовеinterrupt(workers).ClusterManagerдолжен сигнализировать соответствующему рабочему процессу с помощью сигнала прерывания. - со значением
:finalizeдля целей очистки.
Менеджеры кластеров с пользовательскими транспортными средствами
Замена стандартных соединений сокетов TCP/IP всеми рабочими процессами на пользовательский транспортный уровень несколько сложнее. У каждого процесса Julia столько задач связи, сколько рабочих процессов к нему подключено. Например, рассмотрим кластер Julia из 32 процессов в сети с полносвязным топологией:
- Таким образом, у каждого процесса Julia есть 31 задача связи.
- Каждая задача обрабатывает все входящие сообщения от одного удаленного рабочего процесса в цикле обработки сообщений.
- Цикл обработки сообщений ожидает объект
IO(например,TCPSocketв реализации по умолчанию), считывает все сообщение, обрабатывает его и ожидает следующего. - Отправка сообщений процессу выполняется напрямую из любой задачи Julia — а не только задач связи — опять же через соответствующий объект
IO.
Замена стандартного транспорта требует, чтобы новая реализация устанавливала подключения к удаленным рабочим процессам и предоставляла соответствующие объекты IO , на которых могут ожидать циклы обработки сообщений. Менеджер-специфические обратные вызовы, которые необходимо реализовать:
connect(manager::FooManager, pid::Integer, config::WorkerConfig) kill(manager::FooManager, pid::Int, config::WorkerConfig)
Реализация по умолчанию (которая использует сокеты TCP/IP) реализована как connect(manager::ClusterManager, pid::Integer, config::WorkerConfig).
connect должен вернуть пару объектов IO , один для чтения данных, отправленных рабочим процессом pid, а другой для записи данных, которые необходимо отправить рабочему процессу pid. Пользовательские менеджеры кластеров могут использовать в памяти BufferStream в качестве коммуникационного канала для передачи данных между пользовательским, возможно, не-IO транспортом и встроенной параллельной инфраструктурой Julia.
BufferStream — это в памяти IOBuffer, который ведет себя как IO — это поток, который может обрабатываться асинхронно.
Папка clustermanager/0mq в репозитории примеров содержит пример использования ZeroMQ для соединения рабочих процессов Julia в звездообразной топологии с брокером 0MQ посередине. Примечание: процессы Julia всё ещё логически соединены друг с другом — любой рабочий процесс может отправлять сообщения любому другому рабочему процессу без какого-либо понимания использования 0MQ в качестве транспортного уровня.
При использовании пользовательских транспортов:
- Рабочие процессы Julia не должны запускаться с
--worker. Запуск с--workerприведет к тому, что вновь запущенные рабочие процессы перейдут к реализации транспортного уровня сокетов TCP/IP. - Для каждого входящего логического соединения с рабочим процессом должен быть вызван
Base.process_messages(rd::IO, wr::IO)(). Это запускает новую задачу, которая обрабатывает чтение и запись сообщений от/в рабочий процесс, представленный объектамиIO. -
init_worker(cookie, manager::FooManager)обязательно должен быть вызван во время инициализации процесса рабочего процесса. - Поле
connect_at::AnyвWorkerConfigможет быть установлено менеджером кластера при вызовеlaunch. Значение этого поля передаётся во все обратные вызовыconnect. Как правило, оно содержит информацию о том, как подключиться к рабочему процессу. Например, транспортный уровень сокетов TCP/IP использует это поле для указания кортежа(host, port)для подключения к рабочему процессу.
kill(manager, pid, config) вызывается для удаления рабочего процесса из кластера. В мастер-процессе соответствующие объекты IO должны быть закрыты реализацией для обеспечения надлежащей очистки. Реализация по умолчанию просто выполняет вызов exit() для указанного удаленного рабочего процесса.
Папка примеров clustermanager/simple — это пример, демонстрирующий простую реализацию с использованием сокетов UNIX для настройки кластера.
Требования к сети для LocalManager и SSHManager
Кластеры Julia разработаны для выполнения в уже защищенных средах на инфраструктуре, такой как локальные ноутбуки, кластеры департаментов или облака. В этом разделе рассматриваются требования к безопасности сети для встроенных LocalManager и SSHManager:
Мастер-процесс не прослушивает порты. Он подключается только к рабочим процессам.
Каждый рабочий процесс привязывается только к одному из локальных интерфейсов и прослушивает назначенный операционной системой временный номер порта.
LocalManager, используемыйaddprocs(N), по умолчанию связывается только с петлевым интерфейсом. Это означает, что рабочие процессы, запущенные позже на удаленных хостах (или любым лицом с вредоносными намерениями), не могут подключиться к кластеру. Вызовaddprocs(4)и последующийaddprocs(["remote_host"])потерпит неудачу. Некоторые пользователям может потребоваться создать кластер, включающий их локальную систему и несколько удалённых систем. Это можно сделать, явно указавLocalManagerпривязываться к внешнему сетевому интерфейсу с помощью ключевого аргументаrestrict:addprocs(4; restrict=false).SSHManager, используемыйaddprocs(list_of_remote_hosts), запускает рабочие процессы на удаленных хостах через SSH. По умолчанию SSH используется только для запуска рабочих процессов Julia. Последующие соединения мастер-рабочий процесс и рабочий процесс-рабочий процесс используют обычные, незашифрованные сокеты TCP/IP. На удалённых хостах должна быть включена возможность входа без пароля. Дополнительные флаги SSH или учетные данные могут быть указаны с помощью ключевого аргументаsshflags.-
addprocs(list_of_remote_hosts; tunnel=true, sshflags=<ssh keys and other flags>)полезно, когда нужно использовать SSH-соединения и для мастер-рабочего процесса. Типичный сценарий для этого — локальный ноутбук, на котором запущен интерпретатор Julia (т. е. мастер), с остальным кластером в облаке, например, на 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 для графических процессоров, которая включает:
Операции на низком уровне (ядра C) OpenCL.jl и CUDAdrv.jl, которые соответственно являются интерфейсом OpenCL и оберткой CUDA.
Интерфейсы низкого уровня (ядро Julia) такие как CUDAnative.jl, который представляет собой собственное реализацию CUDA на Julia.
Абстракции высокого уровня от поставщиков, такие как CuArrays.jl и CLArrays.jl
Библиотеки высокого уровня, такие как ArrayFire.jl и GPUArrays.jl
В следующем примере мы будем использовать как DistributedArrays.jl, так и CuArrays.jl, чтобы распределить массив по нескольким процессам, сначала преобразовав его через distribute() и CuArray().
Помните, при импорте DistributedArrays.jl импортируйте его во все процессы, используя @everywhere
$ ./julia -p 4
julia> addprocs()
julia> @everywhere using DistributedArrays
julia> using CuArrays
julia> B = ones(10_000) ./ 2;
julia> A = ones(10_000) .* π;
julia> C = 2 .* A ./ B;
julia> all(C .≈ 4*π)
true
julia> typeof(C)
Array{Float64,1}
julia> dB = distribute(B);
julia> dA = distribute(A);
julia> dC = 2 .* dA ./ dB;
julia> all(dC .≈ 4*π)
true
julia> typeof(dC)
DistributedArrays.DArray{Float64,1,Array{Float64,1}}
julia> cuB = CuArray(B);
julia> cuA = CuArray(A);
julia> cuC = 2 .* cuA ./ cuB;
julia> all(cuC .≈ 4*π);
true
julia> typeof(cuC)
CuArray{Float64,1}
Обратите внимание, что некоторые функции Julia в настоящее время не поддерживаются CUDAnative.jl[2], особенно некоторые функции, такие как sin должны быть заменены на CUDAnative.sin(cc: @maleadt).
В следующем примере мы будем использовать как DistributedArrays.jl, так и CuArrays.jl, чтобы распределить массив по нескольким процессам и вызвать на нём общую функцию.
function power_method(M, v)
for i in 1:100
v = M*v
v /= norm(v)
end
return v, norm(M*v) / norm(v) # or (M*v) ./ v
end
power_method многократно создаёт новый вектор и нормализует его. Мы не указали сигнатуру типа в объявлении функции, давайте посмотрим, будет ли это работать с вышеупомянутыми типами данных:
julia> M = [2. 1; 1 1];
julia> v = rand(2)
2-element Array{Float64,1}:
0.40395
0.445877
julia> power_method(M,v)
([0.850651, 0.525731], 2.618033988749895)
julia> cuM = CuArray(M);
julia> cuv = CuArray(v);
julia> curesult = power_method(cuM, cuv);
julia> typeof(curesult)
CuArray{Float64,1}
julia> dM = distribute(M);
julia> dv = distribute(v);
julia> dC = power_method(dM, dv);
julia> typeof(dC)
Tuple{DistributedArrays.DArray{Float64,1,Array{Float64,1}},Float64}
Чтобы завершить это краткое знакомство с внешними пакетами, рассмотрим MPI.jl, оболочку Julia протокола MPI. Поскольку рассмотрение каждой внутренней функции займёт слишком много времени, лучше просто оценить подход, используемый для реализации протокола.
Рассмотрим этот игрушечный скрипт, который просто вызывает каждый подпроцесс, инициализирует его ранг, и когда достигает главного процесса, производит сумму рангов
import MPI
MPI.Init()
comm = MPI.COMM_WORLD
MPI.Barrier(comm)
root = 0
r = MPI.Comm_rank(comm)
sr = MPI.Reduce(r, MPI.SUM, root, comm)
if(MPI.Comm_rank(comm) == root)
@printf("sum of ranks: %s\n", sr)
end
MPI.Finalize()
mpirun -np 4 ./julia example.jl
- 1В этом контексте MPI относится к стандарту MPI-1. Начиная с MPI-2, комитет стандартов MPI ввёл новый набор механизмов связи, объединённых под общим названием «Доступ к удалённой памяти» (RMA). Мотивацией для добавления rma в стандарт MPI было упрощение односторонних схем связи. Для получения дополнительной информации о последнем стандарте MPI см. https://mpi-forum.org/docs.
- 2Страницы справки Julia GPU
© 2009–2020 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.5.3/manual/distributed-computing/