Spec-Zone.ru › Julia 1.4

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

Distributed.addprocsФункция

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

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

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

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

Для запуска рабочих процессов без блокировки REPL или содержащей функции, если запуск рабочих процессов происходит программно, выполните addprocs в отдельной задаче.

Примеры

# On busy clusters, call `addprocs` asynchronously
t = @async addprocs(...)
# Utilize workers as and when they come online
if nprocs() > 1   # Ensure at least one new worker is available
   ....   # perform distributed execution
end
# Retrieve newly launched worker IDs, or any error messages
if istaskdone(t)   # Check if `addprocs` has completed to ensure `fetch` doesn't block
    if nworkers() == N
        new_pids = fetch(t)
    else
        fetch(t)
    end
  end
исходный код
addprocs(machines; tunnel=false, sshflags=``, max_parallel=10, kwargs...) -> List of process identifiers

Добавляет процессы на удалённые машины через SSH. Требуется julia для установки в одном и том же месте на каждом узле или для доступности через общую файловую систему.

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

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

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

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

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

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

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

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

  • dir: задаёт рабочую директорию на рабочих процессах. По умолчанию текущая директория узла (определена по pwd())

  • enable_threaded_blas: если true , то BLAS будет работать на нескольких потоках в добавленных процессах. По умолчанию false.

  • exename: имя исполняемого файла julia. По умолчанию "$(Sys.BINDIR)/julia" или "$(Sys.BINDIR)/julia-debug" в зависимости от случая.

  • exeflags: дополнительные флаги, передаваемые рабочим процессам.

  • topology: Определяет, как рабочие процессы подключаются друг к другу. Отправка сообщения между неподключенными рабочими процессами приводит к ошибке.

    • topology=:all_to_all: Все процессы подключены друг к другу. По умолчанию.

    • topology=:master_worker: Подключение только к рабочим процессам осуществляется через драйверский процесс, т.е. pid 1. Рабочие процессы не подключены друг к другу.

    • topology=:custom: Метод launch менеджера кластера определяет топологию соединения через поля ident и connect_idents в WorkerConfig. Рабочий процесс с идентификатором кластера ident будет подключен ко всем рабочим процессам, указанным в connect_idents.

  • lazy: Применимо только с topology=:all_to_all. Если true, соединения между рабочими процессами устанавливаются лениво, т.е. они устанавливаются при первом удалённом вызове между рабочими процессами. По умолчанию true.

Переменные среды:

Если главный процесс не сможет установить соединение с запущенным рабочим процессом в течение 60,0 секунд, рабочий процесс рассматривает это как критическую ситуацию и завершается. Этот таймаут можно контролировать с помощью переменной среды JULIA_WORKER_TIMEOUT. Значение JULIA_WORKER_TIMEOUT в главном процессе определяет количество секунд, в течение которого запущенный рабочий процесс ожидает установления соединения.

исходный код
addprocs(; kwargs...) -> List of process identifiers

Эквивалентно addprocs(Sys.CPU_THREADS; kwargs...)

Обратите внимание, что рабочие процессы не выполняют скрипт запуска .julia/config/startup.jl, а также не синхронизируют своё глобальное состояние (такое как глобальные переменные, определения новых методов и загруженные модули) с какими-либо другими запущенными процессами.

исходный код
addprocs(np::Integer; restrict=true, kwargs...) -> List of process identifiers

Запускает рабочие процессы с помощью встроенного LocalManager, который запускает рабочие процессы только на локальном узле. Это можно использовать для эффективного использования нескольких ядер. addprocs(4) добавит 4 процесса на локальную машину. Если restrict — true, привязка ограничена 127.0.0.1. Ключевые аргументы dir, exename, exeflags, topology, lazy и enable_threaded_blas имеют тот же эффект, что и описано для addprocs(machines).

исходный код

Distributed.nprocsФункция

nprocs()

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

Примеры

julia> nprocs()
3

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

Distributed.nworkersФункция

nworkers()

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

Примеры

$ julia -p 5

julia> nprocs()
6

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

Distributed.procsМетод

procs()

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

Примеры

$ julia -p 5

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

Distributed.procsМетод

procs(pid::Integer)

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

исходный код

Distributed.workersФункция

workers()

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

Примеры

$ julia -p 5

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

Distributed.rmprocsФункция

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

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

Аргумент waitfor определяет время ожидания завершения рабочих процессов:

  • Если не указано, rmprocs будет ждать, пока все запрошенные pids не будут удалены.
  • Исключение 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)

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

исходный код

Distributed.FutureТип

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

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

исходный код

Distributed.RemoteChannelТип

RemoteChannel(pid::Integer=myid())

Создать ссылку на Channel{Any}(1) на процессе pid. По умолчанию pid — это текущий процесс.

RemoteChannel(f::Function, pid::Integer=myid())

Создайте ссылки на удалённые каналы определённого размера и типа. f — это функция, которая, когда выполняется на pid, должна вернуть реализацию AbstractChannel.

