Spec-Zone.ru › Julia 1.7

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

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

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

  • 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 (bash, sh и т. д.). По умолчанию.

    • 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 не будут удалены.
  • Если все рабочие процессы не будут завершены за указанные 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.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)

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

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 в разных кластерных средах. В Базе есть два типа менеджеров: 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–2021 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.7.0/stdlib/Distributed/

Spec-Zone.ru

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