Spec-Zone.ru › Julia 1.1

Параллельное вычисление

Для новичков в многопоточности и параллельных вычислениях полезно сначала понять разные уровни параллелизма, предлагаемые Julia. Мы можем разделить их на три основные категории:

  1. Корутины Julia (зелёные потоки)
  2. Многопоточность
  3. Многоядерная или распределённая обработка

Сначала мы рассмотрим задачи Julia (задачи (также корутины)) и другие модули, которые полагаются на библиотеку времени выполнения Julia, которые позволяют нам приостанавливать и возобновлять вычисления с полным контролем меж-Tasks взаимодействия без необходимости ручного взаимодействия с планировщиком операционной системы. Julia также поддерживает взаимодействие между Tasks через операции, такие как wait и fetch. Взаимодействие и синхронизация данных управляются через Channel, которые являются каналами, обеспечивающими меж-Tasks взаимодействие.

Julia также поддерживает экспериментальную многопоточность, где выполнение разветвляется, и анонимная функция выполняется во всех потоках. Известный как подход «fork-join», параллельные потоки выполняют независимые действия, и в конечном итоге должны быть объединены в главном потоке Julia, чтобы позволить продолжению последовательного выполнения. Многопоточность поддерживается с помощью модуля Base.Threads , который все ещё считается экспериментальным, так как Julia пока не полностью потокобезопасна. В частности, кажется, возникают ошибки segfaults во время операций ввода-вывода и переключения задач. Для актуальной справки следите за следите за трекером задач. Многопоточность следует использовать только в том случае, если вы учитываете глобальные переменные, блокировки и атомарные операции, все из которых будут объяснены позже.

В конечном итоге мы представим подход 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 потоками:

export JULIA_NUM_THREADS=4

(Эта команда работает в оболочках Bourne на Linux и OSX. Обратите внимание, что если вы используете оболочку C на этих платформах, вы должны использовать ключевое слово set вместо export. Если вы работаете в Windows, запустите командную строку в расположении julia.exe и используйте set вместо export.)

Давайте проверим, что у нас есть 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, так как это приведёт к ошибке segfault.

@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 должно быть равно количеству потоков CPU (логическим ядрам) на машине. Обратите внимание, что аргумент -p неявно загружает модуль Distributed.

$ ./julia -p 2

julia> r = remotecall(rand, 2, 2, 2)
Future(2, 1, 4, nothing)

julia> s = @spawnat 2 1 .+ fetch(r)
Future(2, 1, 5, nothing)

julia> fetch(s)
2×2 Array{Float64,2}:
 1.18526  1.50912
 1.16296  1.60607

Первый аргумент к remotecall — это вызываемая функция. Большинство параллельных вычислений в Julia не ссылаются на конкретные процессы или количество доступных процессов, но remotecall считается низкоуровневым интерфейсом, обеспечивающим более точный контроль. Второй аргумент к remotecall — это id процесса, который выполнит работу, а оставшиеся аргументы будут переданы вызываемой функции.

Как видите, в первой строке мы попросили процесс 2 построить 2x2 случайную матрицу, а во второй строке мы попросили его прибавить 1 к ней. Результат обоих вычислений доступен в двух будущих значениях, r и s. Макрос @spawnat вычисляет выражение во втором аргументе на процессе, указанном первым аргументом.

Иногда вам может потребоваться значение, вычисленное удалённо, немедленно. Это обычно происходит, когда вы считываете значение из удалённого объекта для получения данных, необходимых для следующей локальной операции. Для этой цели существует функция remotecall_fetch. Она эквивалентна fetch(remotecall(...)) , но более эффективна.

julia> remotecall_fetch(getindex, 2, r, 1, 1)
0.18526337335308085

Помните, что getindex(r,1,1) эквивалентно r[1,1], поэтому этот вызов извлекает первый элемент из будущего r.

Синтаксис remotecall не очень удобен. Макрос @spawn упрощает задачу. Он работает с выражением, а не с функцией и выбирает, где выполнить операцию за вас:

julia> r = @spawn rand(2,2)
Future(2, 1, 4, nothing)

julia> s = @spawn 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 на процесс, выполняющий сложение. В этом случае @spawn достаточно умён, чтобы выполнить вычисление на процессе, владеющем r, поэтому fetch будет бесполезным (работа не выполняется).

