Spec-Zone.ru › Julia 1.3

Параллельное программирование

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

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

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

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

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

Корутины

Платформа параллельного программирования Julia использует задачи (также корутины) для переключения между несколькими вычислениями. Для выражения порядка выполнения между легкими потоками необходимы примитивы взаимодействия. Julia предлагает Channel(func::Function, ctype=Any, csize=0, taskref=nothing), которая создает новую задачу из func, связывает ее с новым каналом типа ctype и размером csize и планирует задачу. Channels может служить способом взаимодействия между задачами, поскольку Channel{T}(sz::Int) создает буферизованный канал типа T и размера sz. Всякий раз, когда код выполняет операцию взаимодействия, такую как fetch или wait, текущая задача приостанавливается, а планировщик выбирает другую задачу для выполнения. Задача перезапускается, когда завершается событие, которое она ожидает.

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

Каналы

Раздел о Task в Управлении потоком обсуждал выполнение нескольких функций в кооперативном режиме. Channel могут быть полезны для передачи данных между запущенными задачами, особенно при работе с операциями ввода-вывода.

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

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

  • Несколько писателей в разных задачах могут одновременно записывать в один и тот же канал через вызовы put!.

  • Несколько читателей в разных задачах могут одновременно читать данные через вызовы take!.

  • В качестве примера:

    # Given Channels c1 and c2,
    c1 = Channel(32)
    c2 = Channel(32)
    
    # and a function `foo` which reads items from c1, processes the item read
    # and writes a result to c2,
    function foo()
        while true
            data = take!(c1)
            [...]               # process data
            put!(c2, result)    # write out result
        end
    end
    
    # we can schedule `n` instances of `foo` to be active concurrently.
    for _ in 1:n
        @async foo()
    end
  • Каналы создаются с помощью конструктора Channel{T}(sz). Канал будет содержать только объекты типа T. Если тип не указан, канал может содержать объекты любого типа. sz относится к максимальному количеству элементов, которые могут храниться в канале в любой момент времени. Например, Channel(32) создает канал, который может хранить максимум 32 объекта любого типа. Channel{MyType}(64) может хранить до 64 объектов типа MyType в любой момент времени.

  • Если канал Channel пуст, читатели (при вызове take!) будут блокироваться до тех пор, пока данные не станут доступны.

  • Если канал Channel полон, писатели (при вызове put!) будут блокироваться до тех пор, пока место не освободится.

  • isready проверяет наличие какого-либо объекта в канале, а wait ожидает, пока объект не станет доступным.

  • Канал Channel изначально находится в открытом состоянии. Это означает, что к нему можно свободно читать и писать через вызовы take! и put!. close закрывает канал Channel. В закрытом канале Channel вызов put! завершится с ошибкой. Например:

    julia> c = Channel(2);
    
    julia> put!(c, 1) # `put!` on an open channel succeeds
    1
    
    julia> close(c);
    
    julia> put!(c, 2) # `put!` on a closed channel throws an exception.
    ERROR: InvalidStateException("Channel is closed.",:closed)
    Stacktrace:
    [...]
  • take! и fetch (получает, но не удаляет значение) в закрытом канале успешно возвращают все существующие значения до тех пор, пока он не опустеет. Продолжая предыдущий пример:

    julia> fetch(c) # Any number of `fetch` calls succeed.
    1
    
    julia> fetch(c)
    1
    
    julia> take!(c) # The first `take!` removes the value.
    1
    
    julia> take!(c) # No more data available on a closed channel.
    ERROR: InvalidStateException("Channel is closed.",:closed)
    Stacktrace:
    [...]

Канал Channel может использоваться как итерируемый объект в цикле for, в этом случае цикл работает до тех пор, пока канал Channel содержит данные или открыт. Переменная цикла принимает все значения, добавленные в канал Channel. Цикл for завершается, когда канал Channel закрыт и опустел.

Например, следующее приведет к тому, что цикл for будет ждать дополнительных данных:

julia> c = Channel{Int}(10);

julia> foreach(i->put!(c, i), 1:3) # add a few entries

julia> data = [i for i in c]

в то время как это вернет данные после прочтения всех данных:

julia> c = Channel{Int}(10);

