Spec-Zone.ru › Julia 1.2

Распределённые вычисления

Distributed.addprocsФункция

addprocs(manager::ClusterManager; kwargs...) -> List of process identifiers

Запускает рабочие процессы через указанный менеджер кластера.

Например, кластеры Beowulf поддерживаются через пользовательский менеджер кластера, реализованный в пакете ClusterManagers.jl.

Количество секунд, в течение которого только что запущенный рабочий процесс ожидает установления соединения с мастером, может быть указано через переменную JULIA_WORKER_TIMEOUT в среде рабочего процесса. Актуально только при использовании TCP/IP в качестве транспорта.

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

Примеры

# On busy clusters, call `addprocs` asynchronously
t = @async addprocs(...)
# Utilize workers as and when they come online
if nprocs() > 1   # Ensure at least one new worker is available
   ....   # perform distributed execution
end
# Retrieve newly launched worker IDs, or any error messages
if istaskdone(t)   # Check if `addprocs` has completed to ensure `fetch` doesn't block
    if nworkers() == N
        new_pids = fetch(t)
    else
        fetch(t)
    end
  end
исходный код
addprocs(machines; tunnel=false, sshflags=``, max_parallel=10, kwargs...) -> List of process identifiers

Добавление процессов на удалённые машины через SSH. Требуется julia для установки в одном и том же месте на каждом узле или для доступа через общую файловую систему.

machines — вектор спецификаций машин. Для каждой спецификации запускается рабочий процесс.

Спецификация машины — это либо строка machine_spec , либо кортеж — (machine_spec, count).

machine_spec — строка вида [user@]host[:port] [bind_addr[:port]]. user по умолчанию — текущий пользователь, port — стандартный порт SSH. Если указан [bind_addr[:port]], другие рабочие процессы будут подключаться к этому рабочему процессу на указанном bind_addr и port.

count — количество рабочих процессов, которые необходимо запустить на указанном узле. Если указано как :auto, будет запущено столько рабочих процессов, сколько потоков ЦП на данном узле.

Ключевые аргументы:

  • tunnel: если true, для подключения к рабочему процессу из процесса-мастера будет использоваться SSH-туннелирование. По умолчанию false.

  • sshflags: указывает дополнительные опции SSH, например, sshflags=`-i /home/foo/bar.pem`

  • max_parallel: указывает максимальное количество одновременно подключённых рабочих процессов к одному узлу. По умолчанию 10.

  • dir: указывает рабочую директорию на рабочих процессах. По умолчанию — текущая директория узла (полученная с помощью pwd()).

  • enable_threaded_blas: если true, BLAS будет работать в нескольких потоках в добавленных процессах. По умолчанию false.

  • exename: имя исполняемого файла julia. По умолчанию "$(Sys.BINDIR)/julia" или "$(Sys.BINDIR)/julia-debug", в зависимости от ситуации.

  • exeflags: дополнительные флаги, передаваемые рабочим процессам.

  • topology: Указывает, как рабочие процессы подключаются друг к другу. Отправка сообщения между незапущенными рабочими процессами приводит к ошибке.

    • topology=:all_to_all: Все процессы подключены друг к другу. По умолчанию.

    • topology=:master_worker: Подключается только процесс-драйвер, т.е. pid 1, к рабочим процессам. Рабочие процессы не подключаются друг к другу.

    • topology=:custom: Метод launch менеджера кластера определяет топологию соединения через поля ident и connect_idents в WorkerConfig. Рабочий процесс с идентификатором менеджера кластера ident подключится ко всем рабочим процессам, указанным в connect_idents.

  • lazy: Применимо только с topology=:all_to_all. Если true, соединения между рабочими процессами устанавливаются лениво, т.е. при первом вызове удалённого вызова между рабочими процессами. По умолчанию true.

Переменные окружения :

Если процесс-мастер не сможет установить соединение с только что запущенным рабочим процессом в течение 60.0 секунд, рабочий процесс расценит это как критическую ситуацию и завершит работу. Этот таймаут можно контролировать через переменную окружения JULIA_WORKER_TIMEOUT. Значение JULIA_WORKER_TIMEOUT в процессе-мастере задаёт количество секунд, в течение которого только что запущенный рабочий процесс ожидает установления соединения.

исходный код
addprocs(; kwargs...) -> List of process identifiers

