Spec-Zone.ru › Julia 1.10

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

Инструменты для распределённой параллельной обработки.

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 — это вектор «спецификаций машины», которые задаются строками вида [user@]host[:port] [bind_addr[:port]]. user по умолчанию равен текущему пользователю, а port — стандартному порту SSH. Если [bind_addr[:port]] указан, другие рабочие процессы будут подключаться к этому рабочему процессу по указанному bind_addr и port.

Возможен запуск нескольких процессов на удалённом хосте, используя кортеж в векторе machines или форме (machine_spec, count), где count — количество рабочих процессов, которые необходимо запустить на указанном хосте. Передача :auto в качестве количества рабочих процессов запустит столько рабочих процессов, сколько потоков ЦП на удалённом хосте.

Примеры:

addprocs([
    "remote1",               # one worker on 'remote1' logging in with the current username
    "user@remote2",          # one worker on 'remote2' logging in with the 'user' username
    "user@remote3:2222",     # specifying SSH port to '2222' for 'remote3'
    ("user@remote4", 4),     # launch 4 workers on 'remote4'
    ("user@remote5", :auto), # launch as many workers as CPU threads on 'remote5'
])

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

  • 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: оболочка Unix/Linux, совместимая с POSIX (sh, ksh, bash, dash, zsh и т. д.). По умолчанию.

    • shell=:csh: оболочка Unix C shell (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" по ситуации. Рекомендуется использовать общую версию Julia на всех удалённых машинах, поскольку иначе могут возникнуть ошибки сериализации и распространения кода.

  • 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(np::Integer=Sys.CPU_THREADS; restrict=true, kwargs...) -> List of process identifiers

Запускает np рабочих процессов на локальном хосте с помощью встроенного LocalManager.

Локальные рабочие процессы наследуют текущую среду пакетов (то есть активный проект, LOAD_PATH и DEPOT_PATH) из главного процесса.

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

  • restrict::Bool: если true (по умолчанию) привязка ограничена 127.0.0.1.
  • dir, exename, exeflags, env, topology, lazy, enable_threaded_blas: тот же эффект, что и для SSHManager, см. документацию для addprocs(machines::AbstractVector).

Наследование среды пакетов и ключевой аргумент env были добавлены в Julia 1.9.

исходный код

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

ProcessExitedException(worker_id::Int)

После завершения процесса Julia-клиента дальнейшие попытки обратиться к умершему дочернему процессу вызовут это исключение.

исходный код

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. Не удаляет извлечённый элемент.

исходный код
fetch(x::Any)

Вернуть x.

исходный код

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 was called with a negative real argument but 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::Union{Vector{Int},AbstractRange{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))

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

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

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

Примеры

$ julia -p 3

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

Distributed.clear!Метод

clear!(pool::CachingPool) -> pool

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

исходный код

Distributed.remoteФункция

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

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

исходный код

Distributed.remotecallМетод

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

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

Примеры

$ julia -p 3

julia> wp = WorkerPool([2, 3]);

julia> A = rand(3000);

julia> f = remotecall(maximum, wp, A)
Future(2, 1, 6, nothing)

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

исходный код

Distributed.remotecall_waitМетод

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

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

Примеры

$ julia -p 3

julia> wp = WorkerPool([2, 3]);

julia> A = rand(3000);

julia> f = remotecall_wait(maximum, wp, A)
Future(3, 1, 9, nothing)

julia> fetch(f)
0.9995177101692958
исходный код

Distributed.remotecall_fetchМетод

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

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

Примеры

$ julia -p 3

julia> wp = WorkerPool([2, 3]);

julia> A = rand(3000);

julia> remotecall_fetch(maximum, wp, A)
0.9995177101692958
исходный код

Distributed.remote_doМетод

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

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

исходный код

Distributed.@spawnatМакрос

@spawnat p expr

Создаёт замыкание вокруг выражения и асинхронно выполняет замыкание на процессе p. Возвращает Future с результатом. Если p — это символьное значение :any, система автоматически выбирает процессор для использования.

Примеры

julia> addprocs(3);

julia> f = @spawnat 2 myid()
Future(2, 1, 3, nothing)

julia> fetch(f)
2

julia> f = @spawnat :any myid()
Future(3, 1, 7, nothing)

julia> fetch(f)
3

Аргумент :any доступен начиная с Julia 1.3.

исходный код

Distributed.@fetchМакрос

@fetch expr

Эквивалентно fetch(@spawnat :any expr). См. fetch и @spawnat.

Примеры

julia> addprocs(3);

julia> @fetch myid()
2

julia> @fetch myid()
3

julia> @fetch myid()
4

julia> @fetch myid()
2
исходный код

Distributed.@fetchfromМакрос

@fetchfrom

Эквивалентно fetch(@spawnat p expr). См. fetch и @spawnat.

Примеры

julia> addprocs(3);

julia> @fetchfrom 2 myid()
2

julia> @fetchfrom 4 myid()
4
исходный код

Distributed.@distributedМакрос

@distributed

Параллельный цикл for в распределённой памяти, вида:

@distributed [reducer] for var = range
    body
end

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

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

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

Distributed.@everywhereМакрос

@everywhere [procs()] expr

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

@everywhere bar = 1

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

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

foo = 1
@everywhere bar = $foo

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

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

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

Distributed.clear!Метод

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

Очищает глобальные привязки в модулях, инициализируя их значением nothing. syms должен быть типа Symbol или набором Symbol . pids и mod идентифицируют процессы и модуль, в котором нужно повторно инициализировать глобальные переменные. Очищаются только те имена, которые определены в mod.

Исключение генерируется, если требуется очистить глобальную константу.

исходный код

Distributed.remoteref_idФункция

remoteref_id(r::AbstractRemoteRef) -> RRID

Future и RemoteChannel идентифицируются полями:

  • where - указывает узел, на котором фактически существует подлежащий объекту/хранилищу объект, на который ссылается ссылка.

  • whence - указывает узел, с которого была создана удалённая ссылка. Обратите внимание, что это отличается от узла, на котором фактически существует подлежащий объекту, на который ссылается ссылка. Например, вызов RemoteChannel(2) с главного процесса приведет к значению where 2 и значению whence 1.

  • id уникален для всех ссылок, созданных с рабочего процесса, указанного whence.

Вместе whence и id уникально идентифицируют ссылку на всех рабочих процессах.

remoteref_id — это низкоуровневый API, который возвращает объект RRID, который оборачивает whence и id значения удалённой ссылки.

исходный код

Distributed.channel_from_idФункция

channel_from_id(id) -> c

Низкоуровневый API, который возвращает базовое AbstractChannel для id , возвращённого remoteref_id. Вызов допустим только на узле, где существует базовое соединение.

исходный код

Distributed.worker_id_from_socketФункция

worker_id_from_socket(s) -> pid

Низкоуровневый API, который, получив соединение IO или Worker, возвращает pid рабочего процесса, с которым оно соединено. Это полезно при написании пользовательских методов serialize для типа, которые оптимизируют выводимые данные в зависимости от идентификатора принимающего процесса.

исходный код

Distributed.cluster_cookieМетод

cluster_cookie() -> cookie

Возвращает куку для кластера.

исходный код

Distributed.cluster_cookieМетод

cluster_cookie(cookie) -> cookie

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

исходный код

Интерфейс менеджера кластера

Этот интерфейс предоставляет механизм для запуска и управления рабочими процессами Julia в различных кластерных средах. В Base существуют два типа менеджеров: LocalManager, для запуска дополнительных рабочих процессов на том же хосте, и SSHManager, для запуска на удалённых хостах через ssh. Для подключения и передачи сообщений между процессами используются TCP/IP сокеты. Менеджеры кластера могут предоставлять и другой транспорт.

Distributed.ClusterManagerТип

ClusterManager

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

исходный код

Distributed.WorkerConfigТип

WorkerConfig

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

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

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

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

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

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

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

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

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

Distributed.launchФункция

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

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

исходный код

Distributed.manageФункция

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

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

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

Base.killМетод

kill(manager::ClusterManager, pid::Int, config::WorkerConfig)

Реализуется менеджерами кластеров. Вызывается в главном процессе функцией rmprocs. Она должна заставить удалённый рабочий процесс, указанный в pid, выйти. kill(manager::ClusterManager.....) выполняет удалённую exit() для pid.

исходный код

Sockets.connectМетод

connect(manager::ClusterManager, pid::Int, config::WorkerConfig) -> (instrm::IO, outstrm::IO)

Реализуется менеджерами кластеров с использованием пользовательского транспорта. Она должна установить логическое соединение с рабочим процессом с id pid, указанным в config, и вернуть пару объектов IO. Сообщения от pid до текущего процесса будут читаться из instrm, а сообщения для отправки в pid будут записываться в outstrm. Реализация пользовательского транспорта должна гарантировать, что сообщения передаются и принимаются полностью и в правильном порядке. connect(manager::ClusterManager.....) устанавливает соединения TCP/IP между рабочими процессами.

исходный код

Distributed.init_workerФункция

init_worker(cookie::AbstractString, manager::ClusterManager=DefaultClusterManager())

Вызывается менеджерами кластеров, реализующими пользовательский транспорт. Она инициализирует только что запущенный процесс как рабочий процесс. Флаг командной строки --worker[=<cookie>] инициализирует процесс как рабочий процесс, используя TCP/IP сокеты для транспорта. cookie – это cluster_cookie.

исходный код

Distributed.start_workerФункция

start_worker([out::IO=stdout], cookie::AbstractString=readline(stdin); close_stdin::Bool=true, stderr_to_stdout::Bool=true)

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

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

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

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

исходный код

Distributed.process_messagesФункция

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

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

См. также cluster_cookie.

исходный код

Distributed.default_addprocs_paramsФункция

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

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

исходный код

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

Spec-Zone.ru

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