(Стоит отметить, что @spawn не является встроенным, а определён в Julia как макрос. Можно определить свои собственные подобные конструкции.)

Важно помнить, что после получения Future кэширует своё значение локально. Дальнейшие вызовы fetch не требуют сетевого соединения. После того, как все ссылки на Future будут обработаны, удаляется значение, хранящееся на удалённом сервере.

@async похож на @spawn, но выполняет задачи только на локальном процессе. Мы используем его для создания задачи-«поставщика» для каждого процесса. Каждая задача выбирает следующий индекс, который нужно вычислить, затем ждёт завершения своего процесса, а затем повторяет это, пока не закончатся индексы. Обратите внимание, что задачи-поставщики не начинают выполняться до тех пор, пока основная задача не дойдёт до конца блока @sync, в этот момент она уступает управление и ждёт завершения всех локальных задач перед возвратом из функции. Что касается версий 0.7 и выше, задачи-поставщики могут обмениваться состоянием через nextidx, так как все они выполняются на одном процессе. Даже если Tasks планируются кооперативно, в некоторых контекстах, например, в асинхронном вводе-выводе, всё равно может потребоваться блокировка. Это означает, что переключение контекста происходит только в определённых точках: в этом случае, когда вызывается remotecall_fetch. Это текущее состояние реализации, и оно может измениться в будущих версиях Julia, поскольку предполагается, что это позволит запускать до N Tasks на M Process, что известно как модель многопоточности M:N M:N Threading. Тогда потребуется модель приобретения/освобождения блокировки для 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(@spawn 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, и они не синхронизируют своё глобальное состояние (такое как глобальные переменные, новые определения методов и загруженные модули) с другими работающими процессами.

Другие типы кластеров могут быть поддерживаются написанием собственного пользовательского ClusterManager, как описано ниже в разделе ClusterManagers.

Перемещение данных

Отправка сообщений и перемещение данных составляют большую часть накладных расходов в распределённой программе. Снижение количества сообщений и объёма передаваемых данных имеет решающее значение для достижения производительности и масштабируемости. Для этого важно понимать перемещение данных, выполняемое различными конструкциями распределённого программирования Julia.

fetch можно рассматривать как явную операцию перемещения данных, так как она напрямую запрашивает перемещение объекта на локальную машину. @spawn (и некоторые связанные с ним конструкции) также перемещает данные, но это не так очевидно, поэтому это можно назвать неявной операцией перемещения данных. Рассмотрим два подхода к построению и возведению в квадрат случайной матрицы:

Метод 1:

julia> A = rand(1000,1000);

julia> Bref = @spawn A^2;

[...]

julia> fetch(Bref);

Метод 2:

julia> Bref = @spawn rand(1000,1000)^2;

[...]

julia> fetch(Bref);

Разница кажется незначительной, но на самом деле она довольно значительна из-за поведения @spawn. В первом методе случайная матрица строится локально, затем отправляется на другой процесс, где она возводится в квадрат. Во втором методе случайная матрица строится и возводится в квадрат на другом процессе. Поэтому во втором методе передаётся гораздо меньше данных, чем в первом.

В этом примере различия между двумя методами легко различимы и выбор между ними очевиден. Однако в реальной программе проектирование перемещения данных может потребовать большего внимания и, вероятно, некоторых измерений. Например, если первому процессу нужна матрица A, то первый метод может быть лучше. Или, если вычисление A дорогостояще, и только текущий процесс его имеет, то перемещение его на другой процесс может быть неизбежным. Или, если текущий процесс имеет очень мало задач между @spawn и fetch(Bref), то, возможно, лучше полностью отказаться от параллелизма. Или представьте, что rand(1000,1000) заменено на более дорогостоящую операцию. Тогда может быть целесообразно добавить ещё одну инструкцию @spawn только для этого этапа.

Глобальные переменные

Выражения, выполняемые удалённо через @spawn, или лямбда-выражения, указанные для удалённого выполнения, с помощью 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.

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

К счастью, многие полезные параллельные вычисления не требуют перемещения данных. Типичный пример — моделирование Монте-Карло, где несколько процессов могут одновременно обрабатывать независимые испытания моделирования. Мы можем использовать @spawn для подбрасывания монет на двух процессах. Сначала напишите следующую функцию в 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 = @spawn count_heads(100000000)
Future(2, 1, 6, nothing)

julia> b = @spawn 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 может быть обобщено. Мы использовали два явных оператора @spawn, что ограничивает параллелизм двумя процессами. Чтобы запустить его на любом количестве процессов, мы можем использовать параллельный цикл for, работающий в распределенной памяти, который можно записать на Julia с помощью @distributed так:

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

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

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

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

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

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

using SharedArrays

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

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

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

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

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

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

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

julia> pmap(svdvals, M);

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

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

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

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

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

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

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

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

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

Каналы и RemoteChannels

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

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

Мы запускаем 4 рабочих процесса для обработки одного jobs удаленного канала. Задачи, идентифицированные по идентификатору (job_id), записываются в канал. Каждая удаленно выполняемая задача в этой модели читает job_id, ждет случайное количество времени и записывает обратно кортеж из job_id, затраченного времени и собственного pid в канал результатов. Наконец, все results выводятся на главном процессе.

julia> addprocs(4); # add worker processes

julia> const jobs = RemoteChannel(()->Channel{Int}(32));

julia> const results = RemoteChannel(()->Channel{Tuple}(32));

julia> @everywhere function do_work(jobs, results) # define work function everywhere
           while true
               job_id = take!(jobs)
               exec_time = rand()
               sleep(exec_time) # simulates elapsed time doing actual work
               put!(results, (job_id, exec_time, myid()))
           end
       end

julia> function make_jobs(n)
           for i in 1:n
               put!(jobs, i)
           end
       end;

julia> n = 12;

julia> @async make_jobs(n); # feed the jobs channel with "n" jobs

julia> for p in workers() # start tasks on the workers to process requests in parallel
           remote_do(do_work, p, jobs, results)
       end

julia> @elapsed while n > 0 # print out results
           job_id, exec_time, where = take!(results)
           println("$job_id finished in $(round(exec_time; digits=2)) seconds on worker $where")
           global n = n - 1
       end
1 finished in 0.18 seconds on worker 4
2 finished in 0.26 seconds on worker 5
6 finished in 0.12 seconds on worker 4
7 finished in 0.18 seconds on worker 4
5 finished in 0.35 seconds on worker 5
4 finished in 0.68 seconds on worker 2
3 finished in 0.73 seconds on worker 3
11 finished in 0.01 seconds on worker 3
12 finished in 0.02 seconds on worker 3
9 finished in 0.26 seconds on worker 5
8 finished in 0.57 seconds on worker 4
10 finished in 0.58 seconds on worker 2
0.055971741

Удаленные ссылки и распределенное сборка мусора

Объекты, на которые ссылаются удаленные ссылки, могут быть освобождены только тогда, когда все удерживающие ссылки в кластере удалены.

Узел, где хранится значение, отслеживает, у каких рабочих процессов есть ссылка на него. Каждый раз, когда RemoteChannel или (неизвлеченное) Future сериализуется в рабочий процесс, узел, на который указывает ссылка, уведомляется. И каждый раз, когда RemoteChannel или (неизвлеченное) Future собираются мусором локально, узел, владеющий значением, снова уведомляется. Это реализовано в внутреннем кластерно-осознающем сериализаторе. Удаленные ссылки действительны только в контексте работающего кластера. Сериализация и десериализация ссылок на обычные IO объекты не поддерживаются.

Уведомления выполняются путем отправки сообщений отслеживания — сообщения «добавить ссылку», когда ссылка сериализуется в другой процесс, и сообщение «удалить ссылку», когда ссылка локально собирается мусором.

Поскольку Future одноразовые и кешируются локально, действие fetch для Future также обновляет информацию о отслеживании ссылок на узле, владеющем значением.

Узел, который владеет значением, освобождает его, как только все ссылки на него будут удалены.

С Future, сериализация уже извлеченного Future в другой узел также отправляет значение, так как исходное удаленное хранилище может собрать значение к этому моменту.

END_OF_DOCUMENT_MARKER

Важно отметить, что когда объект подвергается локальному сборке мусора, зависит от размера объекта и текущего давления на память в системе.

В случае удаленных ссылок размер объекта локальной ссылки довольно мал, в то время как значение, хранящееся на удаленном узле, может быть довольно большим. Поскольку локальный объект может не быть собран немедленно, рекомендуется явно вызывать finalize для локальных экземпляров RemoteChannel или невыгруженных Future. Поскольку вызов fetch для Future также удаляет его ссылку из удаленного хранилища, это не требуется для выгруженных Future. Явное вызов finalize приводит к немедленному сообщению на удаленный узел о снятии ссылки на значение.

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

Общие массивы

Общие массивы используют общую системную память для отображения одного и того же массива во многих процессах. Несмотря на некоторое сходство с 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 рабочих процессов (и для создающего процесса, если он находится на том же хосте).

Если указана функция 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)] = myid())
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)] = myid())
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! заключается в минимизации трафика между рабочими процессами, что позволяет каждому из них вычислять на назначенной части более длительное время.