Эквивалентно addprocs(Sys.CPU_THREADS; kwargs...)

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

исходный код
addprocs(np::Integer; restrict=true, kwargs...) -> List of process identifiers

Запускает рабочие процессы с помощью встроенного LocalManager, который запускает рабочие процессы только на локальном узле. Это можно использовать для эффективного использования нескольких ядер. addprocs(4) добавит 4 процесса на локальной машине. Если restrict — true, привязка ограничена 127.0.0.1. Ключевые аргументы dir, exename, exeflags, topology, lazy и enable_threaded_blas имеют тот же эффект, что и документировано для addprocs(machines).

исходный код

Distributed.nprocsФункция

nprocs()

Получение количества доступных процессов.

Примеры

julia> nprocs()
3

julia> workers()
5-element Array{Int64,1}:
 2
 3
исходный код

Distributed.nworkersФункция

nworkers()

Получение количества доступных рабочих процессов. Это значение на единицу меньше nprocs(). Равно nprocs() если nprocs() == 1.

Примеры

$ julia -p 5

julia> nprocs()
6

julia> nworkers()
5
исходный код

Distributed.procsМетод

procs()

Возвращает список всех идентификаторов процессов, включая pid 1 (который не включён в workers()).

Примеры

$ julia -p 5

julia> procs()
3-element Array{Int64,1}:
 1
 2
 3
исходный код

Distributed.procsМетод

procs(pid::Integer)

Возвращает список всех идентификаторов процессов на том же физическом узле. Возвращаются все рабочие процессы, привязанные к тому же ip-адресу, что и pid.

исходный код

Distributed.workersФункция

workers()

Возвращает список идентификаторов всех рабочих процессов.

Примеры

$ julia -p 5

julia> workers()
2-element Array{Int64,1}:
 2
 3
исходный код

Distributed.rmprocsФункция

rmprocs(pids...; waitfor=typemax(Int))

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

Аргумент waitfor задаёт время ожидания завершения работы рабочих процессов:

  • Если не указано, rmprocs будет ждать, пока все запрошенные pids не будут удалены.
  • Исключение ErrorException генерируется, если все рабочие процессы не могут быть завершены до истечения запрошенных waitfor секунд.
  • При значении waitfor 0, вызов возвращается немедленно, а рабочие процессы планируются для удаления в другой задаче. Возвращается объект запланированной Task задачи. Пользователь должен вызвать wait на задаче перед вызовом других параллельных операций.

Примеры

$ julia -p 5

julia> t = rmprocs(2, 3, waitfor=0)
Task (runnable) @0x0000000107c718d0

julia> wait(t)

julia> workers()
3-element Array{Int64,1}:
 4
 5
 6
исходный код

Distributed.interruptФункция

interrupt(pids::Integer...)

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

исходный код
interrupt(pids::AbstractVector=workers())

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

исходный код

Distributed.myidФункция

myid()

Получение идентификатора текущего процесса.

Примеры

julia> myid()
1

julia> remotecall_fetch(() -> myid(), 4)
4
исходный код

Distributed.pmapФункция

pmap(f, [::AbstractWorkerPool], c...; distributed=true, batch_size=1, on_error=nothing, retry_delays=[], retry_check=nothing) -> collection

Преобразование коллекции c путём применения f к каждому элементу с использованием доступных рабочих процессов и задач.

Для нескольких аргументов коллекции применяется f поэлементно.

Обратите внимание, что f должна быть доступна всем рабочим процессам; см. Доступность кода и загрузка пакетов для получения подробной информации.

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

По умолчанию, pmap распределяет вычисления по всем указанным рабочим процессам. Чтобы использовать только локальный процесс и распределять по задачам, укажите distributed=false. Это эквивалентно использованию asyncmap. Например, pmap(f, c; distributed=false) эквивалентно asyncmap(f,c; ntasks=()->nworkers())

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

Любая ошибка останавливает pmap от обработки остальной части коллекции. Чтобы переопределить это поведение, можно указать функцию обработки ошибок через аргумент on_error, который принимает один аргумент — исключение. Функция может остановить обработку, повторно выбросив исключение, или, чтобы продолжить, вернуть любое значение, которое затем возвращается в вызывающую функцию вместе с результатами.

Рассмотрим следующие два примера. Первый возвращает объект исключения в строке, второй — 0 вместо любого исключения:

