Распределённые вычисления
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 Windowscmd.exe.
dir: указывает рабочую директорию на рабочих процессах. По умолчанию — текущая директория хоста (определённая с помощьюpwd()).enable_threaded_blas: еслиtrue, то BLAS будет работать в нескольких потоках в добавленных процессах. По умолчанию —false.exename: имя исполняемого файлаjulia. По умолчанию —"$(Sys.BINDIR)/julia"или"$(Sys.BINDIR)/julia-debug"в зависимости от ситуации.exeflags: дополнительные флаги, передаваемые рабочим процессам.-
topology: Определяет, как рабочие процессы подключаются друг к другу. Отправка сообщения между неподключенными рабочими процессами приводит к ошибке.topology=:all_to_all: Все процессы подключены друг к другу. По умолчанию.topology=:master_worker: Только процесс-драйвер, т. е.pid1, подключается к рабочим процессам. Рабочие процессы не подключаются друг к другу.topology=:custom: Методlaunchменеджера кластера определяет топологию подключений через поляidentиconnect_identsвWorkerConfig. Рабочий процесс с идентификатором менеджера кластераidentподключится ко всем рабочим процессам, указанным вconnect_idents.
lazy: Применимо только сtopology=:all_to_all. Еслиtrue, подключения между рабочими процессами устанавливаются лениво, т. е. они устанавливаются при первом удалённом вызове между рабочими процессами. По умолчанию — true.env: предоставьте массив пар строк, напримерenv=["JULIA_DEPOT_PATH"=>"/depot"], чтобы запросить установку переменных окружения на удалённой машине. По умолчанию автоматически передаётся только переменная окруженияJULIA_WORKER_TIMEOUTс локальной среды в удалённую.cmdline_cookie: передать токен аутентификации через параметр командной строки--worker. Более безопасное по умолчанию поведение, передающее токен через ssh stdio, может зависнуть при использовании рабочих процессов Windows, использующих более старые (до ConPTY) версии Julia или Windows. В этом случаеcmdline_cookie=trueпредлагает обходной путь.
Ключевые аргументы ssh, shell, env и cmdline_cookie были добавлены в Julia 1.6.
Переменные окружения:
Если главный процесс не сможет установить соединение с только что запущенным рабочим процессом в течение 60,0 секунд, рабочий процесс расценивает это как критическую ситуацию и завершает работу. Этот таймаут можно настроить с помощью переменной среды JULIA_WORKER_TIMEOUT. Значение JULIA_WORKER_TIMEOUT в главном процессе определяет количество секунд, в течение которого только что запущенный рабочий процесс ожидает установления соединения.
addprocs(; kwargs...) -> List of process identifiers
Эквивалентно addprocs(Sys.CPU_THREADS; kwargs...)
Обратите внимание, что рабочие процессы не выполняют скрипт запуска .julia/config/startup.jl, и не синхронизируют своё глобальное состояние (такое как глобальные переменные, новые определения методов и загруженные модули) ни с одним из других работающих процессов.
addprocs(np::Integer; restrict=true, kwargs...) -> List of process identifiers
Запускает рабочие процессы с помощью встроенного LocalManager, который запускает рабочие процессы только на локальном хосте. Это позволяет использовать несколько ядер. addprocs(4) добавит 4 процесса на локальной машине. Если restrict равно true, привязка ограничена 127.0.0.1. Ключевые аргументы dir, exename, exeflags, topology, lazy и enable_threaded_blas имеют тот же эффект, что и документировано для addprocs(machines).
Distributed.nprocsФункция
nprocs()
Получить количество доступных процессов.
Примеры
julia> nprocs()
3
julia> workers()
2-element Array{Int64,1}:
2
3
исходный код
Distributed.nworkersФункция
nworkers()
Получить количество доступных рабочих процессов. Это на единицу меньше, чем nprocs(). Равно nprocs() , если nprocs() == 1.
Примеры
$ julia -p 2 julia> nprocs() 3 julia> nworkers() 2исходный код
Distributed.procsМетод
procs()
Возвращает список всех идентификаторов процессов, включая pid 1 (который не входит в workers()).
Примеры
$ julia -p 2
julia> procs()
3-element Array{Int64,1}:
1
2
3
исходный код
Distributed.procsМетод
procs(pid::Integer)
Возвращает список всех идентификаторов процессов на том же физическом узле. В частности, возвращаются все рабочие процессы, привязанные к тому же IP-адресу, что и pid.
Distributed.workersФункция
workers()
Возвращает список идентификаторов всех рабочих процессов.
Примеры
$ julia -p 2
julia> workers()
2-element Array{Int64,1}:
2
3
исходный код
Distributed.rmprocsФункция
rmprocs(pids...; waitfor=typemax(Int))
Удаляет указанные рабочие процессы. Обратите внимание, что только процесс 1 может добавлять или удалять рабочие процессы.
Аргумент waitfor указывает, сколько времени ожидать завершения работы рабочих процессов:
- Если не указано,
rmprocsбудет ожидать, пока все запрошенныеpidsне будут удалены. - Если все рабочие процессы не будут завершены за указанные
waitforсекунд, будет выброшено исключениеErrorException. - Со значением
waitfor0 вызов возвращается немедленно, а рабочие процессы планируются для удаления в другой задаче. Возвращается объект запланированной задачи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– опция перенаправления, используемая для опции-Lssh -
bind_addr– адрес на удалённом хосте для привязки -
sshflags– флаги, используемые при установлении SSH-соединения -
max_parallel– максимальное количество работников, которые могут одновременно подключиться к хосту
Некоторые поля используются как для LocalManager, так и для SSHManager:
-
connect_at– определяет, является ли это вызов для связи между работниками или между драйвером и работником -
process– процесс, к которому будет установлено соединение (обычно менеджер назначает его во времяaddprocs) -
ospid– идентификатор процесса согласно операционной системе хоста, используемый для прерывания процессов работников -
environ– приватный словарь, используемый для хранения временной информации менеджерами Local/SSH -
ident– работник, идентифицированный менеджеромClusterManager -
connect_idents– список идентификаторов работников, к которым работник должен подключиться, если используется настраиваемая топология -
enable_threaded_blas–true,false, илиnothing, использовать ли потоковый BLAS на работниках
Distributed.launchФункция
launch(manager::ClusterManager, params::Dict, launched::Array, launch_ntfy::Condition)
Реализуется менеджерами кластеров. Для каждого запущенного Julia-работника этой функцией, она должна добавить запись WorkerConfig в launched и уведомить launch_ntfy. Функция ОБЯЗАНА завершиться, как только все работники, запрошенные manager, будут запущены. params — словарь всех именованных аргументов, с которыми вызывалась addprocs.
Distributed.manageФункция
manage(manager::ClusterManager, id::Integer, config::WorkerConfig. op::Symbol)
Реализуется менеджерами кластеров. Она вызывается на главном процессе во время жизни работника с соответствующими значениями op:
- со значениями
:register/:deregister, когда работник добавляется/удаляется из пула Julia-работников. - со значением
:interrupt, когда вызываетсяinterrupt(workers).ClusterManagerдолжен послать соответствующему работнику сигнал прерывания. - со значением
:finalizeдля целей очистки.
Base.killМетод
kill(manager::ClusterManager, pid::Int, config::WorkerConfig)
Реализуется менеджерами кластеров. Вызывается на главном процессе функцией rmprocs. Она должна заставить удалённый работник, указанный pid, выйти. kill(manager::ClusterManager.....) выполняет удалённый exit() на pid.
Sockets.connectМетод
connect(manager::ClusterManager, pid::Int, config::WorkerConfig) -> (instrm::IO, outstrm::IO)
Реализуется менеджерами кластеров с использованием настраиваемых транспортов. Она должна установить логическое соединение с работником с идентификатором pid, указанным config, и вернуть пару объектов IO. Сообщения от pid текущему процессу будут считываться из instrm, а сообщения для отправки pid будут записываться в outstrm. Реализация настраиваемого транспорта должна гарантировать, что сообщения передаются и принимаются полностью и в правильном порядке. connect(manager::ClusterManager.....) устанавливает TCP/IP-соединения между работниками.
Distributed.init_workerФункция
init_worker(cookie::AbstractString, manager::ClusterManager=DefaultClusterManager())
Вызывается менеджерами кластеров, реализующими настраиваемые транспорты. Она инициализирует недавно запущенный процесс как работник. Аргумент командной строки --worker[=<cookie>] инициализирует процесс как работника, используя TCP/IP-сокеты для транспорта. cookie является cluster_cookie.
Distributed.start_workerФункция
start_worker([out::IO=stdout], cookie::AbstractString=readline(stdin); close_stdin::Bool=true, stderr_to_stdout::Bool=true)
start_worker — внутренняя функция, являющаяся точкой входа по умолчанию для процессов работников, подключающихся через TCP/IP. Она настраивает процесс как работник кластера Julia.
Информация о хосте:порту записывается в поток out (по умолчанию 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/