Spec-Zone.ru › Julia 0.7

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

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

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

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

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

В конце мы представим способ 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 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
        @schedule 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> @schedule make_jobs(n); # feed the jobs channel with "n" jobs

julia> for i in 1:4 # start 4 tasks to process requests in parallel
           @schedule 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,2)) seconds")
           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 shell на этих платформах, вам следует использовать ключевое слово 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, m, fill(big(10)^20, nthreads()-1))]
            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 создать 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 планируются кооперативно, в некоторых контекстах, таких как асинхронное ввод-вывод asynchronous I\O, всё равно может потребоваться блокировка. Это означает, что переключения контекста происходят только в определённых точках: в данном случае, когда вызывается remotecall_fetch. Это текущее состояние реализации (dev v0.7), и оно может измениться в будущих версиях Julia, так как это позволит запустить до N Tasks на M Process, также известное как 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 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 похожи на последовательные циклы for, их поведение значительно отличается. В частности, итерации не происходят в заданном порядке, а записи в переменные или массивы не будут глобально видимыми, так как итерации выполняются на разных процессах. Любые переменные, используемые внутри параллельного цикла, будут скопированы и распространены на каждый процесс.

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

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

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

using SharedArrays

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

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

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

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

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

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

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

julia> pmap(svdvals, M);

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

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

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

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

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

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

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

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

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

Каналы и RemoteChannels

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

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

Мы запускаем 4 рабочих узла для обработки одного jobs удалённого канала. Задачи, идентифицированные по 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,2)) seconds on worker $where")
           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 приводит к немедленной отправке сообщения на удалённый узел для удаления его ссылки на значение.

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

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

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

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

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

Управляющие кластерами

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

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

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

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

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

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

В то время как слой транспорта по умолчанию использует обычный TCPSocket, возможно, что кластер Julia предоставляет свой собственный транспорт.

Julia предоставляет два встроенных управляющих кластерами:

  • LocalManager, используется, когда вызываются addprocs() или addprocs(np::Integer)
  • SSHManager, используется, когда addprocs(hostnames::Array) вызывается со списком имён хостов

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

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

  • быть подтипом абстрактного типа ClusterManager
  • реализовывать launch, метод, отвечающий за запуск новых рабочих
  • реализовывать manage, который вызывается при различных событиях в течение жизненного цикла рабочего (например, отправка сигнала прерывания)

addprocs(manager::FooManager) требует, чтобы FooManager реализовал:

function launch(manager::FooManager, params::Dict, launched::Array, c::Condition)
    [...]
end

function manage(manager::FooManager, id::Integer, config::WorkerConfig, op::Symbol)
    [...]
end

В качестве примера давайте посмотрим, как реализован LocalManager, менеджер, отвечающий за запуск рабочих на одном хосте:

struct LocalManager <: ClusterManager
    np::Integer
end

function launch(manager::LocalManager, params::Dict, launched::Array, c::Condition)
    [...]
end

function manage(manager::LocalManager, id::Integer, config::WorkerConfig, op::Symbol)
    [...]
end

Метод launch принимает следующие аргументы:

  • manager::ClusterManager: менеджер кластера, с которым вызывается addprocs
  • params::Dict: все ключевые аргументы, переданные в addprocs
  • launched::Array: массив, к которому добавляется один или несколько объектов WorkerConfig
  • c::Condition: переменная условия, которая уведомляется по мере запуска рабочих процессов

Метод launch вызывается асинхронно в отдельной задаче. Завершение этой задачи сигнализирует о том, что все запрошенные рабочие процессы были запущены. Следовательно, функция launch ДОЛЖНА завершиться сразу же, как только все запрошенные рабочие процессы будут запущены.

Недавно запущенные рабочие процессы подключены друг к другу и к мастер-процессу по принципу «все со всеми». Указание командной строки аргумента --worker[=<cookie>] приводит к тому, что запущенные процессы инициализируются как рабочие процессы, а подключения устанавливаются через сокеты TCP/IP.

Все рабочие процессы в кластере используют тот же куки, что и мастер-процесс. Когда куки не указан, т. е. с опцией --worker, рабочий процесс пытается прочитать его из стандартного ввода. LocalManager и SSHManager передают куки новым запущенным рабочим процессам через их стандартный ввод.

По умолчанию рабочий процесс будет прослушивать свободный порт по адресу, возвращаемому вызовом getipaddr(). Указать конкретный адрес для прослушивания можно с помощью необязательного аргумента --bind-to bind_addr[:port]. Это полезно для многоадресных хостов.

В качестве примера транспортного средства, отличного от TCP/IP, реализация может выбрать использование MPI, в этом случае --worker не должно быть указано. Вместо этого новые запущенные рабочие процессы должны вызвать init_worker(cookie) перед использованием каких-либо параллельных конструкций.

Для каждого запущенного рабочего процесса метод launch должен добавить объект WorkerConfig (с соответствующими инициализированными полями) в launched

