Spec-Zone.ru › Julia 1.0

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

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(e)), 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([x])

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

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

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

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

source
wait(r::Future)

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

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

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

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

Вариант 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.

source

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
source

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
source

Distributed.remote_doМетод

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

Вариант remote_do(f, pid, ....), использующий WorkerPool. Ожидает и берёт свободный рабочий узел из 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 в разных средах кластеров. В Base есть два типа менеджеров: 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). Менеджер должен сигнализировать соответствующему рабочему процессу об прерывании.
  • с :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 (по умолчанию — стандартный вывод).

Функция закрывает 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.0.4/stdlib/Distributed/

Spec-Zone.ru

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