Spec-Zone.ru › Julia 1.2

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

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

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

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

Julia также поддерживает экспериментальную многопоточность, где выполнение ветвится, и анонимная функция выполняется на всех потоках. Известный как подход fork-join, параллельные потоки выполняются независимо и, в конечном итоге, должны быть объединены в главном потоке 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 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. Задачи, идентифицированные по id (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, так как это приведёт к ошибке сегментирования.

@threadcall может быть удалена/изменена в будущих версиях Julia.

Обработка на нескольких ядрах или распределённая обработка

Реализация распределённого параллельного вычисления с памятью на разных узлах предоставляется модулем Distributed в составе стандартной библиотеки Julia.

Большинство современных компьютеров имеют более одного процессора, и несколько компьютеров могут быть объединены в кластер. Использование нескольких процессоров позволяет ускорить многие вычисления. На производительность влияют два основных фактора: скорость самих процессоров и скорость доступа к памяти. В кластере очевидно, что данный процессор имеет самый быстрый доступ к ОЗУ на том же компьютере (узле). Возможно, более удивительно, что аналогичные проблемы актуальны и для типичных многоядерных ноутбуков из-за различий в скорости оперативной памяти и кэша. Следовательно, хорошая среда многопроцессорной обработки должна позволять управлять «владением» куском памяти конкретным процессором. Julia предоставляет многопроцессорную среду, основанную на передаче сообщений, позволяя программам запускаться одновременно на нескольких процессах в отдельных областях памяти.

Реализация передачи сообщений в Julia отличается от других сред, таких как MPI [1]. Общение в Julia обычно «одностороннее», что означает, что программисту нужно явно управлять только одним процессом в двухпроцессной операции. Кроме того, эти операции обычно не выглядят как «отправка сообщения» и «приём сообщения», а скорее напоминают операции более высокого уровня, такие как вызовы пользовательских функций.

Распределённое программирование в Julia основано на двух примитивах: удалённые ссылки и удалённые вызовы. Удалённая ссылка — это объект, который может использоваться из любого процесса для ссылки на объект, хранящийся в определённом процессе. Удалённый вызов — это запрос одного процесса на вызов определённой функции с определёнными аргументами в другом (возможно, том же) процессе.

Удалённые ссылки бывают двух типов: Future и RemoteChannel.

Удаленный вызов возвращает Future для своего результата. Удалённые вызовы возвращают немедленно; процесс, который сделал вызов, переходит к следующей операции, в то время как удалённый вызов происходит где-то ещё. Вы можете дождаться завершения удалённого вызова, вызвав wait на возвращённой Future, а полное значение результата можно получить с помощью fetch.

С другой стороны, RemoteChannel — это переписываемые ссылки. Например, несколько процессов могут координировать свою обработку, ссылаясь на одну и ту же удалённую Channel.

Каждый процесс имеет связанный с ним идентификатор. Процесс, предоставляющий интерактивный приглашающий интерфейс Julia, всегда имеет идентификатор, равный id . Процессы, используемые по умолчанию для параллельных операций, называются «работниками». Если процессов только один, процесс 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. Затем потребуется модель приобретения/освобождения блокировок для 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. Для запуска рабочих процессов Julia (из той же директории, что и текущий хост) на указанных машинах используется бесключевой ssh вход.

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

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

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

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

using SharedArrays

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

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

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

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

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

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

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

julia> pmap(svdvals, M);

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

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

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

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

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

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

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

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

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

Каналы и RemoteChannels

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

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

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

julia> addprocs(4); # add worker processes

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

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

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

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

julia> n = 12;

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

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

Локальные вызовы(@id man-distributed-local-invocations)

Данные обязательно копируются на удаленный узел для выполнения. Это относится как к удаленным вызовам, так и к хранению данных в RemoteChannel / Future на другом узле. Как ожидалось, это приводит к созданию копии сериализованных объектов на удаленном узле. Однако, когда целевой узел является локальным узлом, т.е. идентификатор вызывающего процесса совпадает с идентификатором удалённого узла, выполняется локальный вызов. Обычно (но не всегда) он выполняется в другой задаче — но нет сериализации/десериализации данных. Следовательно, вызов относится к тем же экземплярам объектов, что и переданные — копии не создаются. Это поведение показано ниже:

julia> using Distributed;

julia> rc = RemoteChannel(()->Channel(3));   # RemoteChannel created on local node

julia> v = [0];

julia> for i in 1:3
           v[1] = i                          # Reusing `v`
           put!(rc, v)
       end;

julia> result = [take!(rc) for _ in 1:3];

julia> println(result);
Array{Int64,1}[[3], [3], [3]]

julia> println("Num Unique objects : ", length(unique(map(objectid, result))));
Num Unique objects : 1

julia> addprocs(1);

julia> rc = RemoteChannel(()->Channel(3), workers()[1]);   # RemoteChannel created on remote node

julia> v = [0];

julia> for i in 1:3
           v[1] = i
           put!(rc, v)
       end;

julia> result = [take!(rc) for _ in 1:3];

julia> println(result);
Array{Int64,1}[[1], [2], [3]]

julia> println("Num Unique objects : ", length(unique(map(objectid, result))));
Num Unique objects : 3

Как видно, вызов put! для локального RemoteChannel с тем же объектом v модифицированным между вызовами, приводит к сохранению того же единственного экземпляра объекта. В отличие от копий v при создании, когда узел, владеющий rc, является другим узлом.

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

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

julia> using Distributed; addprocs(1);

julia> v = [0];

julia> v2 = remotecall_fetch(x->(x[1] = 1; x), myid(), v);     # Executed on local node

julia> println("v=$v, v2=$v2, ", v === v2);
v=[1], v2=[1], true

julia> v = [0];

julia> v2 = remotecall_fetch(x->(x[1] = 1; x), workers()[1], v); # Executed on remote node

julia> println("v=$v, v2=$v2, ", v === v2);
v=[0], v2=[1], false

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

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

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

Общие массивы используют общую системную память для отображения одного массива во многих процессах. Несмотря на некоторые сходства с DArray, поведение SharedArray существенно отличается. В DArray каждый процесс имеет локальный доступ только к части данных, и ни два процесса не разделяют одну и ту же часть; в отличие от SharedArray, каждый «участвующий» процесс имеет доступ ко всему массиву. SharedArray — хороший выбор, когда вам нужно иметь большое количество данных, совместно доступных двум или более процессам на одном компьютере.

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

SharedArray индексирование (присвоение и доступ к значениям) работает так же, как и для обычных массивов, и является эффективным, так как подлежащая память доступна для локального процесса. Следовательно, большинство алгоритмов работают естественно с SharedArray, хотя и в режиме одного процесса. В случаях, когда алгоритм настаивает на Array в качестве входных данных, подлежащий массив можно извлечь из SharedArray с помощью вызова sdata. Для других AbstractArray типов, sdata просто возвращает сам объект, поэтому использование sdata с любым Array-типом объекта безопасно.

Конструктор для общего массива имеет вид:

SharedArray{T,N}(dims::NTuple; init=false, pids=Int[])

что создаёт N-мерный общий массив типа T и размера dims для процессов, указанных pids. В отличие от распределённых массивов, общий массив доступен только для указанных участвующих рабочих процессов параметром pids (и создающий процесс тоже, если он на том же хосте).

Если функция init, с сигнатурой initfn(S::SharedArray), указана, она вызывается на всех участвующих рабочих процессах. Вы можете указать, чтобы каждый рабочий процесс выполнял функцию init на отдельной части массива, тем самым распараллелив инициализацию.

Вот краткий пример:

julia> using Distributed

julia> addprocs(3)
3-element Array{Int64,1}:
 2
 3
 4

julia> @everywhere using SharedArrays

julia> S = SharedArray{Int,2}((3,4), init = S -> S[localindices(S)] = 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 запустит столько рабочих процессов, сколько потоков 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 для циклов обработки сообщений. Менеджер-специфические обратные вызовы, которые необходимо реализовать:

init_worker(cookie, manager::FooManager)

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

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

  • Рабочие процессы 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() возвращает куки, а cluster_cookie(cookie)() устанавливает его и возвращает новый куки.
  • Все соединения аутентифицируются с обеих сторон, чтобы гарантировать, что к друг другу могут подключаться только рабочие процессы, запущенные главным процессом.
  • Куки можно передать рабочим процессам при запуске с помощью аргумента --worker=<cookie>. Если аргумент --worker указан без куки, рабочий процесс пытается прочитать куки из стандартного ввода (stdin). stdin закрывается сразу после получения куки.
  • Менеджеры кластера могут получить куки в главном процессе, вызвав 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, как показано в Разделяемые массивы. Необходимо упомянуть экосистему Julia для программирования на GPU, которая включает:

  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 см. 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.2.0/manual/parallel-computing/

Spec-Zone.ru

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