Spec-Zone.ru › Julia 1.3

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

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, оно запустит столько рабочих процессов, сколько потоков CPU на конкретном узле.

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

  • 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 не будут удалены.
  • Если все рабочие процессы не могут быть завершены до запрошенных waitfor секунд, генерируется ErrorException.
  • Со значением 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))
источник

Distributed.RemoteExceptionТип

RemoteException(captured)

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

источник

Distributed.FutureТип

Future(w::Int, rrid::RRID, v::Union{Some, Nothing}=nothing)

Future — это заглушка для одного вычисления с неизвестным статусом завершения и временем. Для нескольких потенциальных вычислений см. RemoteChannel. См. remoteref_id для определения AbstractRemoteRef.

источник

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 — текущий процесс.

источник

Base.fetchМетод

fetch(x::Future)

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

источник

Base.fetchМетод

fetch(c::RemoteChannel)

Дожидается и получает значение из RemoteChannel. Брошенные исключения такие же, как и для Future. Не удаляет полученный элемент.

источник

Distributed.remotecallМетод

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

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

источник

Distributed.remotecall_waitМетод

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

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

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

источник

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)).
...
источник

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.

источник

Base.put!Метод

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

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

источник

Base.put!Метод

put!(rr::Future, v)

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

источник

Base.take!Метод

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

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

источник

Base.isreadyМетод

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

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

исходный код

Base.isreadyМетод

isready(rr::Future)

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

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

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

Distributed.AbstractWorkerPoolТип

AbstractWorkerPool

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

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

Стандартные реализации выше (в AbstractWorkerPool) требуют полей channel::Channel{Int} workers::Set{Int}, где channel содержит идентификаторы свободных рабочих, а 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

Вариант WorkerPool remotecall(f, pid, ....). Ожидает и берёт свободный рабочий узел из 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)

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

исходный код

Distributed.remotecall_waitМетод

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

Вариант WorkerPool remotecall_wait(f, pid, ....). Ожидает и берёт свободный рабочий узел из 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

Вариант WorkerPool remotecall_fetch(f, pid, ....). Ожидает и берёт свободный рабочий узел из 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

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

исходный код

Distributed.@spawnatМакрос

@spawnat p expr

Создаёт замыкание вокруг выражения и асинхронно выполняет его на процессе p. Возвращает Future результата. Если p — это символьная литеральная строка :any, система автоматически выберет процессор для использования.

Примеры

julia> addprocs(3);

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

julia> fetch(f)
2

julia> f = @spawnat :any myid()
Future(3, 1, 7, nothing)

julia> fetch(f)
3
Julia 1.3

Аргумент :any доступен начиная с Julia 1.3.

исходный код

Distributed.@fetchМакрос

@fetch expr

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

Примеры

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

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

@sync @distributed for var = range
    body
end
source

Distributed.@everywhereМакрос

@everywhere [procs()] expr

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

@everywhere bar = 1

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

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

foo = 1
@everywhere bar = $foo

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

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

source

Distributed.clear!Метод

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

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

Исключение генерируется, если запрошено очищение глобальной константы.

source

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 значения удалённой ссылки.

source

Distributed.channel_from_idФункция

channel_from_id(id) -> c

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

source

Distributed.worker_id_from_socketФункция

worker_id_from_socket(s) -> pid

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

source

Distributed.cluster_cookieМетод

cluster_cookie() -> cookie

Возвращает куки кластера.

source

Distributed.cluster_cookieМетод

cluster_cookie(cookie) -> cookie

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

source

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

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

Distributed.ClusterManagerТип

ClusterManager

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

source

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 на рабочих узлах
source

Distributed.launchФункция

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

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

source

Distributed.manageФункция

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

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

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

source

Base.killМетод

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

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

source

Sockets.connectМетод

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

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

source

Distributed.init_workerФункция

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

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

source

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.

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

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

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

source

Distributed.process_messagesФункция

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

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

См. также cluster_cookie.

source

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

Spec-Zone.ru

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