julia> pmap(x->iseven(x) ? error("foo") : x, 1:4; on_error=identity)
4-element Array{Any,1}:
 1
  ErrorException("foo")
 3
  ErrorException("foo")

julia> pmap(x->iseven(x) ? error("foo") : x, 1:4; on_error=ex->0)
4-element Array{Int64,1}:
 1
 0
 3
 0

Ошибки также могут обрабатываться путем повторной попытки не удавшихся вычислений. Ключевые аргументы retry_delays и retry_check передаются в retry как ключевые аргументы delays и check соответственно. Если указано разделение на пакеты, и весь пакет завершается ошибкой, все элементы в пакете повторно обрабатываются.

Обратите внимание, что если оба on_error и retry_delays указаны, обработчик on_error вызывается перед повторной попыткой. Если on_error не вызывает (или не перебрасывает) исключение, элемент не будет повторно обрабатываться.

Пример: при ошибках повторить f для элемента максимум 3 раза без задержки между повторами.

pmap(f, c; retry_delays = zeros(3))

Пример: повторить f только если исключение не является экземпляром InexactError, с экспоненциально возрастающими задержками до 3 раз. Вернуть NaN вместо всех InexactError случаев.

pmap(f, c; on_error = e->(isa(e, InexactError) ? NaN : rethrow()), retry_delays = ExponentialBackOff(n = 3))
source

Distributed.RemoteExceptionТип

RemoteException(captured)

Исключения при удалённых вычислениях перехватываются и повторно выбрасываются локально. RemoteException оборачивает pid рабочего процесса и перехваченное исключение. CapturedException перехватывает удалённое исключение и сериализуемую форму стека вызовов во время возникновения исключения.

source

Distributed.FutureТип

Future(pid::Integer=myid())

Создать Future на процессе pid. По умолчанию pid — текущий процесс.

source

Distributed.RemoteChannelТип

RemoteChannel(pid::Integer=myid())

Создать ссылку на Channel{Any}(1) на процессе pid. По умолчанию pid — текущий процесс.

RemoteChannel(f::Function, pid::Integer=myid())

Создавать ссылки на удалённые каналы определённого размера и типа. f — функция, которая, когда выполняется на pid, должна возвращать реализацию AbstractChannel.

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

По умолчанию pid — текущий процесс.

source

Base.fetchМетод

fetch(x::Future)

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

source

Base.fetchМетод

fetch(c::RemoteChannel)

Подождать и получить значение из RemoteChannel. Исключения обрабатываются так же, как и для Future. Элемент, полученный из канала, не удаляется.

source

Distributed.remotecallМетод

remotecall(f, id::Integer, args...; kwargs...) -> Future

Асинхронно вызвать функцию f с заданными аргументами на указанном процессе. Вернуть Future. Ключевые аргументы, если таковые имеются, передаются в f.

source

Distributed.remotecall_waitМетод

remotecall_wait(f, id::Integer, args...; kwargs...)

Выполнить более быстрый wait(remotecall(...)) в одном сообщении на Worker с идентификатором рабочего процесса id. Ключевые аргументы, если таковые имеются, передаются в f.

См. также wait и remotecall.

source

Distributed.remotecall_fetchМетод

remotecall_fetch(f, id::Integer, args...; kwargs...)

Выполнить fetch(remotecall(...)) в одном сообщении. Ключевые аргументы, если таковые имеются, передаются в f. Любые удалённые исключения перехватываются в RemoteException и выбрасываются.

См. также fetch и remotecall.

Примеры

$ julia -p 2

julia> remotecall_fetch(sqrt, 2, 4)
2.0

julia> remotecall_fetch(sqrt, 2, -4)
ERROR: On worker 2:
DomainError with -4.0:
sqrt will only return a complex result if called with a complex argument. Try sqrt(Complex(x)).
...
source

Distributed.remote_doМетод

remote_do(f, id::Integer, args...; kwargs...) -> nothing

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

Успешное выполнение указывает, что запрос был принят для выполнения на удалённом узле.

Хотя последовательные remotecall на тот же рабочий процесс сериализуются в порядке их вызова, порядок выполнения на удалённом рабочем процессе неопределён. Например, remote_do(f1, 2); remotecall(f2, 2); remote_do(f3, 2) сериализует вызов f1, за которым следуют f2 и f3 в этом порядке. Однако нет гарантии, что f1 выполняется до f3 на рабочем процессе 2.

