Spec-Zone.ru › Julia 0.7

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

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

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

исходный код

Distributed.nworkersФункция

nworkers()

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

исходный код

Distributed.procsМетод

procs()

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

исходный код

Distributed.procsМетод

procs(pid::Integer)

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

исходный код

Distributed.workersФункция

workers()

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

исходный код

Distributed.rmprocsФункция

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

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

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

исходный код

Distributed.interruptФункция

interrupt(pids::Integer...)

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

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

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

исходный код

Distributed.myidФункция

myid()

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

исходный код

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 обрабатывает элементы из пакета с использованием нескольких одновременных задач.

END_OF_DOCUMENT_MARKER

Любая ошибка останавливает 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))
исходный код

Distributed.RemoteExceptionТип

RemoteException(captured)

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

исходный код

Distributed.FutureТип

Future(pid::Integer=myid())

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

исходный код

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.waitФункция

wait([x])

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

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

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

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

исходный код
wait(r::Future)

Дождаться появления значения для указанного будущего.

исходный код
wait(r::RemoteChannel, args...)

Дождаться появления значения в указанном удалённом канале.

исходный код

Base.fetchМетод

fetch(x)

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

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

    исходный код

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

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

    Distributed.WorkerPoolТип

    WorkerPool(workers::Vector{Int})

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

    исходный код

    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 (по умолчанию).

    исходный код

    Distributed.clear!Метод

    clear!(pool::CachingPool) -> pool

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

    исходный код

    Distributed.remoteФункция

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

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

    исходный код

    Distributed.remotecallМетод

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

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

    исходный код

    Distributed.remotecall_waitМетод

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

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

    исходный код

    Distributed.remotecall_fetchМетод

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

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

    исходный код

    Distributed.remote_doМетод

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

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

    исходный код

    Base.timedwaitФункция

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

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

    исходный код

    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
    исходный код

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

    Информация host:port записывается в поток 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/v0.7.0/stdlib/Distributed/

    Spec-Zone.ru

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