Spec-Zone.ru › Julia 1.10

Многопроцессорная обработка и распределённые вычисления

Реализация параллельных вычислений с распределённой памятью предоставляется модулем 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(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, в котором она уступает управление и ждёт завершения всех локальных задач, прежде чем возвращаться из функции. Что касается версии 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 worker (из той же директории, что и текущий хост) на указанных машинах. Каждое определение машины имеет вид [count*][user@]host[:port] [bind_addr[:port]]. user по умолчанию равен текущему пользователю, port — стандартному порту ssh. count — количество рабочих процессов, запускаемых на узле, по умолчанию равно 1. Необязательный параметр bind-to bind_addr[:port] задаёт IP-адрес и порт, которые другие рабочие процессы должны использовать для подключения к этому рабочему процессу.

Хотя Julia в целом стремится к обратной совместимости, распределение кода рабочим процессам зависит от Serialization.serialize. Как указано в соответствующей документации, это не гарантирует работоспособность в разных версиях 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.

Параллельное применение и циклы

К счастью, многие полезные параллельные вычисления не требуют перемещения данных. Типичный пример — симуляция Монте-Карло, где несколько процессов могут обрабатывать независимые испытания симуляции одновременно. Мы можем использовать @spawnat для подбрасывания монет на двух процессах. Сначала напишите следующую функцию в count_heads.jl:

function count_heads(n)
    c::Int = 0
    for i = 1:n
        c += rand(Bool)
    end
    c
end

Функция count_heads просто складывает n случайных битов. Вот как мы можем выполнить несколько испытаний на двух машинах и сложить результаты:

julia> @everywhere include_string(Main, $(read("count_heads.jl", String)), "count_heads.jl")

julia> a = @spawnat :any count_heads(100000000)
Future(2, 1, 6, nothing)

julia> b = @spawnat :any count_heads(100000000)
Future(3, 1, 7, nothing)

julia> fetch(a)+fetch(b)
100001564

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

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

nheads = @distributed (+) for i = 1:200000000
    Int(rand(Bool))
end

Эта конструкция реализует шаблон назначения итераций многим процессам и их объединения с помощью указанной редукции (в данном случае (+)). Результат каждой итерации принимается как значение последнего выражения внутри цикла. Само выражение параллельного цикла оценивается в итоговый ответ.

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

Например, следующий код не будет работать как ожидается:

a = zeros(100000)
@distributed for i = 1:100000
    a[i] = i
end

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

using SharedArrays

a = SharedArray{Float64}(10)
@distributed for i = 1:10
    a[i] = i
end

Использование «внешних» переменных в параллельных циклах вполне допустимо, если переменные являются только для чтения:

a = randn(1000)
@distributed (+) for i = 1:100000
    f(a[rand(1:end)])
end

Здесь каждая итерация применяет f к случайно выбранной выборке из вектора a, используемого всеми процессами.

Как вы можете видеть, оператор сокращения можно опустить, если он не нужен. В этом случае цикл выполняется асинхронно, т.е. он запускает независимые задачи на всех доступных рабочих процессах и возвращает массив Future сразу же без ожидания завершения. Вызывающий процесс может дождаться завершения Future в более поздний момент, вызвав fetch на них, или дождаться завершения в конце цикла, добавив @sync, как @sync @distributed for.

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

julia> M = Matrix{Float64}[rand(1000,1000) for i = 1:10];

julia> pmap(svdvals, M);

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

Удаленные ссылки и абстрактные каналы

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

Для реализации put!, take!, fetch, isready и wait требуется конкретная реализация AbstractChannel (например, Channel). Объект удаленного объекта, на который ссылается Future, хранится в Channel{Any}(1), т.е. в Channel размером 1, способном хранить объекты типа Any.

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

Конструктор RemoteChannel(f::Function, pid)() позволяет создавать ссылки на каналы, содержащие более одного значения определенного типа. f — это функция, выполняемая на pid, и она должна возвращать AbstractChannel.

Например, RemoteChannel(()->Channel{Int}(10), pid), вернёт ссылку на канал типа Int и размера 10. Канал существует на рабочем процессе pid.

Методы put!, take!, fetch, isready и wait для RemoteChannel делегируются на хранилище на удаленном процессе.

RemoteChannel таким образом может ссылаться на реализованные пользователем объекты AbstractChannel. Простой пример этого показан в dictchannel.jl в репозитории Примеров, который использует словарь в качестве своего удаленного хранилища.

Каналы и RemoteChannels

  • Канал Channel локален для процесса. Рабочий процесс 2 не может напрямую обратиться к Channel на рабочем процессе 3 и наоборот. Однако RemoteChannel может передавать и получать значения между рабочими процессами.
  • RemoteChannel можно рассматривать как дескриптор для Channel.
  • Идентификатор процесса pid, связанный с RemoteChannel, идентифицирует процесс, где находится хранилище, то есть основной Channel.
  • Любой процесс с ссылкой на RemoteChannel может помещать и извлекать элементы из канала. Данные автоматически отправляются (или извлекаются) в процесс, к которому относится RemoteChannel.
  • Сериализация Channel также сериализует любые данные, присутствующие в канале. Поэтому десериализация фактически создаёт копию исходного объекта.
  • С другой стороны, сериализация RemoteChannel включает только сериализацию идентификатора, который идентифицирует местоположение и экземпляр Channel, на который ссылается дескриптор. Таким образом, десериализованный объект RemoteChannel (на любом рабочем процессе) также указывает на то же хранилище, что и исходный.

Пример каналов, приведенный выше, можно изменить для межпроцессного взаимодействия, как показано ниже.

Мы запускаем 4 рабочих процесса для обработки одного jobs удаленного канала. Задачи, идентифицируемые по идентификатору (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]. Это полезно для хостов с несколькими IP-адресами.

В качестве примера не-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 запустит столько рабочих процессов, сколько потоков 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: только процесс-драйвер, т.е. pid 1, имеет подключения к рабочим процессам.
  • :custom: метод launch менеджера кластера указывает топологию подключения через поля ident и connect_idents в WorkerConfig. Рабочий процесс с идентификатором, предоставленным менеджером кластера, ident подключится ко всем рабочим процессам, указанным в connect_idents.

Ключевой аргумент lazy=true|false влияет только на опцию topology опцию :all_to_all. Если true, кластер запускается с мастером, подключенным ко всем рабочим процессам. Установка конкретных подключений между рабочими процессами осуществляется при первом удаленном вызове между двумя рабочими процессами. Это помогает снизить начальные ресурсы, выделенные для внутрикластерной коммуникации. Подключения устанавливаются в зависимости от временных требований параллельной программы. Значение по умолчанию для lazy равно true.

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

Заслуживающие внимания внешние пакеты

Помимо параллелизма в Julia, есть много внешних пакетов, которые стоит упомянуть. Например, MPI.jl — оболочка Julia для протокола MPI, Dagger.jl предоставляет функциональность, аналогичную Python's Dask, а DistributedArrays.jl предоставляет операции с массивами, распределенные по рабочим процессам, как показано в Объединенных массивах.

Следует упомянуть и экосистему Julia для программирования на GPU, которая включает:

  1. CUDA.jl оборачивает различные библиотеки CUDA и поддерживает компиляцию ядер Julia для графических процессоров Nvidia.

  2. oneAPI.jl оборачивает унифицированную модель программирования oneAPI и поддерживает выполнение ядер Julia на поддерживаемых ускорителях. В настоящее время поддерживается только Linux.

  3. AMDGPU.jl оборачивает библиотеки AMD ROCm и поддерживает компиляцию ядер Julia для графических процессоров AMD. В настоящее время поддерживается только Linux.

  4. Высокоуровневые библиотеки, такие как KernelAbstractions.jl, Tullio.jl и ArrayFire.jl.

В следующем примере мы будем использовать DistributedArrays.jl и CUDA.jl для распределения массива по нескольким процессам, сначала преобразовав его через distribute() и CuArray().

Помните, что при импорте DistributedArrays.jl следует импортировать его во все процессы с помощью @everywhere

$ ./julia -p 4

julia> addprocs()

julia> @everywhere using DistributedArrays

julia> using CUDA

julia> B = ones(10_000) ./ 2;

julia> A = ones(10_000) .* π;

julia> C = 2 .* A ./ B;

julia> all(C .≈ 4*π)
true

julia> typeof(C)
Array{Float64,1}

julia> dB = distribute(B);

julia> dA = distribute(A);

julia> dC = 2 .* dA ./ dB;

julia> all(dC .≈ 4*π)
true

julia> typeof(dC)
DistributedArrays.DArray{Float64,1,Array{Float64,1}}

julia> cuB = CuArray(B);

julia> cuA = CuArray(A);

julia> cuC = 2 .* cuA ./ cuB;

julia> all(cuC .≈ 4*π);
true

julia> typeof(cuC)
CuArray{Float64,1}

В следующем примере мы будем использовать DistributedArrays.jl и CUDA.jl для распределения массива по нескольким процессам и вызова на нём общей функции.

function power_method(M, v)
    for i in 1:100
        v = M*v
        v /= norm(v)
    end

    return v, norm(M*v) / norm(v)  # or  (M*v) ./ v
end

power_method многократно создаёт новый вектор и нормализует его. Мы не указали тип данных в объявлении функции, посмотрим, будет ли она работать с упомянутыми типами данных:

julia> M = [2. 1; 1 1];

julia> v = rand(2)
2-element Array{Float64,1}:
0.40395
0.445877

julia> power_method(M,v)
([0.850651, 0.525731], 2.618033988749895)

julia> cuM = CuArray(M);

julia> cuv = CuArray(v);

julia> curesult = power_method(cuM, cuv);

julia> typeof(curesult)
CuArray{Float64,1}

julia> dM = distribute(M);

julia> dv = distribute(v);

julia> dC = power_method(dM, dv);

julia> typeof(dC)
Tuple{DistributedArrays.DArray{Float64,1,Array{Float64,1}},Float64}

Для завершения этого краткого обзора внешних пакетов, рассмотрим MPI.jl, оболочку Julia для протокола MPI. Так как рассмотрение каждой внутренней функции займет слишком много времени, лучше просто оценить подход, используемый для реализации протокола.

Рассмотрим этот пример скрипта, который просто вызывает каждый дочерний процесс, инициализирует его ранг, и когда достигает главного процесса, производит суммирование рангов.

import MPI

MPI.Init()

comm = MPI.COMM_WORLD
MPI.Barrier(comm)

root = 0
r = MPI.Comm_rank(comm)

sr = MPI.Reduce(r, MPI.SUM, root, comm)

if(MPI.Comm_rank(comm) == root)
   @printf("sum of ranks: %s\n", sr)
end

MPI.Finalize()
mpirun -np 4 ./julia example.jl
  • 1В этом контексте MPI относится к стандарту MPI-1. Начиная с MPI-2, комитет по стандартам MPI ввёл новую группу механизмов связи, объединённых под названием удалённый доступ к памяти (RMA). Цель добавления rma в стандарт MPI заключалась в облегчении односторонних шаблонов связи. Для получения дополнительной информации о последнем стандарте MPI см. https://mpi-forum.org/docs.

© 2009–2024 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.10/manual/distributed-computing/

Spec-Zone.ru

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