Любые исключения, выброшенные f, выводятся в stderr на удалённом рабочем процессе.

Ключевые аргументы, если таковые имеются, передаются в f.

source

Base.put!Метод

put!(rr::RemoteChannel, args...)

Сохранить набор значений в RemoteChannel. Если канал заполнен, блокируется до тех пор, пока не освободится место. Возвращает первый аргумент.

source

Base.put!Метод

put!(rr::Future, v)

Сохранить значение в Future rr. Future — это ссылки на удалённые значения, которые можно записать только один раз. Вызов put! для уже заданного Future вызывает Exception. Все асинхронные удалённые вызовы возвращают Future и устанавливают значение в возвращаемое значение вызова по завершении.

source

Base.take!Метод

take!(rr::RemoteChannel, args...)

Извлечь значение(я) из RemoteChannel rr, удалив значение(я) в процессе.

source

Base.isreadyМетод

isready(rr::RemoteChannel, args...)

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

исходный код

Base.isreadyМетод

isready(rr::Future)

Определите, имеет ли Future сохранённое значение.

Если аргумент Future принадлежит другому узлу, этот вызов заблокируется, ожидая ответа. Рекомендуется ожидать rr в отдельной задаче или использовать локальный Channel в качестве прокси:

c = Channel(1)
@async put!(c, remotecall_fetch(long_computation, p))
isready(c)  # will not block
исходный код

Distributed.AbstractWorkerPoolТип

AbstractWorkerPool

Базовый тип для пулов рабочих процессов, таких как WorkerPool и CachingPool. Пул должен реализовывать:

  • push! - добавить нового рабочего процесса в общий пул (доступный + занятый)
  • put! - вернуть рабочего процесса в доступный пул
  • take! - взять рабочего процесса из доступного пула (для выполнения удалённой функции)
  • length - количество доступных рабочих процессов в общем пуле
  • isready - вернуть false, если добавление take! в пул заблокирует его, иначе true

Реализации выше (в AbstractWorkerPool) требуют полей channel::Channel{Int} workers::Set{Int}, где channel содержит PID свободных рабочих процессов, а workers — набор всех рабочих процессов, связанных с этим пулом.

исходный код

Distributed.WorkerPoolТип

WorkerPool(workers::Vector{Int})

Создать пул WorkerPool из вектора идентификаторов рабочих процессов.

Примеры

$ julia -p 3

julia> WorkerPool([2, 3])
WorkerPool(Channel{Int64}(sz_max:9223372036854775807,sz_curr:2), Set([2, 3]), RemoteChannel{Channel{Any}}(1, 1, 6))
исходный код

Distributed.CachingPoolТип

CachingPool(workers::Vector{Int})

Реализация пула AbstractWorkerPool. remote, remotecall_fetch, pmap (и другие удалённые вызовы, которые выполняют функции удалённо) получают выгоду от кэширования сериализованных/десериализованных функций на узлах рабочих процессов, особенно замыканий (которые могут захватывать большое количество данных).

Удаленный кэш сохраняется на протяжении всего жизненного цикла возвращаемого объекта CachingPool. Чтобы очистить кэш раньше, используйте clear!(pool).

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

Примеры

const foo = rand(10^8);
wp = CachingPool(workers())
let foo = foo
    pmap(wp, i -> sum(foo) + i, 1:100);
end

Выше было бы передано foo только один раз каждому рабочему процессу.

исходный код

Distributed.default_worker_poolФункция

default_worker_pool()

WorkerPool, содержащий неиспользуемые workers — используется remote(f) и pmap (по умолчанию).

Примеры

$ julia -p 3

julia> default_worker_pool()
WorkerPool(Channel{Int64}(sz_max:9223372036854775807,sz_curr:3), Set([4, 2, 3]), RemoteChannel{Channel{Any}}(1, 1, 4))
исходный код

Distributed.clear!Метод

clear!(pool::CachingPool) -> pool

Удаляет все кэшированные функции со всех участвующих рабочих процессов.

исходный код

Distributed.remoteФункция

remote([p::AbstractWorkerPool], f) -> Function

Возвращает анонимную функцию, которая выполняет функцию f на доступном рабочем процессе (взятом из WorkerPool p если он предоставлен) с использованием remotecall_fetch.

