Распределённые вычисления
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
endaddprocs(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, будет запущено столько рабочих процессов, сколько потоков CPU на конкретном хосте.
Ключевые аргументы:
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 оболочка (sh, ksh, bash, dash, zsh и т.д.). По умолчанию.shell=:csh: Unix C оболочка (csh, tcsh).shell=:wincmd: Microsoft Windowscmd.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: Только драйверный процесс, т.е.pid1, подключается к рабочим процессам. Рабочие процессы не подключаются друг к другу.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()
2-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)
Исключения при удаленных вычислениях捕获 и повторно генерируются локально. 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) errormonitor(@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
Вариант 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
Указанный диапазон разбивается и выполняется локально на всех рабочих процессах. Если указана функция-редуктор, @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– опция перенаправления, используемая для опции-Lssh -
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 (по умолчанию — стандартный вывод).
Функция считывает куки из стандартного ввода, если это необходимо, и прослушивает свободный порт (или, если указан, порт в опции командной строки --bind-to) и планирует задачи для обработки входящих TCP-соединений и запросов. Также (при необходимости) закрывает стандартный ввод и перенаправляет стандартную ошибку в стандартный вывод.
Функция не возвращает значение.
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–2022 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.8/stdlib/Distributed/