Spec-Zone.ru › Julia 1.6

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

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. См. exename для установки пути к установке julia на удалённых машинах.

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

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

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

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

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

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

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

  • ssh: имя или путь к исполняемому файлу SSH клиента, используемого для запуска рабочих процессов. По умолчанию "ssh".

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

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

  • shell: задаёт тип оболочки, к которой подключается ssh на рабочих процессах.

    • shell=:posix: POSIX-совместимая Unix/Linux оболочка (bash, sh и т. д.). По умолчанию.

    • shell=:wincmd: Microsoft Windows cmd.exe.

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

  • env: предоставляет массив пар строк, таких как env=["JULIA_DEPOT_PATH"=>"/depot"], чтобы запросить установку переменных окружения на удалённой машине. По умолчанию автоматически передаётся только переменная окружения JULIA_WORKER_TIMEOUT из локальной в удалённую среду.

  • cmdline_cookie: передаёт куки аутентификации через опцию командной строки --worker. Более безопасное поведение по умолчанию, которое передаёт куки через ssh stdio, может зависнуть на Windows рабочих процессах, использующих более старые (до ConPTY) версии Julia или Windows. В этом случае cmdline_cookie=true предоставляет обходной путь.

Ключевые аргументы ssh, shell, env и cmdline_cookie были добавлены в Julia 1.6.

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

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

julia> nprocs()
3

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

Distributed.procsМетод

procs()

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

Примеры

$ julia -p 2

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

Distributed.procsМетод

procs(pid::Integer)

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

исходный код

Distributed.workersФункция

workers()

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

Примеры

$ julia -p 2

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

Distributed.RemoteExceptionТип

RemoteException(captured)

Исключения при удаленных вычислениях捕获并重新抛出到本地。 A RemoteException содержит pid рабочего узла и captured исключение. A CapturedException captures удаленное исключение и сериализуемую форму стека вызовов, когда было вызвано исключение.

исходный код

Distributed.FutureТип

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

A 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, которое captures удаленное исключение и трассировку стека.

исходный код

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. AbstractWorkerPool должен реализовывать:

  • 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(i -> sum(foo) + i, wp, 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)

В этом примере задача была запущена на узле с pid 2, вызвана с pid 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

Аргумент :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
исходный код

Distributed.@everywhereМакрос

@everywhere [procs()] expr

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

@everywhere bar = 1

определит Main.bar на всех текущих процессах. Любые добавленные позже процессы (например, с помощью addprocs()) не будут иметь это выражение.

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

foo = 1
@everywhere bar = $foo

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

Аналогично вызову remotecall_eval(Main, procs, expr), но с двумя дополнительными функциями:

- `using` and `import` statements run on the calling process first, to ensure
  packages are precompiled.
- The current source file path used by `include` is propagated to other processes.
исходный код

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.ClusterManagerТип

ClusterManager

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

исходный код

Distributed.WorkerConfigТип

WorkerConfig

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

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

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

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

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

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

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

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

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

Distributed.launchФункция

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

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

исходный код

Distributed.manageФункция

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

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

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

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); close_stdin::Bool=true, stderr_to_stdout::Bool=true)

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

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

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

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

исходный код

Distributed.process_messagesФункция

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

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

См. также cluster_cookie.

исходный код

Distributed.default_addprocs_paramsФункция

default_addprocs_params(mgr::ClusterManager) -> Dict{Symbol, Any}

Реализуется менеджерами кластеров. Параметры по умолчанию, передаваемые при вызове addprocs(mgr). Минимальный набор опций доступен при вызове default_addprocs_params()

исходный код

© 2009–2021 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.6.0/stdlib/Distributed/

Spec-Zone.ru

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