исходный код

Distributed.remotecallМетод

remotecall(f, pool::AbstractWorkerPool, args...; kwargs...) -> Future

Вариант remotecall(f, pid, ....), использующий WorkerPool. Ожидает и берёт свободный рабочий процесс из pool и выполняет remotecall на нём.

Примеры

$ julia -p 3

julia> wp = WorkerPool([2, 3]);

julia> A = rand(3000);

julia> f = remotecall(maximum, wp, A)
Future(2, 1, 6, nothing)

В этом примере задача выполнялась на процессе с PID 2, вызванная с PID 1.

исходный код

Distributed.remotecall_waitМетод

remotecall_wait(f, pool::AbstractWorkerPool, args...; kwargs...) -> Future

Вариант remotecall_wait(f, pid, ....), использующий WorkerPool. Ожидает и берёт свободный рабочий процесс из pool и выполняет remotecall_wait на нём.

Примеры

$ julia -p 3

julia> wp = WorkerPool([2, 3]);

julia> A = rand(3000);

julia> f = remotecall_wait(maximum, wp, A)
Future(3, 1, 9, nothing)

julia> fetch(f)
0.9995177101692958
исходный код

Distributed.remotecall_fetchМетод

remotecall_fetch(f, pool::AbstractWorkerPool, args...; kwargs...) -> result

Вариант remotecall_fetch(f, pid, ....), использующий WorkerPool. Ожидает и берёт свободный рабочий процесс из pool и выполняет remotecall_fetch на нём.

Примеры

$ julia -p 3

julia> wp = WorkerPool([2, 3]);

julia> A = rand(3000);

julia> remotecall_fetch(maximum, wp, A)
0.9995177101692958
исходный код

Distributed.remote_doМетод

remote_do(f, pool::AbstractWorkerPool, args...; kwargs...) -> nothing

Вариант remote_do(f, pid, ....), использующий WorkerPool. Ожидает и берёт свободный рабочий процесс из pool и выполняет remote_do на нём.

исходный код

Distributed.@spawnМакрос

@spawn

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

Примеры

julia> addprocs(3);

julia> f = @spawn myid()
Future(2, 1, 5, nothing)

julia> fetch(f)
2

julia> f = @spawn myid()
Future(3, 1, 7, nothing)

julia> fetch(f)
3
исходный код

Distributed.@spawnatМакрос

@spawnat

Создаёт замыкание вокруг выражения и выполняет замыкание асинхронно на процессе p. Возвращает Future результата. Принимает два аргумента, p и выражение.

Примеры

julia> addprocs(1);

julia> f = @spawnat 2 myid()
Future(2, 1, 3, nothing)

julia> fetch(f)
2
исходный код

Distributed.@fetchМакрос

@fetch

Эквивалентно fetch(@spawn expr). См. fetch и @spawn.

Примеры

julia> addprocs(3);

julia> @fetch myid()
2

julia> @fetch myid()
3

julia> @fetch myid()
4

julia> @fetch myid()
2
исходный код

Distributed.@fetchfromМакрос

@fetchfrom

Эквивалентно fetch(@spawnat p expr). См. fetch и @spawnat.

Примеры

julia> addprocs(3);

julia> @fetchfrom 2 myid()
2

julia> @fetchfrom 4 myid()
4
исходный код

Distributed.@distributedМакрос

@distributed

Параллельный цикл for с распределённой памятью вида:

@distributed [reducer] for var = range
    body
end

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

Обратите внимание, что без функции reduce, @distributed выполняется асинхронно, т.е. запускает независимые задачи на всех доступных рабочих процессах и возвращается немедленно, не дожидаясь завершения. Чтобы дождаться завершения, добавьте префикс @sync, как в примере:

@sync @distributed for var = range
    body
end
исходный код

Distributed.@everywhereМакрос

@everywhere [procs()] expr

Выполняет выражение под Main на всех procs. Ошибки на любом из процессов собираются в CompositeException и генерируются. Например:

@everywhere bar = 1

определит Main.bar на всех процессах.

В отличие от @spawn и @spawnat, @everywhere не захватывает локальные переменные. Вместо этого, локальные переменные могут быть переданы с помощью интерполяции:

foo = 1
@everywhere bar = $foo

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

Эквивалентно вызову remotecall_eval(Main, procs, expr).

исходный код