mutable struct WorkerConfig
    # Common fields relevant to all cluster managers
    io::Union{IO, Nothing}
    host::Union{AbstractString, Nothing}
    port::Union{Integer, Nothing}

    # Used when launching additional workers at a host
    count::Union{Int, Symbol, Nothing}
    exename::Union{AbstractString, Cmd, Nothing}
    exeflags::Union{Cmd, Nothing}

    # External cluster managers can use this to store information at a per-worker level
    # Can be a dict if multiple fields need to be stored.
    userdata::Any

    # SSHManager / SSH tunnel connections to workers
    tunnel::Union{Bool, Nothing}
    bind_addr::Union{AbstractString, Nothing}
    sshflags::Union{Cmd, Nothing}
    max_parallel::Union{Integer, Nothing}

    # Used by Local/SSH managers
    connect_at::Any

    [...]
end

Большинство полей в WorkerConfig используются встроенными менеджерами. Пользовательские менеджеры кластеров обычно указывают только io или host / port:

  • Если io указано, оно используется для чтения информации о хосте/порту. Рабочий процесс Julia печатает свой адрес привязки и порт при запуске. Это позволяет рабочим процессам Julia прослушивать любой доступный свободный порт вместо необходимости ручной настройки портов рабочих процессов.

  • Если io не указано, host и port используются для подключения.

  • count, exename и exeflags имеют отношение к запуску дополнительных рабочих процессов из рабочего процесса. Например, менеджер кластера может запускать один рабочий процесс на узел и использовать его для запуска дополнительных рабочих процессов.

    • count со значением целого числа n запустит общее количество n рабочих процессов.
    • count со значением :auto запустит столько рабочих процессов, сколько потоков ЦП (логических ядер) на этом компьютере.
    • exename — это имя исполняемого файла julia, включая полный путь.
    • exeflags должно быть установлено на необходимые аргументы командной строки для новых рабочих процессов.
  • tunnel, bind_addr, sshflags и max_parallel используются, когда требуется SSH-туннель для подключения к рабочим процессам из мастер-процесса.

  • userdata предоставляется для пользовательских менеджеров кластеров для хранения собственной информации, специфичной для рабочих процессов.

manage(manager::FooManager, id::Integer, config::WorkerConfig, op::Symbol) вызывается в разное время в течение жизненного цикла рабочего процесса с соответствующими значениями op:

  • со значениями :register/:deregister при добавлении/удалении рабочего процесса из пула рабочих процессов Julia.
  • со значением :interrupt при вызове interrupt(workers). ClusterManager должно сигнализировать соответствующему рабочему процессу с помощью сигнала прерывания.
  • со значением :finalize в целях очистки.

Менеджеры кластеров с пользовательскими транспортными средствами

Замена стандартных подключений сокетов TCP/IP всех процессов на пользовательский транспортный уровень немного сложнее. Каждый процесс Julia имеет столько задач связи, сколько рабочих процессов он подключен к нему. Например, рассмотрим кластер Julia из 32 процессов в сетях с мешаной сетью «все со всеми»:

  • Каждый процесс Julia имеет таким образом 31 задачу связи.
  • Каждая задача обрабатывает все входящие сообщения от одного удаленного рабочего процесса в цикле обработки сообщений.
  • Цикл обработки сообщений ожидает объект IO (например, TCPSocket в реализации по умолчанию), считывает всё сообщение, обрабатывает его и ожидает следующего.
  • Отправка сообщений в процесс выполняется напрямую из любой задачи Julia — не только задач связи — опять же, через соответствующий объект IO.

Замена стандартного транспорта требует от новой реализации настройки подключений к удаленным рабочим процессам и предоставления соответствующих объектов IO для ожидания циклов обработки сообщений. Реализуемые обратные вызовы, специфичные для менеджера:

connect(manager::FooManager, pid::Integer, config::WorkerConfig)
kill(manager::FooManager, pid::Int, config::WorkerConfig)

Реализация по умолчанию (использующая сокеты TCP/IP) реализована как connect(manager::ClusterManager, pid::Integer, config::WorkerConfig).

connect должен вернуть пару объектов IO—один для чтения данных, отправленных рабочим процессом pid, и другой для записи данных, которые необходимо отправить рабочему процессу pid. Пользовательские менеджеры кластеров могут использовать интерактивный BufferStream в качестве артерии для пересылки данных между пользовательским, возможно, не-IO транспортом и встроенной параллельной инфраструктурой Julia.

BufferStream — это интерактивный IOBuffer, который ведет себя как IO — это поток, который может обрабатываться асинхронно.

Папка clustermanager/0mq в репозитории примеров содержит пример использования ZeroMQ для подключения рабочих процессов Julia в топологии звезды с брокером 0MQ посередине. Примечание: процессы Julia все еще логически подключены друг к другу — любой рабочий процесс может послать сообщение любому другому рабочему процессу без каких-либо знаний об использовании 0MQ в качестве транспортного слоя.