Общие массивы и распределенная сборка мусора

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

ClusterManagers

Запуск, управление и взаимодействие процессов Julia в логическом кластере осуществляется через менеджеры кластеров. ClusterManager отвечает за

  • запуск рабочих процессов в среде кластера
  • управление событиями в течение жизни каждого рабочего процесса
  • по желанию, предоставление транспорта данных

У кластера Julia есть следующие характеристики:

  • Изначальный процесс Julia, также называемый master, является специальным и имеет id равным 1.
  • Только процесс master может добавлять или удалять рабочие процессы.
  • Все процессы могут напрямую взаимодействовать друг с другом.

Соединение между рабочими процессами (с использованием встроенного транспорта TCP/IP) устанавливается следующим образом:

  • addprocs вызывается в главном процессе с объектом ClusterManager.
  • addprocs вызывает соответствующий метод launch, который запускает требуемое количество рабочих процессов на соответствующих машинах.
  • Каждый рабочий процесс начинает прослушивать свободный порт и записывает информацию о хосте и порте в stdout.
  • Менеджер кластера захватывает stdout каждого рабочего процесса и делает его доступным для главного процесса.
  • Главный процесс анализирует эту информацию и устанавливает соединения TCP/IP с каждым рабочим процессом.
  • Каждый рабочий процесс также получает уведомления об остальных рабочих процессах в кластере.
  • Каждый рабочий процесс подключается ко всем рабочим процессам, чья id меньше, чем собственная id рабочего процесса.
  • Таким образом, устанавливается сеть «mesh», в которой каждый рабочий процесс напрямую подключен к каждому другому.