Distributed.clear!Метод

clear!(syms, pids=workers(); mod=Main)

Очищает глобальные связи в модулях, инициализируя их значением nothing. syms должен иметь тип Symbol или коллекцию Symbol. pids и mod идентифицируют процессы и модуль, в котором нужно переинициализировать глобальные переменные. Очищаются только те имена, которые определены в mod.

Если требуется очистить глобальную константу, генерируется исключение.

исходный код

Distributed.remoteref_idФункция

remoteref_id(r::AbstractRemoteRef) -> RRID

Удаленные ссылки и каналы идентифицируются полями:

  • where - указывает на узел, где фактически находится объект/хранилище, на которое ссылается ссылка.

  • whence - указывает на узел, с которого была создана удаленная ссылка. Обратите внимание, что это отличается от узла, где фактически находится объект, на который ссылаются. Например, вызов RemoteChannel(2) с главного процесса приведет к значению where равным 2 и значению whence равным 1.

  • id уникален для всех ссылок, созданных с рабочего процесса, указанного в whence.

Вместе whence и id уникально идентифицируют ссылку по всем рабочим процессам.

remoteref_id - это низкоуровневый API, который возвращает объект RRID, который оборачивает значения whence и id удаленной ссылки.

исходный код

Distributed.channel_from_idФункция

channel_from_id(id) -> c

Низкоуровневый API, который возвращает базовое AbstractChannel для id, возвращённого remoteref_id. Вызов допустим только на узле, где существует базовое значение канала.

исходный код

Distributed.worker_id_from_socketФункция

worker_id_from_socket(s) -> pid

Низкоуровневый API, который, получив соединение IO или Worker, возвращает pid рабочего процесса, к которому оно подключено. Это полезно при написании пользовательских методов serialize для типа, что оптимизирует данные, записываемые в зависимости от идентификатора принимающего процесса.

исходный код

Distributed.cluster_cookieМетод

cluster_cookie() -> cookie

Возвращает куку-ключ кластера.

исходный код

Distributed.cluster_cookieМетод

cluster_cookie(cookie) -> cookie

Устанавливает переданный ключ в качестве куки-ключа кластера, затем возвращает его.

исходный код

Интерфейс менеджера кластера

Этот интерфейс предоставляет механизм для запуска и управления рабочими процессами Julia в различных кластерных средах. В Base есть два типа менеджеров: LocalManager, для запуска дополнительных рабочих процессов на том же хосте, и SSHManager, для запуска на удалённых хостах через ssh. TCP/IP сокеты используются для подключения и передачи сообщений между процессами. Менеджеры кластеров могут предоставлять другой транспорт.

Distributed.ClusterManagerТип

ClusterManager

Базовый тип для менеджеров кластеров, которые управляют рабочими процессами как кластером. Менеджеры кластеров реализуют способы добавления, удаления и взаимодействия с рабочими процессами. SSHManager и LocalManager являются подтипами этого типа.

исходный код

Distributed.WorkerConfigТип

WorkerConfig

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

  • io – соединение для доступа к рабочему процессу (подтип IO или Nothing)
  • host – адрес хоста (либо AbstractString или Nothing)
  • port – порт на хосте для подключения к рабочему процессу (либо Int или Nothing)

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

  • count – количество рабочих процессов, которые необходимо запустить на хосте
  • exename – путь к исполняемому файлу Julia на хосте, по умолчанию "$(Sys.BINDIR)/julia" или "$(Sys.BINDIR)/julia-debug"
  • exeflags – флаги, которые нужно использовать при удалённом запуске Julia

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

Некоторые поля используются менеджерами SSHManager и аналогичными:

  • tunnel – true (использовать туннелирование), false (не использовать туннелирование) или nothing (использовать значения по умолчанию для менеджера)
  • bind_addr – адрес на удалённом хосте для привязки
  • sshflags – флаги, используемые при установлении SSH-соединения
  • max_parallel – максимальное количество рабочих процессов, которые можно подключить параллельно к хосту