Например, RemoteChannel(()->Channel{Int}(10), pid), вернёт ссылку на канал типа Int и размера 10 на процессе pid.

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

исходный код

Base.fetchМетод

fetch(x::Future)

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

исходный код

Base.fetchМетод

fetch(c::RemoteChannel)

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

исходный код

Distributed.remotecallМетод

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

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

исходный код

Distributed.remotecall_waitМетод

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

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

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

исходный код

Distributed.remotecall_fetchМетод

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

Выполнить fetch(remotecall(...)) в одном сообщении. Параметры-ключевые слова, если таковые имеются, передаются в f. Любые удалённые исключения перехватываются в RemoteException и выбрасываются.

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

Примеры

$ julia -p 2

julia> remotecall_fetch(sqrt, 2, 4)
2.0

julia> remotecall_fetch(sqrt, 2, -4)
ERROR: On worker 2:
DomainError with -4.0:
sqrt will only return a complex result if called with a complex argument. Try sqrt(Complex(x)).
...
исходный код

Distributed.remote_doМетод

remote_do(f, id::Integer, args...; kwargs...) -> nothing

Асинхронно выполняет f на рабочем процессе id. В отличие от remotecall, он не сохраняет результат вычисления и нет способа дождаться его завершения.

Успешное выполнение указывает, что запрос принят для выполнения на удалённом узле.

Хотя последовательные remotecall на один и тот же рабочий процесс сериализуются в порядке их вызова, порядок выполнения на удалённом рабочем процессе не определён. Например, remote_do(f1, 2); remotecall(f2, 2); remote_do(f3, 2) сериализует вызов f1, за которым следует f2 и f3 в таком порядке. Однако не гарантируется, что f1 будет выполнено до f3 на рабочем процессе 2.

Любые исключения, сгенерированные f, выводятся в stderr на удалённом рабочем процессе.

Параметры-ключевые слова, если таковые имеются, передаются в f.

исходный код

Base.put!Метод

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

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

исходный код

Base.put!Метод

put!(rr::Future, v)

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

исходный код

Base.take!Метод

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

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

исходный код

Base.isreadyМетод

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

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

исходный код

Base.isreadyМетод

isready(rr::Future)

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

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

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

Distributed.AbstractWorkerPoolТип

AbstractWorkerPool

Базовый тип для пулов рабочих процессов, таких как WorkerPool и CachingPool. 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})

Реализация пула с кэшированием. remote, remotecall_fetch, pmap (и другие удалённые вызовы, выполняющие функции на удалённых узлах) выигрывают от кэширования сериализованных/десериализованных функций на рабочих узлах, особенно замыканий (которые могут захватывать большое количество данных).

Удаленный кэш хранится в течение всего срока жизни возвращаемого объекта CachingPool. Чтобы очистить кэш раньше, используйте clear!(pool).

Для глобальных переменных в замыкании захватываются только связи, а не данные. Можно использовать let для захвата глобальных данных.

Примеры

const foo = rand(10^8);
wp = CachingPool(workers())
let foo = foo
    pmap(wp, i -> sum(foo) + i, 1:100);
end

Вышеперечисленное будет передавать foo только один раз на каждый рабочий узел.

исходный код

Distributed.default_worker_poolФункция

default_worker_pool()

WorkerPool, содержащий свободные workers — используется remote(f) и pmap (по умолчанию).

Примеры

$ julia -p 3

julia> default_worker_pool()
WorkerPool(Channel{Int64}(sz_max:9223372036854775807,sz_curr:3), Set([4, 2, 3]), RemoteChannel{Channel{Any}}(1, 1, 4))
исходный код

Distributed.clear!Метод

clear!(pool::CachingPool) -> pool

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

исходный код

Distributed.remoteФункция

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

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

исходный код

Distributed.remotecallМетод

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

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

В этом примере задача выполнялась на процессе с идентификатором 2, вызвана с процесса с идентификатором 1.

исходный код

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

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

Distributed.remote_doМетод

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

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

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

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

@sync @distributed for var = range
    body
end
исходный код

Distributed.@everywhereМакрос

@everywhere [procs()] expr

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

@everywhere bar = 1

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

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

foo = 1
@everywhere bar = $foo

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

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

исходный код

Distributed.clear!Метод

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

Очищает глобальные привязки в модулях, инициализируя их значением nothing. syms должен иметь тип Symbol или коллекцию Symbols. 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 – адрес хоста (либо AbstractString, либо Nothing)
  • port – порт на хосте, используемый для подключения к рабочему процессу (либо Int, либо Nothing)

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

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

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

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

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

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

  • connect_at – определяет, является ли это вызовом между рабочими процессами или между драйвером и рабочим процессом
  • process – процесс, который будет подключён (обычно менеджер назначит его во время addprocs)
  • ospid – идентификатор процесса согласно операционной системе хоста, используемый для прерывания процессов рабочих процессов
  • environ – приватный словарь, используемый для хранения временной информации менеджерами 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)

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

исходный код

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

Spec-Zone.ru

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