Spec-Zone.ru › Julia 1.1

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

Distributed.addprocsФункция

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

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

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

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

исходный код
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 не будут удалены.
  • Исключение 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.waitФункция

wait(r::Future)

Ожидайте, пока значение станет доступным для указанного Future.

source
wait(r::RemoteChannel, args...)

Ожидайте, пока значение станет доступным в указанном RemoteChannel.

source
wait([x])

Заблокировать текущую задачу до тех пор, пока не произойдет какое-либо событие, в зависимости от типа аргумента:

  • Channel: Ожидайте, пока значение не будет добавлено в канал.
  • Condition: Ожидайте notify для условия.
  • Process: Ожидайте завершения процесса или цепочки процессов. Поле exitcode процесса может использоваться для определения успеха или неудачи.
  • Task: Ожидайте завершения Task. Если задача завершается с исключением, исключение передаётся (перебрасывается) в задачу, вызвавшую wait.
  • RawFD: Ожидайте изменений в дескрипторе файла (см. пакет FileWatching).

Если аргумент не передан, задача блокируется на неопределённый период. Задачу можно перезапустить только с помощью явного вызова schedule или yieldto.

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

source

Base.fetchМетод

fetch(x)

Ожидает и извлекает значение из x в зависимости от типа x:

  • Future: Ожидает и получает значение Future. Извлечённое значение кешируется локально. Дальнейшие вызовы fetch к той же ссылке возвращают кешированное значение. Если удалённое значение является исключением, выбрасывает RemoteException, который перехватывает удалённое исключение и трассировку стека.
  • 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, так как они присваиваются только один раз.

source

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
source

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))
source

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 только один раз на каждый рабочий узел.

source

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))
source

Distributed.clear!Метод

clear!(pool::CachingPool) -> pool

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

source

Distributed.remoteФункция

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

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

source

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)

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

source

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
source

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
source

Distributed.remote_doМетод

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

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

source

Base.timedwaitФункция

timedwait(testcb::Function, secs::Float64; pollint::Float64=0.1)

Ожидает, пока testcb вернёт true или в течение secs секунд, в зависимости от того, что наступит раньше. testcb опрошивается каждые pollint секунд.

source

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
source

Distributed.@spawnatМакрос

@spawnat

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

Примеры

julia> addprocs(1);

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

julia> fetch(f)
2
source

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
source

Distributed.@fetchfromМакрос

@fetchfrom

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

Примеры

julia> addprocs(3);

julia> @fetchfrom 2 myid()
2

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

Base.@asyncМакрос

@async

Оборачивает выражение в Task и добавляет его в очередь планировщика локальной машины.

исходный код

Base.@syncМакрос

@sync

Ожидает завершения всех лексически вложенных обращений к @async, @spawn, @spawnat и @distributed. Все исключения, сгенерированные вложенными асинхронными операциями, собираются и выбрасываются как CompositeException.

исходный код

Distributed.@distributedМакрос

@distributed

Распределённый по памяти, параллельный цикл for следующего вида:

@distributed [reducer] for var = range
    body
end

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

Обратите внимание, что без функции-редуктора @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

Future и RemoteChannel идентифицируются полями:

  • 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 в различных кластерных средах. В базовом модуле существуют два типа менеджеров: LocalManager, для запуска дополнительных рабочих процессов на том же хосте, и SSHManager, для запуска на удалённых хостах через ssh. Для связи и передачи сообщений между процессами используются сокеты TCP/IP. Возможно предоставление менеджерами кластера других видов транспорта.

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 для целей очистки.
исходный код

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)

Реализуется менеджерами кластеров с использованием пользовательских протоколов. Она должна установить логическое соединение с рабочим процессом с id 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))

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

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

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

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

источник

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.1.1/stdlib/Distributed/

Spec-Zone.ru

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