Параллельное вычисление
Для новичков в многопоточности и параллельном программировании полезно сначала понять различные уровни параллелизма, предлагаемые Julia. Мы можем разделить их на три основные категории:
- Корутины Julia (зеленые потоки)
- Многопоточность
- Многоядерная или распределенная обработка
Сначала мы рассмотрим задачи Julia (задачи (также известные как корутины)) и другие модули, которые полагаются на библиотеку выполнения Julia, что позволяет нам приостанавливать и возобновлять вычисления с полным контролем меж-Tasks коммуникации без необходимости ручного взаимодействия с планировщиком операционной системы. Julia также поддерживает коммуникацию между Tasks с помощью операций, таких как wait и fetch. Коммуникация и синхронизация данных управляются Channel каналами, которые являются каналами для меж-Tasks коммуникации.
Julia также поддерживает экспериментальную многопоточность, где выполнение раздваивается, и анонимная функция выполняется на всех потоках. Известный как подход разветвления-слияния, параллельные потоки выполняют свою работу независимо и должны в конечном итоге быть объединены в главном потоке Julia, чтобы позволить продолжению последовательного выполнения. Многопоточность поддерживается с помощью модуля Base.Threads, который все еще считается экспериментальным, так как Julia еще не полностью потокобезопасна. В частности, сбои, связанные с сегментацией памяти, могут возникать во время операций ввода-вывода и переключения задач. Для получения актуальной справки следите за трекером проблем. Многопоточность следует использовать только в том случае, если вы учитываете глобальные переменные, блокировки и атомарные операции, все из которых будут описаны позже.
В конце мы представим подход Julia к распределенным и параллельным вычислениям. Имея в виду научные вычисления, Julia изначально реализует интерфейсы для распределения процесса по нескольким ядрам или машинам. Также мы упомянем полезные внешние пакеты для распределенного программирования, такие как MPI.jl и DistributedArrays.jl.
Корутины
Платформа параллельного программирования Julia использует задачи (также известные как корутины) для переключения между несколькими вычислениями. Для выражения порядка выполнения между легкими потоками необходимы примитивы коммуникации. Julia предлагает Channel(func::Function, ctype=Any, csize=0, taskref=nothing) , который создает новую задачу из func, связывает ее с новым каналом типа ctype и размером csize и планирует задачу. Channels может служить способом связи между задачами, так как Channel{T}(sz::Int) создает буферизованный канал типа T и размера sz. Всякий раз, когда код выполняет операцию связи, такую как fetch или wait, текущая задача приостанавливается, и планировщик выбирает другую задачу для выполнения. Задача перезапускается, когда событие, которое она ожидает, завершается.
Для многих задач не требуется думать о задачах напрямую. Однако они могут использоваться для одновременного ожидания нескольких событий, что обеспечивает *динамическое планирование*. При динамическом планировании программа сама определяет, что вычислить или где это сделать, в зависимости от того, когда завершатся другие задачи. Это необходимо для непредсказуемых или несбалансированных рабочих нагрузок, когда мы хотим назначить больше работы процессам только после завершения ими текущих задач.
Каналы
Раздел о Task в Потоке управления обсуждал выполнение нескольких функций в кооперативном режиме. Channel каналы могут быть очень полезны для передачи данных между запущенными задачами, особенно для задач, связанных с операциями ввода-вывода.
Примеры операций, связанных с вводом-выводом, включают чтение/запись файлов, доступ к веб-сервисам, выполнение внешних программ и т. д. Во всех этих случаях общее время выполнения может быть улучшено, если другие задачи могут выполняться во время чтения файла или ожидания завершения внешнего сервиса/программы.
Канал можно визуализировать как трубу, т. е. у него есть конец записи и конец чтения:
Несколько писателей в разных задачах могут одновременно писать в один канал с помощью вызовов
put!.Несколько читателей в разных задачах могут одновременно считывать данные с помощью вызовов
take!.-
В качестве примера:
# Given Channels c1 and c2, c1 = Channel(32) c2 = Channel(32) # and a function `foo` which reads items from c1, processes the item read # and writes a result to c2, function foo() while true data = take!(c1) [...] # process data put!(c2, result) # write out result end end # we can schedule `n` instances of `foo` to be active concurrently. for _ in 1:n @async foo() end Каналы создаются с помощью конструктора
Channel{T}(sz). Канал будет содержать только объекты типаT. Если тип не указан, канал может содержать объекты любого типа.szотносится к максимальному количеству элементов, которые могут быть размещены в канале в любой момент. Например,Channel(32)создает канал, который может хранить максимум 32 объекта любого типа. АChannel{MyType}(64)может содержать до 64 объектов типаMyTypeв любой момент.Если
Channelпуст, читатели (при вызовеtake!) будут блокироваться до тех пор, пока данные не станут доступны.Если
Channelзаполнен, писатели (при вызовеput!) будут блокироваться до тех пор, пока не освободится место.isreadyпроверяет наличие каких-либо объектов в канале, аwaitожидает, пока объект не станет доступным.-
Канал
Channelизначально открыт. Это означает, что к нему можно свободно читать и писать с помощью вызововtake!иput!.closeзакрываетChannel. В закрытомChannelканалеput!завершится ошибкой. Например:julia> c = Channel(2); julia> put!(c, 1) # `put!` on an open channel succeeds 1 julia> close(c); julia> put!(c, 2) # `put!` on a closed channel throws an exception. ERROR: InvalidStateException("Channel is closed.",:closed) Stacktrace: [...] -
take!иfetch(получение значения без удаления) в закрытом канале успешно возвращают любые имеющиеся значения до тех пор, пока он не опустеет. Продолжение предыдущего примера:julia> fetch(c) # Any number of `fetch` calls succeed. 1 julia> fetch(c) 1 julia> take!(c) # The first `take!` removes the value. 1 julia> take!(c) # No more data available on a closed channel. ERROR: InvalidStateException("Channel is closed.",:closed) Stacktrace: [...]
Канал Channel может использоваться как итерируемый объект в цикле for, в этом случае цикл выполняется до тех пор, пока в Channel есть данные или канал открыт. Переменная цикла принимает все значения, добавленные в Channel. Цикл for завершается, как только Channel закрыт и опустел.
Например, следующее приведет к тому, что цикл for будет ожидать дополнительных данных:
julia> c = Channel{Int}(10);
julia> foreach(i->put!(c, i), 1:3) # add a few entries
julia> data = [i for i in c]
в то время как это вернет результат после чтения всех данных:
julia> c = Channel{Int}(10);
julia> foreach(i->put!(c, i), 1:3); # add a few entries
julia> close(c); # `for` loops can exit
julia> data = [i for i in c]
3-element Array{Int64,1}:
1
2
3
Рассмотрим простой пример использования каналов для межзадачной связи. Мы запускаем 4 задачи для обработки данных из одного jobs канала. Задачи, идентифицированные по идентификатору (job_id), записываются в канал. Каждая задача в этой симуляции читает job_id, ждет случайное время и записывает обратно кортеж из job_id и моделируемого времени в канал результатов. В конце концов, все results выводятся на экран.
julia> const jobs = Channel{Int}(32);
julia> const results = Channel{Tuple}(32);
julia> function do_work()
for job_id in jobs
exec_time = rand()
sleep(exec_time) # simulates elapsed time doing actual work
# typically performed externally.
put!(results, (job_id, exec_time))
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 i in 1:4 # start 4 tasks to process requests in parallel
@async do_work()
end
julia> @elapsed while n > 0 # print out results
job_id, exec_time = take!(results)
println("$job_id finished in $(round(exec_time; digits=2)) seconds")
global n = n - 1
end
4 finished in 0.22 seconds
3 finished in 0.45 seconds
1 finished in 0.5 seconds
7 finished in 0.14 seconds
2 finished in 0.78 seconds
5 finished in 0.9 seconds
9 finished in 0.36 seconds
6 finished in 0.87 seconds
8 finished in 0.79 seconds
10 finished in 0.64 seconds
12 finished in 0.5 seconds
11 finished in 0.97 seconds
0.029772311
В текущей версии Julia все задачи мультиплексируются на один поток ОС. Таким образом, хотя задачи, связанные с операциями ввода-вывода, выигрывают от параллельного выполнения, задачи, связанные с вычислениями, фактически выполняются последовательно в одном потоке ОС.
Многопоточность (экспериментальная)
Помимо задач Julia изначально поддерживает многопоточность. Обратите внимание, что этот раздел является экспериментальным, и интерфейсы могут измениться в будущем.
Настройка
По умолчанию Julia запускается с одним потоком выполнения. Это можно проверить, используя команду Threads.nthreads():
julia> Threads.nthreads() 1
Количество потоков, с которыми запускается Julia, контролируется переменной среды JULIA_NUM_THREADS. Теперь давайте запустим Julia с 4 потоками:
Bash на Linux/OSX:
export JULIA_NUM_THREADS=4
C shell на Linux/OSX, CMD на Windows:
set JULIA_NUM_THREADS=4
Powershell на Windows:
$env:JULIA_NUM_THREADS=4
Давайте проверим, что у нас есть 4 потока.
julia> Threads.nthreads() 4
Но мы сейчас в главном потоке. Чтобы проверить это, воспользуемся функцией Threads.threadid
julia> Threads.threadid() 1
Макрос @threads
Давайте рассмотрим простой пример с использованием наших потоков. Создадим массив нулей:
julia> a = zeros(10)
10-element Array{Float64,1}:
0.0
0.0
0.0
0.0
0.0
0.0
0.0
0.0
0.0
0.0
Давайте выполним операции над этим массивом одновременно с помощью 4 потоков. Каждый поток запишет свой идентификатор потока в каждую позицию.
Julia поддерживает параллельные циклы с помощью макроса Threads.@threads. Этот макрос ставится перед циклом for для обозначения Julia, что этот цикл является многопоточным:
julia> Threads.@threads for i = 1:10
a[i] = Threads.threadid()
end
Пространство итераций разбивается между потоками, после чего каждый поток записывает свой идентификатор потока в свои назначенные позиции:
julia> a
10-element Array{Float64,1}:
1.0
1.0
1.0
2.0
2.0
2.0
3.0
3.0
4.0
4.0
Обратите внимание, что у Threads.@threads нет необязательного параметра редукции, как у @distributed.
Атомарные операции
Julia поддерживает доступ к значениям и их изменение *атомарно*, то есть безопасным способом для предотвращения гонок. Значение (которое должно быть примитивного типа) можно обернуть как Threads.Atomic для указания, что к нему нужно обращаться таким образом. Вот пример:
julia> i = Threads.Atomic{Int}(0);
julia> ids = zeros(4);
julia> old_is = zeros(4);
julia> Threads.@threads for id in 1:4
old_is[id] = Threads.atomic_add!(i, id)
ids[id] = id
end
julia> old_is
4-element Array{Float64,1}:
0.0
1.0
7.0
3.0
julia> ids
4-element Array{Float64,1}:
1.0
2.0
3.0
4.0
Если бы мы попытались выполнить сложение без тега атомарности, мы могли бы получить неправильный результат из-за гонки. Вот пример того, что произошло бы, если бы мы не избежали гонки:
julia> using Base.Threads
julia> nthreads()
4
julia> acc = Ref(0)
Base.RefValue{Int64}(0)
julia> @threads for i in 1:1000
acc[] += 1
end
julia> acc[]
926
julia> acc = Atomic{Int64}(0)
Atomic{Int64}(0)
julia> @threads for i in 1:1000
atomic_add!(acc, 1)
end
julia> acc[]
1000
Не все примитивные типы можно обернуть в тэг Atomic. Поддерживаемые типы — Int8, Int16, Int32, Int64, Int128, UInt8, UInt16, UInt32, UInt64, UInt128, Float16, Float32, и Float64. Кроме того, Int128 и UInt128 не поддерживаются на AAarch32 и ppc64le.
Побочные эффекты и изменяемые аргументы функций
При использовании многопоточности нужно быть осторожными при работе с функциями, которые не являются чистыми, так как мы можем получить неправильный ответ. Например, функции, имена которых по соглашению оканчиваются на ! , изменяют свои аргументы и, следовательно, не являются чистыми. Однако существуют функции, имеющие побочные эффекты, и их имена не оканчиваются на !. Например, findfirst(regex, str) изменяет свой аргумент regex или rand() изменяет Base.GLOBAL_RNG:
julia> using Base.Threads
julia> nthreads()
4
julia> function f()
s = repeat(["123", "213", "231"], outer=1000)
x = similar(s, Int)
rx = r"1"
@threads for i in 1:3000
x[i] = findfirst(rx, s[i]).start
end
count(v -> v == 1, x)
end
f (generic function with 1 method)
julia> f() # the correct result is 1000
1017
julia> function g()
a = zeros(1000)
@threads for i in 1:1000
a[i] = rand()
end
length(unique(a))
end
g (generic function with 1 method)
julia> Random.seed!(1); g() # the result for a single thread is 1000
781
В таких случаях необходимо переработать код, чтобы избежать возможности гонки, или использовать примитивы синхронизации.
Например, для исправления примера findfirst выше необходимо иметь отдельный экземпляр переменной rx для каждого потока:
julia> function f_fix()
s = repeat(["123", "213", "231"], outer=1000)
x = similar(s, Int)
rx = [Regex("1") for i in 1:nthreads()]
@threads for i in 1:3000
x[i] = findfirst(rx[threadid()], s[i]).start
end
count(v -> v == 1, x)
end
f_fix (generic function with 1 method)
julia> f_fix()
1000
Теперь мы используем Regex("1") вместо r"1" , чтобы убедиться, что Julia создает отдельные экземпляры объекта Regex для каждого элемента вектора rx.
Случай rand немного сложнее, так как нам нужно гарантировать, что каждый поток использует непересекающиеся последовательности псевдослучайных чисел. Это можно просто обеспечить, используя функцию Future.randjump:
julia> using Random; import Future
julia> function g_fix(r)
a = zeros(1000)
@threads for i in 1:1000
a[i] = rand(r[threadid()])
end
length(unique(a))
end
g_fix (generic function with 1 method)
julia> r = let m = MersenneTwister(1)
[m; accumulate(Future.randjump, fill(big(10)^20, nthreads()-1), init=m)]
end;
julia> g_fix(r)
1000
Мы передаём вектор r в g_fix , так как генерация нескольких ПСЧ — дорогостоящая операция, и мы не хотим её повторять каждый раз, когда запускаем функцию.
@threadcall (Экспериментальная)
Все задачи ввода/вывода, таймеры, команды REPL и т. д. мультиплексируются на один поток ОС через цикл событий. Модифицированная версия libuv (http://docs.libuv.org/en/v1.x/) предоставляет эту функциональность. Точки приостановки обеспечивают кооперативное планирование нескольких задач в одном потоке ОС. Задачи ввода/вывода и таймеры приостанавливаются неявно, ожидая выполнения события. Вызов yield явно позволяет планировать другие задачи.
Таким образом, задача, выполняющая ccall , фактически предотвращает планировщик Julia от выполнения других задач до возврата вызова. Это верно для всех вызовов внешних библиотек. Исключения составляют вызовы в пользовательский код C, который вызывает обратно в Julia (что может затем приостановить) или код C, который вызывает jl_yield() (C-эквивалент yield).
Обратите внимание, что, хотя код Julia выполняется в одном потоке (по умолчанию), библиотеки, используемые Julia, могут запускать свои внутренние потоки. Например, библиотека BLAS может запустить столько потоков, сколько ядер на машине.
Макрос @threadcall решает сценарии, в которых мы не хотим, чтобы ccall блокировал основной цикл событий Julia. Он планирует функцию C для выполнения в отдельном потоке. Для этого используется пул потоков по умолчанию размером 4. Размер пула потоков контролируется переменной окружения UV_THREADPOOL_SIZE. Во время ожидания свободного потока и во время выполнения функции, как только поток доступен, запрашивающая задача (в основном цикле событий Julia) приостанавливается для других задач. Обратите внимание, что @threadcall не возвращается, пока выполнение не завершится. С точки зрения пользователя, это, следовательно, блокирующий вызов, как и другие API Julia.
Очень важно, чтобы вызываемая функция не вызывала обратно в Julia, так как это приведёт к ошибке.
@threadcall может быть удалена/изменена в будущих версиях Julia.
Многоядерная или распределённая обработка
Реализация параллельного вычисления с распределённой памятью предоставляется модулем Distributed в качестве части стандартной библиотеки Julia.
Большинство современных компьютеров имеют более одного процессора, а несколько компьютеров могут быть объединены в кластер. Использование мощности нескольких процессоров позволяет быстрее выполнять многие вычисления. Существуют два основных фактора, влияющих на производительность: скорость самих процессоров и скорость доступа к памяти. В кластере очевидно, что данный процессор будет иметь наибольшую скорость доступа к ОЗУ на том же компьютере (узле). Возможно, более удивительно, что аналогичные проблемы актуальны на типичном многоядерном ноутбуке из-за различий в скорости основной памяти и кэша. Следовательно, хорошая среда для многопроцессорной обработки должна позволять управлять «владением» частью памяти конкретным процессором. Julia предоставляет многопроцессорную среду, основанную на передаче сообщений, чтобы позволить программам работать на нескольких процессах в отдельных областях памяти одновременно.
Реализация Julia передачи сообщений отличается от других сред, таких как MPI [1]. Общение в Julia обычно является «односторонним», что означает, что программист должен явно управлять только одним процессом в операции с двумя процессами. Кроме того, эти операции обычно не выглядят как «отправка сообщения» и «приём сообщения», а скорее похожи на операции более высокого уровня, такие как вызовы пользовательских функций.
Распределённое программирование в Julia основано на двух примитивах: удалённых ссылках и удалённых вызовах. Удалённая ссылка — это объект, который может быть использован любым процессом для ссылки на объект, хранящийся на конкретном процессе. Удалённый вызов — это запрос от одного процесса вызвать определённую функцию с определёнными аргументами на другом (возможно, на том же) процессе.
Удалённые ссылки бывают двух типов: Future и RemoteChannel.
Удалённый вызов возвращает Future своего результата. Удалённые вызовы возвращаются немедленно; процесс, который сделал вызов, переходит к следующей операции, в то время как удалённый вызов происходит где-то ещё. Вы можете дождаться завершения удалённого вызова, вызвав wait на возвращённой Future, и вы можете получить полное значение результата, используя fetch.
С другой стороны, RemoteChannel — это перезаписываемые объекты. Например, несколько процессов могут координировать свою обработку, ссылаясь на тот же удалённый Channel.
Каждый процесс имеет связанный идентификатор. Процесс, предоставляющий интерактивный приглашение Julia, всегда имеет идентификатор id равный 1. Процессы, используемые по умолчанию для параллельных операций, называются «рабочими». Когда есть только один процесс, процесс 1 считается рабочим. В противном случае рабочими считаются все процессы, кроме процесса 1.
Давайте попробуем это. Начало с julia -p n предоставляет n рабочих процессов на локальном компьютере. Как правило, имеет смысл, чтобы n было равно числу потоков ЦП (логических ядер) на машине. Обратите внимание, что аргумент -p неявно загружает модуль Distributed.
$ ./julia -p 2
julia> r = remotecall(rand, 2, 2, 2)
Future(2, 1, 4, nothing)
julia> s = @spawnat 2 1 .+ fetch(r)
Future(2, 1, 5, nothing)
julia> fetch(s)
2×2 Array{Float64,2}:
1.18526 1.50912
1.16296 1.60607
Первый аргумент для remotecall — это функция для вызова. Большинство параллельных программ в Julia не ссылаются на конкретные процессы или количество доступных процессов, но remotecall рассматривается как интерфейс низкого уровня, обеспечивающий более точный контроль. Второй аргумент для remotecall — это id процесса, который будет выполнять работу, а оставшиеся аргументы будут переданы вызываемой функции.
Как вы видите, в первой строке мы попросили процесс 2 создать 2×2 случайную матрицу, а во второй строке мы попросили добавить к ней 1. Результаты обоих вычислений доступны в двух фьючерсах, r и s.
Макрос @spawnat вычисляет выражение во втором аргументе на процессе, указанном в первом аргументе.
Иногда вам может потребоваться значение, вычисленное удалённо, немедленно. Это обычно происходит, когда вы читаете удалённый объект, чтобы получить данные, необходимые для следующей локальной операции. Функция remotecall_fetch предназначена для этой цели. Она эквивалентна fetch(remotecall(...)) , но более эффективна.
julia> remotecall_fetch(getindex, 2, r, 1, 1) 0.18526337335308085
Помните, что getindex(r,1,1) эквивалентно r[1,1], поэтому этот вызов извлекает первый элемент фьючерса r.
Для упрощения, символ :any может быть передан к [@spawnat], который выбирает место выполнения операции за вас:
julia> r = @spawnat :any rand(2,2)
Future(2, 1, 4, nothing)
julia> s = @spawnat :any 1 .+ fetch(r)
Future(3, 1, 5, nothing)
julia> fetch(s)
2×2 Array{Float64,2}:
1.38854 1.9098
1.20939 1.57158
Обратите внимание, что мы использовали 1 .+ fetch(r) вместо 1 .+ r . Это потому, что мы не знаем, где будет выполняться код, поэтому в общем случае может потребоваться fetch для перемещения r в процесс, выполняющий сложение. В данном случае @spawnat достаточно умён, чтобы выполнить вычисление на процессе, который владеет r, поэтому fetch будет «безоперационной» операцией (работа не выполняется).
(Стоит отметить, что @spawnat не является встроенным, а определён в Julia как макрос. Возможно, определить свои собственные такие конструкции.)
Важно помнить, что, после получения, Future будет кэшировать своё значение локально. Дополнительные вызовы fetch не предполагают сетевого обмена. После того, как все ссылающиеся Future получат данные, удалённое сохранённое значение удаляется.
@async похож на @spawnat, но выполняет задачи только на локальном процессе. Мы используем его для создания задачи «подачи» для каждого процесса. Каждая задача выбирает следующий индекс, который нужно вычислить, затем ждёт завершения своего процесса, а затем повторяет это до тех пор, пока не закончатся индексы. Обратите внимание, что задачи «подачи» не начинают выполняться до тех пор, пока основная задача не достигнет конца блока @sync, в этот момент она уступает управление и ждёт завершения всех локальных задач, прежде чем возвратиться из функции. Начиная с версии 0.7 и далее, задачи «подачи» могут обмениваться состоянием через nextidx, поскольку все они выполняются на одном процессе. Даже если Tasks планируются кооперативно, блокировка всё равно может потребоваться в некоторых контекстах, как в случае с асинхронным вводом-выводом. Это означает, что переключения контекста происходят только в определённые моменты: в данном случае, когда вызывается remotecall_fetch. Это текущее состояние реализации, и оно может измениться в будущих версиях Julia, поскольку предполагается, что это позволит запускать до N Tasks на M Process, то есть многопоточность M:N. Затем потребуется модель получения/освобождения блокировки для nextidx , так как небезопасно позволять нескольким процессам одновременно читать и записывать ресурс.
Доступность кода и загрузка пакетов
Ваш код должен быть доступен на любом процессе, который его выполняет. Например, введите следующее в командной строке Julia:
julia> function rand2(dims...)
return 2*rand(dims...)
end
julia> rand2(2,2)
2×2 Array{Float64,2}:
0.153756 0.368514
1.15119 0.918912
julia> fetch(@spawnat :any rand2(2,2))
ERROR: RemoteException(2, CapturedException(UndefVarError(Symbol("#rand2"))
Stacktrace:
[...]
Процесс 1 знал о функции rand2, но процесс 2 — нет.
Чаще всего вы будете загружать код из файлов или пакетов, и у вас есть значительная гибкость в управлении тем, какие процессы загружают код. Рассмотрим файл DummyModule.jl, содержащий следующий код:
module DummyModule
export MyType, f
mutable struct MyType
a::Int
end
f(x) = x^2+1
println("loaded")
end
Для того, чтобы сослаться на MyType на всех процессах, DummyModule.jl необходимо загрузить на каждом процессе. Вызов include("DummyModule.jl") загружает его только на одном процессе. Для загрузки на каждом процессе используйте макрос @everywhere (запуска Julia с julia -p 2):
julia> @everywhere include("DummyModule.jl")
loaded
From worker 3: loaded
From worker 2: loaded
Как обычно, это не добавляет DummyModule в область видимости ни на одном из процессов, для этого требуется using или import. Кроме того, когда DummyModule входит в область видимости одного процесса, она не присутствует ни в каком другом:
julia> using .DummyModule julia> MyType(7) MyType(7) julia> fetch(@spawnat 2 MyType(7)) ERROR: On worker 2: UndefVarError: MyType not defined ⋮ julia> fetch(@spawnat 2 DummyModule.MyType(7)) MyType(7)
Однако всё ещё возможно, например, отправить MyType процессу, который загрузил DummyModule , даже если он не в области видимости:
julia> put!(RemoteChannel(2), MyType(7))
RemoteChannel{Channel{Any}}(2, 1, 13)
Файл также можно предварительно загрузить на несколько процессов при запуске с флагом -L, и для управления вычислениями можно использовать скрипт драйвера:
julia -p <n> -L file1.jl -L file2.jl driver.jl
Процесс Julia, выполняющий скрипт драйвера в приведённом примере, имеет id равный 1, как и процесс, предоставляющий интерактивную командную строку.
Наконец, если DummyModule.jl не является самостоятельным файлом, а пакетом, то using DummyModule загрузит DummyModule.jl на все процессы, но добавит его в область видимости только на том процессе, где был вызван using.
Запуск и управление рабочими процессами
Базовая установка Julia имеет встроенную поддержку двух типов кластеров:
- Локальный кластер, указанный параметром
-p, как показано выше. - Кластер, охватывающий несколько машин, с помощью параметра
--machine-file. Это использует бесклеточныеsshдля запуска процессов Julia-рабочих (из той же директории, что и текущий хост) на указанных машинах.
Функции addprocs, rmprocs, workers и другие доступны как программно для добавления, удаления и запроса процессов в кластере.
julia> using Distributed
julia> addprocs(2)
2-element Array{Int64,1}:
2
3
Модуль Distributed должен быть явно загружен на главном процессе перед вызовом addprocs. Он автоматически становится доступным на рабочих процессах.
Обратите внимание, что рабочие процессы не запускают скрипт запуска ~/.julia/config/startup.jl, и они не синхронизируют своё глобальное состояние (такое как глобальные переменные, новые определения методов и загруженные модули) ни с одним из других запущенных процессов. Вы можете использовать addprocs(exeflags="--project") для инициализации рабочего процесса с определённой средой, а затем @everywhere using <modulename> или @everywhere include("file.jl").
Другие типы кластеров могут быть поддерживаемы путём написания пользовательских ClusterManager, как описано ниже в разделе ClusterManagers.
Перемещение данных
Отправка сообщений и перемещение данных составляют большую часть накладных расходов в распределённой программе. Снижение количества сообщений и объёма отправленных данных имеет решающее значение для достижения производительности и масштабируемости. Для этого важно понять перемещение данных, выполняемое различными конструкциями распределённого программирования Julia.
fetch можно рассматривать как явную операцию перемещения данных, так как она напрямую запрашивает перемещение объекта на локальную машину. @spawnat (и некоторые связанные конструкции) также перемещает данные, но это не так очевидно, поэтому его можно назвать неявной операцией перемещения данных. Рассмотрим эти два подхода к построению и возведению в квадрат случайной матрицы:
Метод 1:
julia> A = rand(1000,1000); julia> Bref = @spawnat :any A^2; [...] julia> fetch(Bref);
Метод 2:
julia> Bref = @spawnat :any rand(1000,1000)^2; [...] julia> fetch(Bref);
Разница кажется незначительной, но на самом деле она довольно существенна из-за поведения @spawnat. В первом методе случайная матрица создаётся локально, затем отправляется в другой процесс, где она возводится в квадрат. Во втором методе случайная матрица создаётся и возводится в квадрат в другом процессе. Поэтому во втором методе отправляется гораздо меньше данных, чем в первом.
В этом примере два метода легко отличить и выбрать. Однако в реальной программе проектирование перемещения данных может потребовать больше размышлений и, вероятно, некоторых измерений. Например, если первому процессу нужна матрица A, то первый метод может быть предпочтительнее. Или, если вычисление A дорогостоящее и только текущий процесс его имеет, то перемещение его в другой процесс может быть неизбежным. Или, если текущий процесс имеет очень мало задач между @spawnat и fetch(Bref), может быть лучше полностью исключить параллелизм. Или представьте, что rand(1000,1000) заменено более сложной операцией. Тогда может иметь смысл добавить ещё одну инструкцию @spawnat только для этой стадии.
Глобальные переменные
Выражения, выполняемые удалённо через @spawnat, или замыкания, определённые для удалённого выполнения с помощью remotecall, могут ссылаться на глобальные переменные. Глобальные привязки в модуле Main обрабатываются немного иначе по сравнению с глобальными привязками в других модулях. Рассмотрим следующий фрагмент кода:
A = rand(10,10) remotecall_fetch(()->sum(A), 2)
В этом случае sum ДОЛЖЕН быть определён в удалённом процессе. Обратите внимание, что A — это глобальная переменная, определённая в локальном рабочем пространстве. У рабочего процесса 2 нет переменной под названием A в Main. Действие отправки замыкания ()->sum(A) рабочему 2 приводит к определению Main.A на 2. Main.A продолжает существовать на рабочем 2 даже после возвращения вызова remotecall_fetch. Удалённые вызовы с вложенными глобальными ссылками (только в модуле Main ) управляют глобальными переменными следующим образом:
Новые глобальные привязки создаются на целевых рабочих процессах, если они ссылаются в рамках удалённого вызова.
Глобальные константы объявляются как константы и на удалённых узлах.
-
Глобальные переменные передаются на целевой рабочий процесс только в контексте удалённого вызова и только если их значение изменилось. Также кластер не синхронизирует глобальные привязки между узлами. Например:
A = rand(10,10) remotecall_fetch(()->sum(A), 2) # worker 2 A = rand(10,10) remotecall_fetch(()->sum(A), 3) # worker 3 A = nothing
Выполнение вышеприведённого фрагмента кода приводит к тому, что
Main.Aна рабочем процессе 2 имеет другое значение, чемMain.Aна рабочем процессе 3, в то время как значениеMain.Aна узле 1 установлено вnothing.
Как вы могли заметить, в то время как память, связанная с глобальными переменными, может быть собрана, когда они переопределяются на главном процессе, аналогичных действий не предпринимается на рабочих процессах, так как привязки остаются действительными. clear! можно использовать для ручного переопределения конкретных глобальных переменных на удалённых узлах в nothing после того, как они больше не нужны. Это высвободит связанную с ними память в рамках обычного цикла сборки мусора.
Таким образом, программы должны быть осторожны при обращении к глобальным переменным в удалённых вызовах. Фактически, если это возможно, следует избегать этого вообще. Если вам необходимо обратиться к глобальным переменным, рассмотрите использование блоков let для локализации глобальных переменных.
Например:
julia> A = rand(10,10);
julia> remotecall_fetch(()->A, 2);
julia> B = rand(10,10);
julia> let B = B
remotecall_fetch(()->B, 2)
end;
julia> @fetchfrom 2 InteractiveUtils.varinfo()
name size summary
––––––––– ––––––––– ––––––––––––––––––––––
A 800 bytes 10×10 Array{Float64,2}
Base Module
Core Module
Main Module
Как видно, глобальная переменная A определена на рабочем процессе 2, но B захвачена как локальная переменная, и поэтому привязка для B не существует на рабочем процессе 2.
Параллельные отображения и циклы
К счастью, многие полезные параллельные вычисления не требуют перемещения данных. Примером может служить симуляция Монте-Карло, где несколько процессов могут обрабатывать независимые испытания симуляции одновременно. Мы можем использовать @spawnat для подбрасывания монет на двух процессах. Сначала напишите следующую функцию в count_heads.jl:
function count_heads(n)
c::Int = 0
for i = 1:n
c += rand(Bool)
end
c
end
Функция count_heads просто складывает вместе n случайные биты. Вот как мы можем провести несколько испытаний на двух машинах и сложить результаты:
julia> @everywhere include_string(Main, $(read("count_heads.jl", String)), "count_heads.jl")
julia> a = @spawnat :any count_heads(100000000)
Future(2, 1, 6, nothing)
julia> b = @spawnat :any count_heads(100000000)
Future(3, 1, 7, nothing)
julia> fetch(a)+fetch(b)
100001564
Этот пример демонстрирует мощный и часто используемый шаблон параллельного программирования. Много итераций выполняются независимо на нескольких процессах, а затем их результаты комбинируются с использованием некоторой функции. Процесс комбинирования называется редукцией, так как он обычно сокращает ранг тензора: вектор чисел сводится к одному числу, или матрица сводится к одной строке или столбцу и т. д. В коде это обычно выглядит как шаблон x = f(x,v[i]), где x является аккумулятором, f — функцией редукции, а v[i] — элементами, которые сокращаются. Желательно, чтобы f была ассоциативной, чтобы порядок выполнения операций не имел значения.
Обратите внимание, что наше использование этого шаблона с count_heads может быть обобщено. Мы использовали два явных оператора @spawnat, который ограничивает параллелизм двумя процессами. Чтобы запустить на любом количестве процессов, мы можем использовать параллельный цикл for, работающий в распределенной памяти, который можно записать в Julia с помощью @distributed так:
nheads = @distributed (+) for i = 1:200000000
Int(rand(Bool))
end
Этот конструкт реализует шаблон назначения итераций нескольким процессам и объединяет их с заданной редукцией (в данном случае (+)). Результат каждой итерации принимается как значение последнего выражения внутри цикла. Само выражение параллельного цикла оценивается в конечный ответ.
Обратите внимание, что, хотя параллельные циклы for выглядят как последовательные циклы for, их поведение существенно отличается. В частности, итерации не происходят в определенном порядке, а записи в переменные или массивы не будут глобально видимы, поскольку итерации выполняются на разных процессах. Любые переменные, используемые внутри параллельного цикла, будут скопированы и распространены на каждый процесс.
Например, следующий код не будет работать как ожидается:
a = zeros(100000)
@distributed for i = 1:100000
a[i] = i
end
Этот код не инициализирует все a, так как каждый процесс будет иметь отдельную копию. Параллельные циклы for, подобные этим, следует избегать. К счастью, Обменные массивы могут быть использованы, чтобы обойти это ограничение:
using SharedArrays
a = SharedArray{Float64}(10)
@distributed for i = 1:10
a[i] = i
end
Использование "внешних" переменных в параллельных циклах вполне приемлемо, если переменные являются только для чтения:
a = randn(1000)
@distributed (+) for i = 1:100000
f(a[rand(1:end)])
end
Здесь каждая итерация применяет f к случайно выбранному образцу из вектора a, общий для всех процессов.
Как вы могли видеть, оператор редукции можно опустить, если он не нужен. В этом случае цикл выполняется асинхронно, т. е. он запускает независимые задачи на всех доступных рабочих узлах и возвращает массив Future немедленно, не ожидая завершения. Вызывающий процесс может подождать завершения Future позже, вызвав fetch на них, или подождать завершения в конце цикла, добавив перед ним @sync, как @sync @distributed for.
В некоторых случаях оператор редукции не нужен, и мы просто хотим применить функцию ко всем целым числам в определённом диапазоне (или, более общо, ко всем элементам в некотором множестве). Это ещё одна полезная операция, называемая параллельным отображением, реализованная в Julia как функция pmap. Например, мы можем вычислить сингулярные значения нескольких больших случайных матриц параллельно следующим образом:
julia> M = Matrix{Float64}[rand(1000,1000) for i = 1:10];
julia> pmap(svdvals, M);
Функция Julia's pmap предназначена для случая, когда каждый вызов функции выполняет большой объём работы. В отличие от этого, @distributed for может обрабатывать ситуации, когда каждая итерация незначительная, возможно, просто суммирует два числа. Только рабочие процессы используются как pmap, так и @distributed for для параллельного вычисления. В случае @distributed for, окончательная редукция выполняется на вызывающем процессе.
Удаленные ссылки и абстрактные каналы
Удаленные ссылки всегда ссылаются на реализацию AbstractChannel.
Необходима конкретная реализация AbstractChannel (например, Channel), чтобы реализовать put!, take!, fetch, isready и wait. Удаленный объект, на который ссылается Future, хранится в Channel{Any}(1), то есть в Channel размера 1, способном хранить объекты типа Any.
RemoteChannel, который может быть переписан, может указывать на любой тип и размер каналов или на любую другую реализацию AbstractChannel.
Конструктор RemoteChannel(f::Function, pid)() позволяет нам создавать ссылки на каналы, содержащие более одного значения определённого типа. f — функция, выполняемая на pid, и она должна возвращать AbstractChannel.
Например, RemoteChannel(()->Channel{Int}(10), pid), вернёт ссылку на канал типа Int и размера 10. Канал существует на рабочем узле pid.
Методы put!, take!, fetch, isready и wait на RemoteChannel делегируются на хранилище на удалённом процессе.
RemoteChannel может поэтому использоваться для ссылки на объекты AbstractChannel, реализованные пользователем. Простой пример этого приведен в dictchannel.jl в репозитории Примеры, который использует словарь как своё удалённое хранилище.
Каналы и RemoteChannels
- Канал
Channelлокален для процесса. Рабочий узел 2 не может напрямую ссылаться наChannelна рабочем узле 3 и наоборот. Однако,RemoteChannelможет передавать и получать значения между рабочими узлами. RemoteChannelможно рассматривать как указатель наChannel.- Идентификатор процесса,
pid, связанный сRemoteChannel, определяет процесс, где хранится хранилище, то есть базовыйChannel. - Любой процесс, имеющий ссылку на
RemoteChannel, может передавать и получать элементы из канала. Данные автоматически отправляются (или извлекаются) в процесс, с которым связанRemoteChannel. - Сериализация
Channelтакже сериализует любые данные, присутствующие в канале. Десериализация, следовательно, создаёт копию исходного объекта. - С другой стороны, сериализация
RemoteChannelвключает только сериализацию идентификатора, который определяет местоположение и экземплярChannel, на который указывает указатель. Десериализованный объектRemoteChannel(на любом рабочем узле), таким образом, также указывает на то же хранилище, что и исходный.
Пример с каналами из вышеуказанного текста можно изменить для межпроцессного обмена данными, как показано ниже.
Мы запускаем 4 рабочих узла для обработки одного jobs удалённого канала. Задачи, идентифицируемые по id (job_id), записываются в канал. Каждая удалённо выполняемая задача в этом моделировании считывает job_id, ждёт случайное время и записывает обратно кортеж из job_id, затраченного времени и собственного pid в канал результатов. Наконец, все results выводятся на главном процессе.
julia> addprocs(4); # add worker processes
julia> const jobs = RemoteChannel(()->Channel{Int}(32));
julia> const results = RemoteChannel(()->Channel{Tuple}(32));
julia> @everywhere function do_work(jobs, results) # define work function everywhere
while true
job_id = take!(jobs)
exec_time = rand()
sleep(exec_time) # simulates elapsed time doing actual work
put!(results, (job_id, exec_time, myid()))
end
end
julia> function make_jobs(n)
for i in 1:n
put!(jobs, i)
end
end;
julia> n = 12;
julia> @async make_jobs(n); # feed the jobs channel with "n" jobs
julia> for p in workers() # start tasks on the workers to process requests in parallel
remote_do(do_work, p, jobs, results)
end
julia> @elapsed while n > 0 # print out results
job_id, exec_time, where = take!(results)
println("$job_id finished in $(round(exec_time; digits=2)) seconds on worker $where")
global n = n - 1
end
1 finished in 0.18 seconds on worker 4
2 finished in 0.26 seconds on worker 5
6 finished in 0.12 seconds on worker 4
7 finished in 0.18 seconds on worker 4
5 finished in 0.35 seconds on worker 5
4 finished in 0.68 seconds on worker 2
3 finished in 0.73 seconds on worker 3
11 finished in 0.01 seconds on worker 3
12 finished in 0.02 seconds on worker 3
9 finished in 0.26 seconds on worker 5
8 finished in 0.57 seconds on worker 4
10 finished in 0.58 seconds on worker 2
0.055971741
Удалённые ссылки и распределённое сбор мусора
Объекты, на которые ссылаются удалённые ссылки, могут быть освобождены только тогда, когда все удерживаемые ссылки в кластере удаляются.
Узел, где хранится значение, отслеживает, у каких рабочих узлов есть ссылка на него. Каждый раз, когда RemoteChannel или (неизвлечённый) Future сериализуется на рабочий узел, узел, на который указывает ссылка, уведомляется. И каждый раз, когда RemoteChannel или (неизвлечённый) Future собирается мусор локально, узел, владеющий значением, снова уведомляется. Это реализуется с помощью внутреннего кластера, который знает сериализатор. Удаленные ссылки действительны только в контексте работающего кластера. Сериализация и десериализация ссылок на обычные IO объекты не поддерживается.
Уведомления производятся посредством отправки сообщений отслеживания — сообщения "добавить ссылку" при сериализации ссылки на другой процесс и сообщение "удалить ссылку" при локальном сборе мусора ссылки.
Поскольку Future записываются один раз и кэшируются локально, действие fetch объекта Future также обновляет информацию о слежении за ссылками на узле, владеющем значением.
Узел, который владеет значением, освобождает его, как только все ссылки на него будут удалены.
С Future сериализация уже полученного Future на другой узел также отправляет значение, поскольку исходное удалённое хранилище может уже получить значение к этому моменту.
Важно отметить, что время, когда объект собирается сборщиком мусора локально, зависит от размера объекта и текущей нагрузки на память в системе.
В случае удалённых ссылок размер локального объекта ссылки невелик, в то время как значение, хранящееся на удалённом узле, может быть достаточно большим. Поскольку локальный объект может не быть собран немедленно, рекомендуется явно вызывать finalize для локальных экземпляров RemoteChannel или для неполученных Future. Поскольку вызов fetch для Future также удаляет его ссылку из удалённого хранилища, это не требуется для полученных Future. Явный вызов finalize приводит к немедленной отправке сообщения на удалённый узел для удаления ссылки на значение.
После завершения ссылка становится недействительной и не может быть использована в последующих вызовах.
Локальные вызовы
Данные обязательно копируются на удалённый узел для выполнения. Это относится как к удалённым вызовам, так и к случаям, когда данные сохраняются в RemoteChannel / Future на другом узле. Как ожидалось, это приводит к созданию копии сериализованных объектов на удалённом узле. Однако, когда целевой узел является локальным узлом, т.е. идентификатор вызывающего процесса совпадает с идентификатором удалённого узла, он выполняется как локальный вызов. Обычно (но не всегда) он выполняется в другой задаче — но нет сериализации/десериализации данных. Вследствие этого вызов ссылается на те же экземпляры объектов, что и переданные — не создаются копии. Это поведение показано ниже:
julia> using Distributed;
julia> rc = RemoteChannel(()->Channel(3)); # RemoteChannel created on local node
julia> v = [0];
julia> for i in 1:3
v[1] = i # Reusing `v`
put!(rc, v)
end;
julia> result = [take!(rc) for _ in 1:3];
julia> println(result);
Array{Int64,1}[[3], [3], [3]]
julia> println("Num Unique objects : ", length(unique(map(objectid, result))));
Num Unique objects : 1
julia> addprocs(1);
julia> rc = RemoteChannel(()->Channel(3), workers()[1]); # RemoteChannel created on remote node
julia> v = [0];
julia> for i in 1:3
v[1] = i
put!(rc, v)
end;
julia> result = [take!(rc) for _ in 1:3];
julia> println(result);
Array{Int64,1}[[1], [2], [3]]
julia> println("Num Unique objects : ", length(unique(map(objectid, result))));
Num Unique objects : 3
Как видно, put! локально владеящим RemoteChannel, при котором объект v изменяется между вызовами, приводит к хранению того же единственного экземпляра объекта. В отличие от создания копий v, когда узел, владеющий rc, — это другой узел.
Следует отметить, что это обычно не проблема. Это нужно учитывать только в том случае, если объект хранится локально и модифицируется после вызова. В таких случаях может быть уместно сохранить deepcopy объекта.
Это также верно для удалённых вызовов на локальном узле, как показано в следующем примере:
julia> using Distributed; addprocs(1);
julia> v = [0];
julia> v2 = remotecall_fetch(x->(x[1] = 1; x), myid(), v); # Executed on local node
julia> println("v=$v, v2=$v2, ", v === v2);
v=[1], v2=[1], true
julia> v = [0];
julia> v2 = remotecall_fetch(x->(x[1] = 1; x), workers()[1], v); # Executed on remote node
julia> println("v=$v, v2=$v2, ", v === v2);
v=[0], v2=[1], false
Как можно видеть ещё раз, удалённый вызов на локальный узел ведет себя так же, как прямой вызов. Вызов изменяет локальные объекты, переданные в качестве аргументов. В удалённом вызове он работает с копией аргументов.
Повторим, что в целом это не проблема. Если локальный узел также используется как вычислительный узел и аргументы используются после вызова, это поведение необходимо учесть, и при необходимости следует передавать глубокие копии аргументов в вызов, вызванный на локальном узле. Вызовы на удалённые узлы всегда будут работать с копиями аргументов.
Общие массивы
Общие массивы используют общую системную память для отображения одного и того же массива во многих процессах. Несмотря на некоторые сходства с DArray, поведение SharedArray отличается. В DArray каждый процесс имеет локальный доступ только к фрагменту данных, и ни два процесса не делят один и тот же фрагмент; в отличие от SharedArray, каждый «участвующий» процесс имеет доступ ко всему массиву. SharedArray — хороший выбор, когда необходимо иметь большой объём данных, совместно доступный двум или более процессам на одном компьютере.
Поддержка общих массивов доступна через модуль SharedArrays, который необходимо явно загрузить на все участвующие рабочие узлы.
SharedArray индексация (присвоение и доступ к значениям) работает так же, как и с обычными массивами, и является эффективной, поскольку базовая память доступна локальному процессу. Поэтому большинство алгоритмов естественно работают с SharedArray, хотя и в режиме одного процесса. В случаях, когда алгоритм требует ввода Array, базовый массив можно извлечь из SharedArray, вызвав sdata. Для других AbstractArray типов sdata просто возвращает сам объект, поэтому sdata безопасно использовать с любым Array-типом объекта.
Конструктор для общего массива имеет вид:
SharedArray{T,N}(dims::NTuple; init=false, pids=Int[])
который создаёт N-мерный общий массив типа T и размера dims для процессов, указанных в pids. В отличие от распределённых массивов, общий массив доступен только для участвующих рабочих узлов, указанных в параметре pids (и создающему процессу, если он находится на том же узле).
Если функция 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, и замерьте время их выполнения во второй раз:
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, является особым и имеетid1. - Только процесс
masterможет добавлять или удалять рабочие процессы. - Все процессы могут напрямую взаимодействовать друг с другом.
Подключение между рабочими узлами (с использованием встроенного транспортного средства TCP/IP) устанавливается следующим образом:
-
addprocsвызывается на главном процессе с объектомClusterManager. -
addprocsвызывает соответствующий методlaunch, который запускает необходимое количество рабочих процессов на соответствующих машинах. - Каждый рабочий процесс начинает прослушивание на свободном порту и записывает информацию о своём хосте и порту в
stdout. - Управляющая программа кластера получает
stdoutкаждого рабочего процесса и делает его доступным для главного процесса. - Главный процесс анализирует эту информацию и устанавливает TCP/IP-соединения с каждым рабочим процессом.
- Каждый рабочий процесс также получает уведомления о других рабочих процессах в кластере.
- Каждый рабочий процесс подключается ко всем рабочим процессам, чей
idменьше, чемidсобственного рабочего процесса. - Таким образом, устанавливается сеть типа «сетка», в которой каждый рабочий процесс напрямую подключен к каждому другому.
Хотя в качестве транспортного уровня по умолчанию используется обычный TCPSocket, для кластера Julia возможно использование собственного транспорта.
Julia предоставляет два встроенных менеджера кластеров:
-
LocalManager, используется при вызовеaddprocs()илиaddprocs(np::Integer) -
SSHManager, используется при вызовеaddprocs(hostnames::Array)со списком имён хостов
LocalManager используется для запуска дополнительных рабочих процессов на том же хосте, тем самым используя многоядерную и многопроцессорную аппаратную платформу.
Таким образом, минимальный менеджер кластера должен:
- быть подтипом абстрактного
ClusterManager - реализовывать
launch, метод, ответственный за запуск новых рабочих процессов - реализовывать
manage, который вызывается при различных событиях в течение жизненного цикла рабочего процесса (например, для отправки сигнала прерывания)
addprocs(manager::FooManager) требует FooManager для реализации:
function launch(manager::FooManager, params::Dict, launched::Array, c::Condition)
[...]
end
function manage(manager::FooManager, id::Integer, config::WorkerConfig, op::Symbol)
[...]
end
В качестве примера рассмотрим реализацию LocalManager, менеджера, ответственного за запуск рабочих процессов на том же хосте:
struct LocalManager <: ClusterManager
np::Integer
end
function launch(manager::LocalManager, params::Dict, launched::Array, c::Condition)
[...]
end
function manage(manager::LocalManager, id::Integer, config::WorkerConfig, op::Symbol)
[...]
end
Метод launch принимает следующие аргументы:
-
manager::ClusterManager: менеджер кластера, с которым вызываетсяaddprocs -
params::Dict: все ключевые аргументы, переданные вaddprocs -
launched::Array: массив, в который нужно добавить один или несколько объектовWorkerConfig -
c::Condition: переменная состояния, которая должна быть уведомлена по мере запуска рабочих процессов
Метод launch вызывается асинхронно в отдельной задаче. Завершение этой задачи сигнализирует о том, что все запрошенные рабочие процессы были запущены. Следовательно, функция launch ДОЛЖНА завершаться, как только все запрошенные рабочие процессы будут запущены.
Недавно запущенные рабочие процессы подключаются друг к другу и к главному процессу по схеме «все со всеми». Указание командной строки --worker[=<cookie>] приводит к тому, что запущенные процессы инициализируются как рабочие процессы, и соединения устанавливаются через TCP/IP-сокеты.
Все рабочие процессы в кластере используют один и тот же куки, что и главный процесс. Когда куки не указан, т. е. с опцией --worker , рабочий процесс пытается прочитать его из стандартного ввода. LocalManager и SSHManager передают куки новым запущенным рабочим процессам через их стандартные вводы.
По умолчанию рабочий процесс будет прослушивать свободный порт по адресу, возвращаемому функцией getipaddr(). Можно указать конкретный адрес для прослушивания с помощью необязательного аргумента --bind-to bind_addr[:port]. Это полезно для многоадресных хостов.
В качестве примера не-TCP/IP-транспорта, реализация может выбрать использование MPI, в этом случае --worker указывать НЕ нужно. Вместо этого, новые запущенные рабочие процессы должны вызвать init_worker(cookie) перед использованием каких-либо параллельных конструкций.
Для каждого запущенного рабочего процесса метод launch должен добавить объект WorkerConfig (с соответствующими инициализированными полями) в launched
mutable struct WorkerConfig
# Common fields relevant to all cluster managers
io::Union{IO, Nothing}
host::Union{AbstractString, Nothing}
port::Union{Integer, Nothing}
# Used when launching additional workers at a host
count::Union{Int, Symbol, Nothing}
exename::Union{AbstractString, Cmd, Nothing}
exeflags::Union{Cmd, Nothing}
# External cluster managers can use this to store information at a per-worker level
# Can be a dict if multiple fields need to be stored.
userdata::Any
# SSHManager / SSH tunnel connections to workers
tunnel::Union{Bool, Nothing}
bind_addr::Union{AbstractString, Nothing}
sshflags::Union{Cmd, Nothing}
max_parallel::Union{Integer, Nothing}
# Used by Local/SSH managers
connect_at::Any
[...]
end
Большинство полей в WorkerConfig используются встроенными менеджерами. Пользовательские менеджеры кластеров обычно указывают только io или host / port:
Если
ioуказан, он используется для чтения информации о хосте/порту. Рабочий процесс Julia выводит свой адрес и порт при запуске. Это позволяет рабочим процессам Julia прослушивать любой свободный порт, а не требовать ручного конфигурирования портов рабочих процессов.Если
ioне указан, используютсяhostиportдля подключения.-
count,exenameиexeflagsимеют отношение к запуску дополнительных рабочих процессов из рабочего процесса. Например, менеджер кластера может запускать по одному рабочему процессу на узел и использовать его для запуска дополнительных рабочих процессов.-
countсо значением целого числаnзапустит общее количествоnрабочих процессов. -
countсо значением:autoзапустит столько рабочих процессов, сколько потоков 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 в репозитории Examples содержит пример использования 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.
Файл cookie кластера
Все процессы в кластере используют один и тот же файл cookie, который по умолчанию является случайно сгенерированной строкой на главном процессе:
-
cluster_cookie()возвращает файл cookie, аcluster_cookie(cookie)()устанавливает его и возвращает новый файл cookie. - Все подключения аутентифицируются с обеих сторон, чтобы гарантировать, что к друг другу могут подключаться только работники, запущенные мастером.
- Файл cookie может быть передан работникам при запуске через аргумент
--worker=<cookie>. Если аргумент--workerуказан без файла cookie, работник пытается прочитать файл cookie из своего стандартного ввода (stdin). После извлечения файла cookiestdinзакрывается немедленно. -
ClusterManagerмогут получить файл cookie на главном сервере, вызвавcluster_cookie(). Управляющие кластерами, не использующие стандартный транспорт TCP/IP (и, следовательно, не указывающие--worker), должны вызватьinit_worker(cookie, manager)с тем же файлом cookie, что и на главном сервере.
Обратите внимание, что среды, требующие более высокого уровня безопасности, могут реализовать это с помощью настраиваемого ClusterManager. Например, файлы cookie могут быть заранее распределены и, следовательно, не указываться как аргумент запуска.
Определение топологии сети (экспериментальная функция)
Ключевой аргумент topology, передаваемый в addprocs, используется для определения способа подключения работников друг к другу:
-
:all_to_all, по умолчанию: все работники подключены друг к другу. -
:master_worker: только процесс драйвера, т.е.pid1, имеет соединения с работниками. -
:custom: методlaunchуправляющего кластером определяет топологию соединений через поляidentиconnect_identsвWorkerConfig. Работник с идентификатором, предоставленным управляющим кластером,identподключится ко всем работникам, указанным вconnect_idents.
Ключевой аргумент lazy=true|false влияет только на опцию topology :all_to_all. Если true, кластер запускается с подключённым мастером ко всем работникам. Специфические соединения «работник-работник» устанавливаются при первом удалённом вызове между двумя работниками. Это помогает уменьшить начальные ресурсы, выделенные для внутрикластерной связи. Соединения устанавливаются в зависимости от текущих требований параллельной программы. Значение по умолчанию для lazy — true.
В настоящее время отправка сообщения между неподключенными работниками приводит к ошибке. Данное поведение, как и функциональность и интерфейс, следует рассматривать как экспериментальное и оно может измениться в будущих версиях.
Заслуживающие внимания внешние пакеты
Помимо параллелизма в Julia, существует множество внешних пакетов, которые стоит упомянуть. Например, MPI.jl — это оболочка Julia для протокола MPI, или DistributedArrays.jl, как показано в Раздел «Общие массивы». Следует упомянуть о экосистеме Julia для программирования на GPU, которая включает:
Операции низкого уровня (на основе ядра C) OpenCL.jl и CUDAdrv.jl, которые представляют собой соответственно интерфейс OpenCL и обертку CUDA.
Интерфейсы низкого уровня (ядро Julia), такие как CUDAnative.jl, который является реализацией CUDA на основе языка Julia.
Абстракции высокого уровня, специфичные для поставщика, такие как CuArrays.jl и CLArrays.jl
Библиотеки высокого уровня, такие как ArrayFire.jl и GPUArrays.jl
В следующем примере мы будем использовать как DistributedArrays.jl, так и CuArrays.jl, чтобы распределить массив по нескольким процессам, сначала преобразовав его через distribute() и CuArray().
Помните, при импорте DistributedArrays.jl импортировать его во все процессы, используя @everywhere
$ ./julia -p 4
julia> addprocs()
julia> @everywhere using DistributedArrays
julia> using CuArrays
julia> B = ones(10_000) ./ 2;
julia> A = ones(10_000) .* π;
julia> C = 2 .* A ./ B;
julia> all(C .≈ 4*π)
true
julia> typeof(C)
Array{Float64,1}
julia> dB = distribute(B);
julia> dA = distribute(A);
julia> dC = 2 .* dA ./ dB;
julia> all(dC .≈ 4*π)
true
julia> typeof(dC)
DistributedArrays.DArray{Float64,1,Array{Float64,1}}
julia> cuB = CuArray(B);
julia> cuA = CuArray(A);
julia> cuC = 2 .* cuA ./ cuB;
julia> all(cuC .≈ 4*π);
true
julia> typeof(cuC)
CuArray{Float64,1}
Обратите внимание, что некоторые функции Julia в настоящее время не поддерживаются CUDAnative.jl [2], особенно некоторые функции, такие как sin, потребуют замены на CUDAnative.sin(cc: @maleadt).
В следующем примере мы будем использовать как DistributedArrays.jl, так и CuArrays.jl, чтобы распределить массив по нескольким процессам и вызвать на нём общую функцию.
function power_method(M, v)
for i in 1:100
v = M*v
v /= norm(v)
end
return v, norm(M*v) / norm(v) # or (M*v) ./ v
end
power_method многократно создаёт новый вектор и нормализует его. Мы не указали никаких типов данных в объявлении функции, давайте посмотрим, будет ли это работать с вышеупомянутыми типами данных:
julia> M = [2. 1; 1 1];
julia> v = rand(2)
2-element Array{Float64,1}:
0.40395
0.445877
julia> power_method(M,v)
([0.850651, 0.525731], 2.618033988749895)
julia> cuM = CuArray(M);
julia> cuv = CuArray(v);
julia> curesult = power_method(cuM, cuv);
julia> typeof(curesult)
CuArray{Float64,1}
julia> dM = distribute(M);
julia> dv = distribute(v);
julia> dC = power_method(dM, dv);
julia> typeof(dC)
Tuple{DistributedArrays.DArray{Float64,1,Array{Float64,1}},Float64}
Для завершения этого краткого обзора внешних пакетов, мы можем рассмотреть MPI.jl, оболочку Julia для протокола MPI. Так как рассмотрение каждой внутренней функции займет слишком много времени, лучше просто оценить подход, используемый для реализации протокола.
Рассмотрим этот игрушечный скрипт, который просто вызывает каждый подпроцесс, инициализирует его ранг, и когда достигается главный процесс, выполняет сумму рангов.
import MPI
MPI.Init()
comm = MPI.COMM_WORLD
MPI.Barrier(comm)
root = 0
r = MPI.Comm_rank(comm)
sr = MPI.Reduce(r, MPI.SUM, root, comm)
if(MPI.Comm_rank(comm) == root)
@printf("sum of ranks: %s\n", sr)
end
MPI.Finalize()
mpirun -np 4 ./julia example.jl
- 1В данном контексте MPI относится к стандарту MPI-1. Начиная с MPI-2, комитет по стандартам MPI ввёл новую группу механизмов связи, общим названием которых является Доступ к удалённой памяти (RMA). Целью добавления rma в стандарт MPI было упрощение односторонних шаблонов связи. Для получения дополнительной информации о последнем стандарте MPI см. https://mpi-forum.org/docs.
- 2Страницы справки 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.4.2/manual/parallel-computing/