При использовании пользовательских транспортных средств:

  • Рабочие процессы Julia НЕ должны запускаться с --worker. Запуск с --worker приведет к тому, что новые запущенные рабочие процессы по умолчанию будут использовать реализацию транспортного средства TCP/IP.
  • Для каждого входящего логического подключения к рабочему процессу необходимо вызвать Base.process_messages(rd::IO, wr::IO)(). Это запускает новую задачу, которая обрабатывает чтение и запись сообщений с/из рабочего процесса, представленного объектами IO.
  • init_worker(cookie, manager::FooManager) должен быть вызван в качестве части инициализации процесса рабочего процесса.
  • Поле connect_at::Any в WorkerConfig может быть установлено менеджером кластера при вызове launch. Значение этого поля передается во все вызовы connect. Как правило, оно несет информацию о том, как подключиться к рабочему процессу. Например, транспортный механизм TCP/IP использует это поле для указания кортежа (host, port) для подключения к рабочему процессу.

kill(manager, pid, config) вызывается для удаления рабочего процесса из кластера. В мастер-процессе соответствующие объекты IO должны быть закрыты реализацией для обеспечения надлежащей очистки. Реализация по умолчанию просто выполняет вызов exit() на указанном удаленном рабочем процессе.

Папка примеров clustermanager/simple — это пример простой реализации с использованием сокетов доменной системы UNIX для настройки кластера.

Требования к сети для LocalManager и SSHManager

Кластеры Julia предназначены для выполнения в уже защищенных средах на инфраструктуре, такой как локальные ноутбуки, ведомственные кластеры или даже облака. В этом разделе рассматриваются требования к безопасности сети для встроенных LocalManager и SSHManager:

  • Мастер-процесс не прослушивает никаких портов. Он подключается только к рабочим процессам.

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

  • LocalManager, используемый addprocs(N), по умолчанию связывается только с интерфейсом обратной связи. Это означает, что рабочие процессы, запущенные позже на удаленных хостах (или кем-либо с недобросовестными намерениями), не могут подключиться к кластеру. Вызов addprocs(4) за которым следует addprocs(["remote_host"]) завершится с ошибкой. Некоторым пользователям может потребоваться создать кластер, включающий их локальную систему и несколько удаленных систем. Это можно сделать, явно запросив LocalManager для привязки к внешнему сетевому интерфейсу с помощью ключевого аргумента restrict: addprocs(4; restrict=false).

  • SSHManager, используемый addprocs(list_of_remote_hosts), запускает рабочие процессы на удаленных хостах через SSH. По умолчанию SSH используется только для запуска рабочих процессов Julia. Последующие подключения мастер-рабочий и рабочий-рабочий используют обычные, незашифрованные сокеты TCP/IP. На удаленных хостах должна быть включена беспарольная авторизация. Дополнительные флаги SSH или учетные данные могут быть указаны с помощью ключевого аргумента sshflags.

  • addprocs(list_of_remote_hosts; tunnel=true, sshflags=<ssh keys and other flags>) полезно, когда мы хотим использовать SSH-подключения и для мастер-рабочего. Типичный сценарий для этого — локальный ноутбук с запуском Julia REPL (то есть мастер) со всем остальным кластером в облаке, скажем, на Amazon EC2. В этом случае необходимо открыть только порт 22 в удаленном кластере, а также выполнить аутентификацию SSH-клиента с использованием инфраструктуры открытых ключей (PKI). Учетные данные для аутентификации можно предоставить через sshflags, например sshflags=`-e <keyfile>`.

    В топологии «все со всеми» (по умолчанию) все рабочие процессы подключаются друг к другу через обычные TCP-сокеты. Таким образом, политика безопасности на узлах кластера должна обеспечивать свободное подключение между рабочими процессами для диапазона эфемерных портов (различается в зависимости от ОС).

    Защита и шифрование всего трафика между рабочими процессами (через SSH) или шифрование отдельных сообщений можно выполнить с помощью пользовательского ClusterManager.

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

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

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

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

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

Аргумент ключевого слова topology, переданный в addprocs используется для определения того, как рабочие должны подключаться друг к другу:

  • :all_to_all, по умолчанию: все рабочие подключены друг к другу.
  • :master_worker: только процесс драйвера, то есть pid 1, имеет подключения к рабочим.
  • :custom: метод launch управляющего кластером определяет топологию подключения через поля ident и connect_idents в WorkerConfig. Рабочий с предоставленным управляющим кластером идентификатором ident подключится ко всем рабочим, указанным в connect_idents.

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

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

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

Помимо параллельности в Julia существует множество внешних пакетов, которые следует упомянуть. Например, MPI.jl — это обёртка Julia для протокола MPI или DistributedArrays.jl, как показано в Общих массивах. Необходимо упомянуть экосистему программирования GPU Julia, которая включает:

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

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

  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, обратитесь к http://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/v0.7.0/manual/parallel-computing/

Spec-Zone.ru

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