Spec-Zone.ru › Julia 1.9

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

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

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 в качестве количества рабочих процессов запустит столько рабочих процессов, сколько потоков CPU на удалённом хосте.

Примеры:

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: оболочка POSIX-совместимая Unix/Linux (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 не будут удалены.
  • Если все рабочие процессы не могут быть завершены до истечения запрошенных waitfor секунд, генерируется исключение ErrorException.
  • При значении 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(x::Any)

Вернуть x.

исходный код
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::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 (по умолчанию). Если не указано явно через default_worker_pool!(pool), пул рабочих процессов по умолчанию инициализируется как 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)

В этом примере задача выполнялась на процессе с идентификатором 2, вызванная из процесса с идентификатором 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

Возвращает cookie кластера.

Distributed.cluster_cookieМетод

cluster_cookie(cookie) -> 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 (по умолчанию — stdout).

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

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

Distributed.process_messagesФункция

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

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

См. также cluster_cookie.

Distributed.default_addprocs_paramsФункция

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

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

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

Spec-Zone.ru

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