Хотя по умолчанию используется транспортный уровень на основе обычного 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 в репозитории примеров содержит пример использования 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.

Куки кластера

Все процессы в кластере используют один и тот же куки, который по умолчанию является случайной строкой в мастер-процессе:

  • cluster_cookie() возвращает cookie, в то время как cluster_cookie(cookie)() устанавливает его и возвращает новый cookie.
  • Все соединения аутентифицированы с обеих сторон, чтобы гарантировать, что подключаться друг к другу могут только рабочие процессы, запущенные главным процессом.
  • Cookie может быть передан рабочим процессам при запуске через аргумент --worker=<cookie>. Если аргумент --worker указан без cookie, рабочий процесс пытается прочитать cookie из стандартного ввода (stdin). stdin закрывается сразу после получения cookie.
  • ClusterManager могут получить cookie на главном процессе, вызвав cluster_cookie(). Управляющие кластерами, не использующие стандартный транспорт TCP/IP (и, следовательно, не указывающие --worker), должны вызвать init_worker(cookie, manager) с тем же cookie, что и на главном процессе.

Обратите внимание, что среды, требующие более высокого уровня безопасности, могут реализовать это с помощью настраиваемого ClusterManager. Например, cookie могут быть предварительно разделены и, следовательно, не указываться как аргумент запуска.

Определение топологии сети (Экспериментальная версия)

Аргумент ключевого слова 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, или DistributedArrays.jl, как представлено в Разделяемые массивы. Необходимо упомянуть экосистему программирования Julia для графических процессоров, которая включает:

  1. Операции на низком уровне (на основе ядра C) OpenCL.jl и CUDAdrv.jl, которые соответственно являются интерфейсом OpenCL и оболочкой CUDA.

  2. Интерфейсы на низком уровне (ядро Julia), такие как CUDAnative.jl, который представляет собой собственное для Julia реализацию CUDA.

  3. Абстракции на высоком уровне, специфичные для поставщика, такие как CuArrays.jl и CLArrays.jl

  4. Библиотеки на высоком уровне, такие как 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–2019 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.1.1/manual/parallel-computing/

Spec-Zone.ru

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