Распределённые вычисления
Distributed.addprocsФункция
addprocs(manager::ClusterManager; kwargs...) -> List of process identifiers
Запускает рабочие процессы через указанный менеджер кластера.
Например, кластеры Beowulf поддерживаются с помощью пользовательского менеджера кластера, реализованного в пакете ClusterManagers.jl.
Количество секунд, которое новый запущенный рабочий процесс ожидает установления соединения с главным процессом, можно указать через переменную JULIA_WORKER_TIMEOUT в среде рабочего процесса. Актуально только при использовании TCP/IP в качестве транспорта.
addprocs(machines; tunnel=false, sshflags=``, max_parallel=10, kwargs...) -> List of process identifiers
Добавляет процессы на удалённые машины через SSH. Требуется julia в том же месте на каждом узле или доступным через общую файловую систему.
machines - это вектор спецификаций машин. Рабочие процессы запускаются для каждой спецификации.
Спецификация машины — это либо строка machine_spec , либо кортеж - (machine_spec, count).
machine_spec — это строка вида [user@]host[:port] [bind_addr[:port]]. user по умолчанию - текущий пользователь, port - стандартный порт SSH. Если указано [bind_addr[:port]] , другие рабочие процессы будут подключаться к этому рабочему процессу по указанному bind_addr и port.
count — количество рабочих процессов, которые необходимо запустить на указанном узле. Если указано :auto , будет запущено столько рабочих процессов, сколько потоков CPU на конкретном узле.
Ключевые аргументы:
tunnel: еслиtrue, то будет использоваться SSH туннелирование для подключения к рабочему процессу из главного процесса. По умолчаниюfalse.sshflags: указывает дополнительные параметры ssh, напримерsshflags=`-i /home/foo/bar.pem`max_parallel: указывает максимальное количество рабочих процессов, подключенных параллельно к узлу. По умолчанию 10.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.
Переменные среды:
Если главный процесс не сможет установить соединение с новым запущенным рабочим процессом в течение 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()
5-element Array{Int64,1}:
2
3
исходный код
Distributed.nworkersФункция
nworkers()
Получить количество доступных рабочих процессов. Это на единицу меньше, чем nprocs(). Равно nprocs() , если nprocs() == 1.
Примеры
$ julia -p 5 julia> nprocs() 6 julia> nworkers() 5исходный код
Distributed.procsМетод
procs()
Возвращает список всех идентификаторов процессов, включая pid 1 (который не включён в workers()).
Примеры
$ julia -p 5
julia> procs()
3-element Array{Int64,1}:
1
2
3
исходный код
Distributed.procsМетод
procs(pid::Integer)
Возвращает список всех идентификаторов процессов на одном физическом узле. Возвращаются все рабочие процессы, привязанные к одному IP-адресу, как у pid.
Distributed.workersФункция
workers()
Возвращает список идентификаторов всех рабочих процессов.
Примеры
$ julia -p 5
julia> workers()
2-element Array{Int64,1}:
2
3
исходный код
Distributed.rmprocsФункция
rmprocs(pids...; waitfor=typemax(Int))
Удаляет указанные рабочие процессы. Обратите внимание, что только процесс 1 может добавлять или удалять рабочие процессы.
Аргумент waitfor указывает, как долго ждать завершения работы рабочих процессов:
- Если не указано,
rmprocsбудет ждать, пока все запрошенныеpidsне будут удалены. - Возникает исключение
ErrorException, если все рабочие процессы не могут быть завершены до истечения запрошенныхwaitforсекунд. - При значении
waitforравном 0, вызов возвращает немедленно с рабочими процессами, запланированными на удаление в другой задаче. Возвращается объект запланированной задачиTask. Пользователь должен вызватьwaitна задаче перед вызовом других параллельных вызовов.
Примеры
$ julia -p 5
julia> t = rmprocs(2, 3, waitfor=0)
Task (runnable) @0x0000000107c718d0
julia> wait(t)
julia> workers()
3-element Array{Int64,1}:
4
5
6
исходный код
Distributed.interruptФункция
interrupt(pids::Integer...)
Прервать текущую выполняющуюся задачу на указанных рабочих процессах. Это эквивалентно нажатию Ctrl-C на локальной машине. Если аргументов нет, все рабочие процессы прерываются.
исходный кодinterrupt(pids::AbstractVector=workers())
Прервать текущую выполняющуюся задачу на указанных рабочих процессах. Это эквивалентно нажатию Ctrl-C на локальной машине. Если аргументов нет, все рабочие процессы прерываются.
исходный код
Distributed.myidФункция
myid()
Получить идентификатор текущего процесса.
Примеры
julia> myid() 1 julia> remotecall_fetch(() -> myid(), 4) 4исходный код
Distributed.pmapФункция
pmap(f, [::AbstractWorkerPool], c...; distributed=true, batch_size=1, on_error=nothing, retry_delays=[], retry_check=nothing) -> collection
Преобразует коллекцию c путём применения f к каждому элементу с использованием доступных рабочих процессов и задач.
Для нескольких коллекций аргументов, применяет f поэлементно.
Обратите внимание, что f должны быть доступны всем рабочим процессам; см. Доступность кода и загрузка пакетов для получения подробностей.
Если пул рабочих процессов не указан, используется весь доступный пул рабочих процессов, т.е. по умолчанию.
По умолчанию, pmap распределяет вычисления по всем указанным рабочим узлам. Для использования только локального процесса и распределения по задачам, укажите distributed=false. Это эквивалентно использованию asyncmap. Например, pmap(f, c; distributed=false) эквивалентно asyncmap(f,c; ntasks=()->nworkers())
pmap также может использовать комбинацию процессов и задач с помощью аргумента batch_size. Для размеров партий больше 1, коллекция обрабатывается в несколько партий, каждая длиной batch_size или меньше. Партия отправляется как один запрос свободному рабочему узлу, где локальный asyncmap обрабатывает элементы из партии с использованием нескольких одновременных задач.
Любая ошибка останавливает pmap от обработки оставшейся части коллекции. Чтобы переопределить это поведение, можно указать функцию обработки ошибок через аргумент on_error, которая принимает один аргумент, т.е. исключение. Функция может остановить обработку, повторно выбросив ошибку, или, чтобы продолжить, вернуть любое значение, которое затем возвращается в строке с результатами вызывающей стороне.
Рассмотрим следующие два примера. Первый возвращает объект исключения в строке, второй — 0 вместо любого исключения:
julia> pmap(x->iseven(x) ? error("foo") : x, 1:4; on_error=identity)
4-element Array{Any,1}:
1
ErrorException("foo")
3
ErrorException("foo")
julia> pmap(x->iseven(x) ? error("foo") : x, 1:4; on_error=ex->0)
4-element Array{Int64,1}:
1
0
3
0
Ошибки также могут обрабатываться путем повторного выполнения не удавшихся вычислений. Ключевые аргументы retry_delays и retry_check передаются в retry как ключевые аргументы delays и check соответственно. Если задано разделение на партии, и вся партия завершается ошибкой, все элементы в партии повторно выполняется.
Обратите внимание, что если оба on_error и retry_delays указаны, обработчик on_error вызывается перед повторным выполнением. Если on_error не вызывает (или не вызывает повторно) исключение, элемент не будет повторно выполняться.
Пример: При ошибках повторно выполнить f для элемента не более 3 раз без задержки между повторными попытками.
pmap(f, c; retry_delays = zeros(3))
Пример: Повторно выполнить f только если исключение не является типа InexactError, с экспоненциально увеличивающимися задержками до 3 раз. Возвращать NaN вместо всех InexactError случаев.
pmap(f, c; on_error = e->(isa(e, InexactError) ? NaN : rethrow(e)), retry_delays = ExponentialBackOff(n = 3))source
Distributed.RemoteExceptionТип
RemoteException(captured)
Исключения при удаленных вычислениях перехватываются и повторно выбрасываются локально. RemoteException оборачивает pid рабочего узла и перехваченное исключение. CapturedException перехватывает удаленное исключение и сериализуемую форму стека вызовов, когда было возбуждено исключение.
Distributed.FutureТип
Future(pid::Integer=myid())
Создать Future на процессе pid. По умолчанию pid — текущий процесс.
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.waitФункция
wait([x])
Блокировать текущую задачу, пока не произойдет какое-либо событие, в зависимости от типа аргумента:
-
Channel: Ожидать добавления значения в канал. -
Condition: Ожидатьnotifyна условии. -
Process: Ожидать выхода процесса или цепочки процессов. Полеexitcodeпроцесса может использоваться для определения успеха или неудачи. -
Task: Ожидать завершенияTask. Если задача завершается с исключением, исключение распространяется (повторно выбрасывается в задаче, которая вызвалаwait). -
RawFD: Ожидать изменений в дескрипторе файла (см. пакетFileWatching).
Если аргумент не передан, задача блокируется на неопределенный период. Задачу можно перезапустить только путем явного вызова schedule или yieldto.
Часто wait вызывается внутри цикла while для обеспечения выполнения ожидаемого условия перед продолжением.
wait(r::Future)
Ожидать, пока значение не станет доступным для указанного Future.
wait(r::RemoteChannel, args...)
Ожидать, пока значение не станет доступным в указанном RemoteChannel.
Base.fetchМетод
fetch(x)
Ожидает и получает значение из x в зависимости от типа x:
-
Future: Ожидать и получить значениеFuture. Полученное значение кешируется локально. Дальнейшие вызовыfetchк той же ссылке возвращают кэшированное значение. Если удаленное значение является исключением, выбрасываетRemoteException, который перехватывает удаленное исключение и трассировку стека. -
RemoteChannel: Ожидать и получить значение удаленной ссылки. Выбрасываемые исключения такие же, как и дляFuture.
Не удаляет полученный элемент.
source
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)). ...source
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 в качестве прокси:
c = Channel(1) @async put!(c, remotecall_fetch(long_computation, p)) isready(c) # will not blocksource
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))
source
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(wp, i -> sum(foo) + i, 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))
source
Distributed.clear!Метод
clear!(pool::CachingPool) -> pool
Удаляет все кэшированные функции со всех участвующих рабочих узлов.
source
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)
В этом примере задача выполнялась на узле с pid 2, вызванной с узла с pid 1.
source
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.9995177101692958source
Distributed.remotecall_fetchМетод
remotecall_fetch(f, pool::AbstractWorkerPool, args...; kwargs...) -> result
Вариант remotecall_fetch(f, pid, ....), использующий WorkerPool. Ожидает и берёт свободный рабочий узел из pool и выполняет remotecall_fetch на нём.
Примеры
$ julia -p 3 julia> wp = WorkerPool([2, 3]); julia> A = rand(3000); julia> remotecall_fetch(maximum, wp, A) 0.9995177101692958source
Distributed.remote_doМетод
remote_do(f, pool::AbstractWorkerPool, args...; kwargs...) -> nothing
Вариант remote_do(f, pid, ....), использующий WorkerPool. Ожидает и берёт свободный рабочий узел из pool и выполняет remote_do на нём.
Base.timedwaitФункция
timedwait(testcb::Function, secs::Float64; pollint::Float64=0.1)
Ожидает, пока testcb вернёт true или на протяжении secs секунд, в зависимости от того, что произойдёт раньше. testcb опрашивается каждые pollint секунд.
Distributed.@spawnМакрос
@spawn
Создаёт замыкание вокруг выражения и запускает его на автоматически выбранном процессе, возвращая Future с результатом.
Примеры
julia> addprocs(3); julia> f = @spawn myid() Future(2, 1, 5, nothing) julia> fetch(f) 2 julia> f = @spawn myid() Future(3, 1, 7, nothing) julia> fetch(f) 3source
Distributed.@spawnatМакрос
@spawnat
Создаёт замыкание вокруг выражения и асинхронно запускает его на процессе p. Возвращает Future с результатом. Принимает два аргумента, p и выражение.
Примеры
julia> addprocs(1); julia> f = @spawnat 2 myid() Future(2, 1, 3, nothing) julia> fetch(f) 2source
Distributed.@fetchМакрос
@fetch
Эквивалентно fetch(@spawn expr). См. fetch и @spawn.
Примеры
julia> addprocs(3); julia> @fetch myid() 2 julia> @fetch myid() 3 julia> @fetch myid() 4 julia> @fetch myid() 2source
Distributed.@fetchfromМакрос
@fetchfrom
Эквивалентно fetch(@spawnat p expr). См. fetch и @spawnat.
Примеры
julia> addprocs(3); julia> @fetchfrom 2 myid() 2 julia> @fetchfrom 4 myid() 4исходный код
Base.@asyncМакрос
@async
Оборачивает выражение в Task и добавляет его в очередь планировщика локальной машины.
Base.@syncМакрос
@sync
Ожидает завершения всех лексически вложенных вызовов @async, @spawn, @spawnat и @distributed. Все исключения, сгенерированные асинхронными операциями, собираются и выбрасываются как CompositeException.
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 на всех процессах.
В отличие от @spawn и @spawnat, @everywhere не захватывает локальные переменные. Вместо этого локальные переменные можно транслировать с помощью интерполяции:
foo = 1 @everywhere bar = $foo
Необязательный аргумент procs позволяет указать подмножество всех процессов для выполнения выражения.
Эквивалентно вызову remotecall_eval(Main, procs, expr).
Distributed.clear!Метод
clear!(syms, pids=workers(); mod=Main)
Очищает глобальные привязки в модулях, инициализируя их значением nothing. syms должен иметь тип Symbol или коллекцию Symbol . pids и mod идентифицируют процессы и модуль, в котором необходимо повторно инициализировать глобальные переменные. Очищаются только те имена, которые определены в mod.
Возникает исключение, если запрос на очистку глобальной константы.
исходный код
Distributed.remoteref_idФункция
remoteref_id(r::AbstractRemoteRef) -> RRID
Future и RemoteChannel идентифицируются полями:
where- указывает узел, где фактически существует подлежащий объект/хранилище, к которому ссылается ссылка.whence- указывает узел, с которого была создана удалённая ссылка. Обратите внимание, что это отличается от узла, где фактически существует подлежащий объект. Например, вызовRemoteChannel(2)с главного процесса приведет к значениюwhereравным 2 и значениюwhenceравным 1.idуникален для всех ссылок, созданных с рабочего процесса, указанногоwhence.
Вместе whence и id однозначно идентифицируют ссылку по всем рабочим процессам.
remoteref_id - это API низкого уровня, который возвращает объект RRID, который оборачивает whence и id значения удалённой ссылки.
Distributed.channel_from_idФункция
channel_from_id(id) -> c
API низкого уровня, который возвращает базовый AbstractChannel для id , возвращенного remoteref_id. Вызов действителен только на узле, где существует базовая очередь.
Distributed.worker_id_from_socketФункция
worker_id_from_socket(s) -> pid
API низкого уровня, который, получив соединение IO или Worker, возвращает идентификатор pid рабочего процесса, к которому оно подключено. Это полезно при написании пользовательских методов serialize для типа, что оптимизирует выводимые данные в зависимости от идентификатора принимающего процесса.
Distributed.cluster_cookieМетод
cluster_cookie() -> cookie
Возвращает куку системы кластеров.
исходный код
Distributed.cluster_cookieМетод
cluster_cookie(cookie) -> cookie
Устанавливает переданную куку как куки системы кластеров, затем возвращает её.
исходный кодИнтерфейс менеджера кластеров
Этот интерфейс предоставляет механизм запуска и управления рабочими процессами Julia в разных средах кластеров. В Base есть два типа менеджеров: LocalManager, для запуска дополнительных рабочих процессов на том же хосте, и SSHManager, для запуска на удалённых хостах через ssh. Для соединения и передачи сообщений между процессами используются сокеты TCP/IP. Менеджеры кластеров могут предоставить и другой транспорт.
Distributed.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). Менеджер должен сигнализировать соответствующему рабочему процессу об прерывании. - с
:finalizeдля целей очистки.
Base.killМетод
kill(manager::ClusterManager, pid::Int, config::WorkerConfig)
Реализуется менеджерами кластеров. Вызывается на главном процессе функцией rmprocs. Он должен заставить удаленный рабочий процесс, указанный pid, завершиться. kill(manager::ClusterManager.....) выполняет удалённый exit() на pid.
Sockets.connectМетод
connect(manager::ClusterManager, pid::Int, config::WorkerConfig) -> (instrm::IO, outstrm::IO)
Реализуется менеджерами кластеров с использованием пользовательских транспортов. Он должен установить логическое соединение с рабочим процессом с id pid, указанным config и вернуть пару IO объектов. Сообщения от pid к текущему процессу будут читаться из instrm, а сообщения для отправки pid будут записываться в outstrm. Реализация пользовательского транспорта должна гарантировать, что сообщения будут передаваться и приниматься полностью и в правильном порядке. connect(manager::ClusterManager.....) настраивает соединения сокетов TCP/IP между рабочими процессами.
Distributed.init_workerФункция
init_worker(cookie::AbstractString, manager::ClusterManager=DefaultClusterManager())
Вызывается менеджерами кластеров, реализующими пользовательские транспортные средства. Он инициализирует только что запущенный процесс как рабочий. Аргумент командной строки --worker[=<cookie>] имеет эффект инициализации процесса как рабочего с использованием сокетов TCP/IP для транспорта. cookie является cluster_cookie.
Distributed.start_workerФункция
start_worker([out::IO=stdout], cookie::AbstractString=readline(stdin))
start_worker — это внутренняя функция, которая является точкой входа по умолчанию для рабочих процессов, подключающихся через TCP/IP. Она настраивает процесс как рабочий процесса кластера Julia.
Информация о хосте и порте записывается в поток out (по умолчанию — стандартный вывод).
Функция закрывает stdin (после чтения куки, если требуется), перенаправляет stderr в stdout, прослушивает свободный порт (или, если указан, порт в параметре командной строки --bind-to) и планирует задачи для обработки входящих TCP-соединений и запросов.
Она не возвращает значение.
исходный код
Distributed.process_messagesФункция
process_messages(r_stream::IO, w_stream::IO, incoming::Bool=true)
Вызывается менеджерами кластеров с использованием пользовательских транспортных средств. Она должна вызываться, когда реализация пользовательского транспортного средства получает первое сообщение от удалённого рабочего процесса. Пользовательская реализация транспортного средства должна управлять логическим соединением с удалённым рабочим процессом и предоставлять два объекта IO, один для входящих сообщений, а другой — для сообщений, адресованных удалённому рабочему процессу. Если incoming равно true, удалённый узел инициировал подключение. Всякий из пары, инициирующий подключение, отправляет куки кластера и номер версии Julia для выполнения процесса аутентификации.
См. также cluster_cookie.
© 2009–2019 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.0.4/stdlib/Distributed/