Некоторые поля используются как менеджерами LocalManager, так и SSHManager.

  • connect_at – определяет, является ли это вызов для взаимодействия между рабочими процессами или между драйвером и рабочим процессом
  • process – процесс, который будет подключён (обычно менеджер назначит это при вызове addprocs)
  • ospid – идентификатор процесса согласно ОС хоста, используется для прерывания рабочих процессов
  • environ – частный словарь, используемый для хранения временной информации локальными/SSH менеджерами
  • ident – рабочий процесс, определённый менеджером ClusterManager
  • connect_idents – список идентификаторов рабочих процессов, к которым должен подключиться рабочий процесс, если используется пользовательская топология
  • enable_threaded_blas – true, false, или nothing, использовать ли многопоточное BLAS на рабочих процессах
исходный код

Distributed.launchФункция

launch(manager::ClusterManager, params::Dict, launched::Array, launch_ntfy::Condition)

Реализуется менеджерами кластеров. Для каждого рабочего процесса Julia, запущенного этой функцией, она должна добавить запись WorkerConfig в launched и уведомить launch_ntfy. Функция ДОЛЖНА завершиться после запуска всех рабочих процессов, запрошенных в manager. params – словарь всех ключевых аргументов, с которыми был вызван addprocs.

исходный код

Distributed.manageФункция

manage(manager::ClusterManager, id::Integer, config::WorkerConfig. op::Symbol)

Реализуется менеджерами кластеров. Она вызывается на главном процессе в течение жизни рабочего процесса с соответствующими значениями op.

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

Base.killМетод

kill(manager::ClusterManager, pid::Int, config::WorkerConfig)

Реализуется менеджерами кластеров. Вызывается в главном процессе, методом rmprocs. Он должен заставить удалённый рабочий процесс, указанный параметром pid, завершиться. kill(manager::ClusterManager.....) выполняет удалённую exit() операцию над pid.

исходный код

Sockets.connectМетод

connect(manager::ClusterManager, pid::Int, config::WorkerConfig) -> (instrm::IO, outstrm::IO)

Реализуется менеджерами кластеров, использующими пользовательские транспортные средства. Он должен установить логическое соединение с рабочим процессом с идентификатором pid, указанным параметром config, и вернуть пару объектов IO. Сообщения от pid в текущий процесс будут читаться из instrm, а сообщения, которые нужно отправить pid, будут записываться в outstrm. Реализация пользовательского транспортного средства должна гарантировать, что сообщения доставляются полностью и в правильном порядке. connect(manager::ClusterManager.....) устанавливает соединения TCP/IP между рабочими процессами.

исходный код

Distributed.init_workerФункция

init_worker(cookie::AbstractString, manager::ClusterManager=DefaultClusterManager())

Вызывается менеджерами кластеров, реализующими пользовательские транспортные средства. Она инициализирует недавно запущенный процесс как рабочий процесс. Аргумент командной строки --worker[=<cookie>] имеет эффект инициализации процесса как рабочего процесса, используя TCP/IP сокеты для транспорта. cookie является cluster_cookie.

исходный код

Distributed.start_workerФункция

start_worker([out::IO=stdout], cookie::AbstractString=readline(stdin); close_stdin::Bool=true, stderr_to_stdout::Bool=true)

start_worker — это внутренняя функция, которая является точкой входа по умолчанию для рабочих процессов, подключающихся через TCP/IP. Она устанавливает процесс как рабочий процесс кластера Julia.

Информация о хосте:порту записывается в поток out (по умолчанию — stdout).

Функция читает куки из stdin, если это необходимо, и слушает на свободном порте (или, если указан, на порте в опции командной строки --bind-to) и планирует задачи для обработки входящих TCP-соединений и запросов. Она также (при необходимости) закрывает stdin и перенаправляет stderr в stdout.

Функция не возвращает значение.

исходный код

Distributed.process_messagesФункция

process_messages(r_stream::IO, w_stream::IO, incoming::Bool=true)

Вызывается менеджерами кластеров, использующими пользовательские транспортные средства. Она должна вызываться, когда реализация пользовательского транспортного средства получает первое сообщение от удалённого рабочего процесса. Пользовательская реализация транспортного средства должна управлять логическим соединением с удалённым рабочим процессом и предоставлять два объекта IO, один для входящих сообщений, а другой — для сообщений, адресованных удалённому рабочему процессу. Если incoming равно true, соединение инициировал удалённый узел. Та из пары, которая инициирует соединение, отправляет куки кластера и номер версии Julia для выполнения процедуры аутентификации.

См. также cluster_cookie.

исходный код

© 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/stdlib/Distributed/

Spec-Zone.ru

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