Spec-Zone.ru › Julia 1.8

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

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, будет запущено столько рабочих процессов, сколько потоков 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 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()
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 – опция перенаправления, используемая для опции -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)

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

Spec-Zone.ru

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