julia> foreach(i->put!(c, i), 1:3); # add a few entries

julia> close(c);                    # `for` loops can exit

julia> data = [i for i in c]
3-element Array{Int64,1}:
 1
 2
 3

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

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

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

julia> function do_work()
           for job_id in jobs
               exec_time = rand()
               sleep(exec_time)                # simulates elapsed time doing actual work
                                               # typically performed externally.
               put!(results, (job_id, exec_time))
           end
       end;

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

julia> n = 12;

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

julia> for i in 1:4 # start 4 tasks to process requests in parallel
           @async do_work()
       end

julia> @elapsed while n > 0 # print out results
           job_id, exec_time = take!(results)
           println("$job_id finished in $(round(exec_time; digits=2)) seconds")
           global n = n - 1
       end
4 finished in 0.22 seconds
3 finished in 0.45 seconds
1 finished in 0.5 seconds
7 finished in 0.14 seconds
2 finished in 0.78 seconds
5 finished in 0.9 seconds
9 finished in 0.36 seconds
6 finished in 0.87 seconds
8 finished in 0.79 seconds
10 finished in 0.64 seconds
12 finished in 0.5 seconds
11 finished in 0.97 seconds
0.029772311

В текущей версии Julia все задачи мультиплексируются на один поток ОС. Таким образом, в то время как задачи, связанные с операциями ввода-вывода, извлекают выгоду из параллельного выполнения, задачи, связанные с вычислениями, фактически выполняются последовательно в одном потоке ОС.

Многопоточность (Экспериментальная)

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

Настройка

По умолчанию Julia запускается с одним потоком выполнения. Это можно проверить, используя команду Threads.nthreads():

julia> Threads.nthreads()
1

Количество потоков, с которыми запускается Julia, контролируется переменной среды с именем JULIA_NUM_THREADS. Теперь запустим Julia с 4 потоками:

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
END_OF_DOCUMENT_MARKER
Примечание

Не все примитивные типы могут быть обернуты в тег Atomic. Поддерживаемые типы — Int8, Int16, Int32, Int64, Int128, UInt8, UInt16, UInt32, UInt64, UInt128, Float16, Float32 и Float64. Кроме того, Int128 и UInt128 не поддерживаются на AAarch32 и ppc64le.

Побочные эффекты и изменяемые аргументы функций

При использовании многопоточности нужно быть осторожными при работе с функциями, которые не являются чистыми, так как мы можем получить неправильный результат. Например, функции, имена которых по соглашению оканчиваются на !, изменяют свои аргументы и поэтому не являются чистыми. Однако существуют функции с побочными эффектами, имена которых не оканчиваются на !. Например, findfirst(regex, str) изменяет свой аргумент regex, или rand() изменяет Base.GLOBAL_RNG:

julia> using Base.Threads

julia> nthreads()
4

julia> function f()
           s = repeat(["123", "213", "231"], outer=1000)
           x = similar(s, Int)
           rx = r"1"
           @threads for i in 1:3000
               x[i] = findfirst(rx, s[i]).start
           end
           count(v -> v == 1, x)
       end
f (generic function with 1 method)

julia> f() # the correct result is 1000
1017

julia> function g()
           a = zeros(1000)
           @threads for i in 1:1000
               a[i] = rand()
           end
           length(unique(a))
       end
g (generic function with 1 method)

julia> Random.seed!(1); g() # the result for a single thread is 1000
781

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

Например, чтобы исправить пример findfirst выше, необходимо иметь отдельную копию переменной rx для каждого потока:

julia> function f_fix()
             s = repeat(["123", "213", "231"], outer=1000)
             x = similar(s, Int)
             rx = [Regex("1") for i in 1:nthreads()]
             @threads for i in 1:3000
                 x[i] = findfirst(rx[threadid()], s[i]).start
             end
             count(v -> v == 1, x)
         end
f_fix (generic function with 1 method)

julia> f_fix()
1000

Теперь мы используем Regex("1") вместо r"1", чтобы убедиться, что Julia создаёт отдельные экземпляры объекта Regex для каждой записи вектора rx.

Случай с rand немного сложнее, поскольку нам необходимо гарантировать, что каждый поток использует неперекрывающиеся псевдослучайные последовательности чисел. Это можно просто обеспечить с помощью функции Future.randjump:

julia> using Random; import Future

julia> function g_fix(r)
           a = zeros(1000)
           @threads for i in 1:1000
               a[i] = rand(r[threadid()])
           end
           length(unique(a))
       end
g_fix (generic function with 1 method)

julia>  r = let m = MersenneTwister(1)
                [m; accumulate(Future.randjump, fill(big(10)^20, nthreads()-1), init=m)]
            end;

julia> g_fix(r)
1000

Мы передаём вектор r в g_fix, так как генерация нескольких псевдослучайных последовательностей — ресурсоёмкая операция, поэтому мы не хотим повторять её каждый раз при вызове функции.

@threadcall (Экспериментальное)

Все задачи ввода-вывода, таймеры, команды REPL и т.д. мультиплексируются на один системный поток через цикл событий. Модифицированная версия libuv (http://docs.libuv.org/en/v1.x/) предоставляет эту функциональность. Точки приостановки позволяют кооперативно планировать несколько задач на одном системном потоке. Задачи ввода-вывода и таймеры приостанавливаются неявно, ожидая выполнения события. Явное вызывание yield позволяет планировать выполнение других задач.

Таким образом, задача, выполняющая ccall, фактически препятствует планировщику Julia от выполнения любых других задач до возврата вызова. Это справедливо для всех вызовов внешних библиотек. Исключениями являются вызовы пользовательского кода C, который вызывает Julia (который может затем приостановиться), или код C, который вызывает jl_yield() (C-эквивалент yield).

Обратите внимание, что, хотя код Julia выполняется в одном потоке (по умолчанию), библиотеки, используемые Julia, могут запускать свои внутренние потоки. Например, библиотека BLAS может запускать столько потоков, сколько ядер на машине.

Макрос @threadcall решает сценарии, где мы не хотим, чтобы ccall блокировал основной цикл событий Julia. Он планирует выполнение функции C в отдельном потоке. Для этого используется пул потоков с размером по умолчанию 4. Размер пула потоков регулируется переменной среды UV_THREADPOOL_SIZE. При ожидании свободного потока и во время выполнения функции, как только поток становится доступным, запрашивающая задача (в основном цикле событий Julia) уступает другим задачам. Обратите внимание, что @threadcall не возвращается до завершения выполнения. С точки зрения пользователя, это, следовательно, блокирующий вызов, как и другие API Julia.

Очень важно, чтобы вызываемая функция не вызывала Julia, так как это приведёт к ошибке segfault.

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

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

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

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

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

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

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

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

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

Каждый процесс имеет связанный идентификатор. Процесс, предоставляющий интерактивный приглашение Julia, всегда имеет id, равное 1. Процессы, используемые по умолчанию для параллельных операций, называются «рабочими». Если процессов только один, процесс 1 считается рабочим. В противном случае рабочими считаются все процессы, кроме процесса 1.

Давайте попробуем это. Начав с julia -p n, мы получаем n рабочих процессов на локальной машине. Обычно имеет смысл, чтобы n было равно числу потоков процессора (логических ядер) на машине. Обратите внимание, что аргумент -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.

Чтобы упростить задачу, символ :any можно передать в [@spawnat], что позволит вам выбрать, где выполнить операцию:

julia> r = @spawnat :any rand(2,2)
Future(2, 1, 4, nothing)

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

julia> fetch(s)
2×2 Array{Float64,2}:
 1.38854  1.9098
 1.20939  1.57158

Обратите внимание, что мы использовали 1 .+ fetch(r) вместо 1 .+ r. Это связано с тем, что нам неизвестно, где будет выполняться код, поэтому в общем случае может потребоваться fetch для перемещения r в процесс, выполняющий сложение. В данном случае, @spawnat достаточно умён, чтобы выполнить вычисление на процессе, который владеет r, поэтому fetch будет бесполезной операцией (работа не выполняется).

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

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

@async похож на @spawnat, но выполняет задачи только на локальном процессе. Мы используем его для создания задачи-«подачи» для каждого процесса. Каждая задача выбирает следующий индекс, который нужно вычислить, затем ожидает завершения своего процесса, а затем повторяет, пока не закончатся индексы. Обратите внимание, что задачи-«подачи» не начинаются до тех пор, пока основная задача не достигнет конца блока @sync, в этот момент она уступает управление и ожидает завершения всех локальных задач перед возвратом из функции. Что касается версий v0.7 и выше, задачи-«подачи» могут обмениваться состоянием через nextidx, так как они все выполняются на одном процессе. Даже если Tasks запланированы кооперативно, в некоторых контекстах, как, например, в асинхронном вводе-выводе, может потребоваться блокировка. Это означает, что переключение контекста происходит только в определенных точках: в этом случае, когда вызывается remotecall_fetch. Это текущее состояние реализации, и оно может измениться в будущих версиях Julia, так как это позволит запускать до N Tasks на M Process, т.е. многопоточность М:Н. Тогда потребуется модель приобретения/освобождения блокировок для nextidx, поскольку одновременное чтение-запись ресурса несколькими процессами небезопасно.

Доступность кода и загрузка пакетов

Ваш код должен быть доступен на любом процессе, который его выполняет. Например, введите следующее в приглашении Julia:

julia> function rand2(dims...)
           return 2*rand(dims...)
       end

julia> rand2(2,2)
2×2 Array{Float64,2}:
 0.153756  0.368514
 1.15119   0.918912

julia> fetch(@spawnat :any rand2(2,2))
ERROR: RemoteException(2, CapturedException(UndefVarError(Symbol("#rand2"))
Stacktrace:
[...]

Процесс 1 знал о функции rand2, но процесс 2 — нет.

Чаще всего вы будете загружать код из файлов или пакетов, и у вас есть значительная гибкость в управлении тем, какие процессы загружают код. Рассмотрим файл DummyModule.jl, содержащий следующий код:

module DummyModule

export MyType, f

mutable struct MyType
    a::Int
end

f(x) = x^2+1

println("loaded")

end

Чтобы ссылаться на MyType на всех процессах, необходимо загрузить DummyModule.jl на каждом процессе. Вызов include("DummyModule.jl") загружает его только на одном процессе. Чтобы загрузить его на каждом процессе, используйте макрос @everywhere (запуская Julia с julia -p 2):

julia> @everywhere include("DummyModule.jl")
loaded
      From worker 3:    loaded
      From worker 2:    loaded

Как обычно, это не добавляет DummyModule в область видимости ни на одном из процессов, что требует using или import. Более того, когда DummyModule попадает в область видимости на одном процессе, она не появляется на других:

julia> using .DummyModule

julia> MyType(7)
MyType(7)

julia> fetch(@spawnat 2 MyType(7))
ERROR: On worker 2:
UndefVarError: MyType not defined
⋮

julia> fetch(@spawnat 2 DummyModule.MyType(7))
MyType(7)

Однако, всё ещё возможно, например, отправить MyType процессу, который загрузил DummyModule, даже если она не находится в области видимости:

julia> put!(RemoteChannel(2), MyType(7))
RemoteChannel{Channel{Any}}(2, 1, 13)

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

julia -p <n> -L file1.jl -L file2.jl driver.jl

Процесс Julia, выполняющий скрипт драйвера в примере выше, имеет id, равный 1, так же как и процесс, предоставляющий интерактивное приглашение.

Наконец, если DummyModule.jl не является отдельным файлом, а пакетом, то using DummyModule загрузит DummyModule.jl на всех процессах, но сделает его доступным только на процессе, где был вызван using.

Запуск и управление рабочими процессами

Базовая установка Julia поддерживает два типа кластеров:

  • Локальный кластер, указанный параметром -p, как показано выше.
  • Кластер, охватывающий машины, использующий параметр --machine-file. Это использует бесклавишный ssh логин для запуска рабочих процессов Julia (из той же директории, что и текущий хост) на указанных машинах.

Функции addprocs, rmprocs, workers и другие доступны в качестве программируемого средства для добавления, удаления и запроса процессов в кластере.

julia> using Distributed

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

Модуль Distributed необходимо явно загрузить на главном процессе перед вызовом addprocs. Он автоматически становится доступным для рабочих процессов.

Обратите внимание, что рабочие процессы не запускают скрипт запуска ~/.julia/config/startup.jl, и не синхронизируют свое глобальное состояние (такое как глобальные переменные, новые определения методов и загруженные модули) с другими работающими процессами.

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

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

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

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

Метод 1:

julia> A = rand(1000,1000);

julia> Bref = @spawnat :any A^2;

[...]

julia> fetch(Bref);

Метод 2:

julia> Bref = @spawnat :any rand(1000,1000)^2;

[...]

julia> fetch(Bref);

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

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

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

Выражения, выполняемые удаленно через @spawnat, или замыкания, определённые для удалённого выполнения с использованием remotecall, могут ссылаться на глобальные переменные. Глобальные привязки внутри модуля Main обрабатываются немного по-другому по сравнению с глобальными привязками в других модулях. Рассмотрим следующий фрагмент кода:

A = rand(10,10)
remotecall_fetch(()->sum(A), 2)

В этом случае sum ДОЛЖЕН быть определён в удалённом процессе. Обратите внимание, что A — это глобальная переменная, определённая в локальном рабочем пространстве. Рабочий процесс 2 не имеет переменной, называемой A в рамках Main. Действие отправки замыкания ()->sum(A) рабочему процессу 2 приводит к определению Main.A на 2. Main.A продолжает существовать на рабочем процессе 2 даже после возвращения вызова remotecall_fetch. Удалённые вызовы со встроенными ссылками на глобальные переменные (только для модуля Main) обрабатывают глобальные переменные следующим образом:

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

  • Глобальные константы объявляются как константы и на удалённых узлах.

  • Глобальные переменные пересылаются целевому рабочему процессу только в контексте удалённого вызова, и только если их значение изменилось. Также кластер не синхронизирует глобальные привязки между узлами. Например:

    A = rand(10,10)
    remotecall_fetch(()->sum(A), 2) # worker 2
    A = rand(10,10)
    remotecall_fetch(()->sum(A), 3) # worker 3
    A = nothing

    Выполнение приведенного выше фрагмента кода приводит к тому, что Main.A на рабочем процессе 2 имеет другое значение, чем Main.A на рабочем процессе 3, в то время как значение Main.A на узле 1 устанавливается в nothing.

Как вы могли заметить, хотя память, связанная с глобальными переменными, может быть собрана при повторной привязке на главном процессе, ничего подобного не происходит на рабочих процессах, так как привязки остаются действительными. clear! может использоваться для ручного переназначения определённых глобальных переменных на удалённых узлах в nothing, когда они больше не требуются. Это освободит всю память, связанную с ними, в рамках обычного цикла сборки мусора.

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

Например:

julia> A = rand(10,10);

julia> remotecall_fetch(()->A, 2);

julia> B = rand(10,10);

julia> let B = B
           remotecall_fetch(()->B, 2)
       end;

julia> @fetchfrom 2 InteractiveUtils.varinfo()
name           size summary
––––––––– ––––––––– ––––––––––––––––––––––
A         800 bytes 10×10 Array{Float64,2}
Base                Module
Core                Module
Main                Module

Как видно, глобальная переменная A определена на рабочем процессе 2, но B захвачена как локальная переменная, и поэтому привязка для B не существует на рабочем процессе 2.

Параллельные карты и циклы

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

using SharedArrays

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

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

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

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

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

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

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

julia> pmap(svdvals, M);

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

Если функция 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)

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

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

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

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

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

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

Кластер Джулии имеет следующие характеристики:

  • Изначальный процесс Джулии, также называемый 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-сокеты.

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

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

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

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

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

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

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

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

    # Used by Local/SSH managers
    connect_at::Any

    [...]
end

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

Папка примеров clustermanager/simple демонстрирует пример простого применения сокетов UNIX для настройки кластера.

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

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

  • Главный процесс не прослушивает порты. Он только подключается к рабочим узлам.

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

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

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

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

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

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

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

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

  • cluster_cookie() возвращает куки, а 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, как показано в Раздел о Shared Arrays. Необходимо упомянуть экосистему программирования 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, см. https://mpi-forum.org/docs.

[2]

Страницы справки Julia GPU

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

Spec-Zone.ru

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