Распределённые вычисления
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()
Получить количество доступных процессов.
исходный код
Distributed.nworkersФункция
nworkers()
Получить количество доступных рабочих процессов. Это на единицу меньше, чем nprocs(). Равно nprocs() , если nprocs() == 1.
Distributed.procsМетод
procs()
Возвращает список всех идентификаторов процессов.
исходный код
Distributed.procsМетод
procs(pid::Integer)
Возвращает список всех идентификаторов процессов на том же физическом узле. В частности, возвращаются все рабочие процессы, привязанные к тому же IP-адресу, что и pid.
Distributed.workersФункция
workers()
Возвращает список всех идентификаторов рабочих процессов.
исходный код
Distributed.rmprocsФункция
rmprocs(pids...; waitfor=typemax(Int))
Удаляет указанные рабочие процессы. Обратите внимание, что только процесс 1 может добавлять или удалять рабочие процессы.
Аргумент waitfor определяет, как долго ждать завершения работы рабочих процессов: - Если не указан, rmprocs будет ждать, пока все запрошенные pids не будут удалены. - Возникает исключение ErrorException , если все рабочие процессы не могут быть завершены до истечения запрошенных waitfor секунд. - При значении waitfor 0 вызов возвращает немедленно с рабочими процессами, запланированными на удаление в другом задании. Возвращается запланированный объект Task. Пользователь должен вызвать wait на задании перед вызовом любых других параллельных вызовов.
Distributed.interruptФункция
interrupt(pids::Integer...)
Прерывает текущее выполняемое задание на указанных рабочих процессах. Это эквивалентно нажатию Ctrl-C на локальной машине. Если аргументы не указаны, все рабочие процессы прерываются.
исходный кодinterrupt(pids::AbstractVector=workers())
Прерывает текущее выполняемое задание на указанных рабочих процессах. Это эквивалентно нажатию Ctrl-C на локальной машине. Если аргументы не указаны, все рабочие процессы прерываются.
исходный код
Distributed.myidФункция
myid()
Получить идентификатор текущего процесса.
исходный код
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))исходный код
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)
Дождаться появления значения для указанного будущего.
исходный кодwait(r::RemoteChannel, args...)
Дождаться появления значения в указанном удалённом канале.
исходный код
Base.fetchМетод
fetch(x)
Ожидает и извлекает значение из x в зависимости от типа x.
Future: Ожидание и получение значения Future. Извлечённое значение кэшируется локально. Дальнейшие вызовы fetch к той же ссылке возвращают кэшированное значение. Если удалённое значение является исключением, выбрасывается RemoteException, который перехватывает удалённое исключение и трассировку стека.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.
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 являются одноразовыми удалёнными ссылками. Попытка установить значение в уже заданный 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 blockисходный код
Distributed.WorkerPoolТип
WorkerPool(workers::Vector{Int})
Создать WorkerPool из вектора идентификаторов рабочих узлов.
исходный код
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 (по умолчанию).
Distributed.clear!Метод
clear!(pool::CachingPool) -> pool
Удаляет все кешированные функции со всех участвующих рабочих узлов.
исходный код
Distributed.remoteФункция
remote([::AbstractWorkerPool], f) -> Function
Возвращает анонимную функцию, которая выполняет функцию f на доступном рабочем узле с использованием remotecall_fetch.
Distributed.remotecallМетод
remotecall(f, pool::AbstractWorkerPool, args...; kwargs...) -> Future
WorkerPool вариант remotecall(f, pid, ....). Ожидает и берёт свободный рабочий узел из pool и выполняет remotecall на нём.
Distributed.remotecall_waitМетод
remotecall_wait(f, pool::AbstractWorkerPool, args...; kwargs...) -> Future
WorkerPool вариант remotecall_wait(f, pid, ....). Ожидает и берёт свободный рабочий узел из pool и выполняет remotecall_wait на нём.
Distributed.remotecall_fetchМетод
remotecall_fetch(f, pool::AbstractWorkerPool, args...; kwargs...) -> result
WorkerPool вариант remotecall_fetch(f, pid, ....). Ожидает и берёт свободный рабочий узел из pool и выполняет remotecall_fetch на нём.
Distributed.remote_doМетод
remote_do(f, pool::AbstractWorkerPool, args...; kwargs...) -> nothing
WorkerPool вариант remote_do(f, pid, ....). Ожидает и берёт свободный рабочий узел из 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) 3исходный код
Distributed.@spawnatМакрос
@spawnat
Создать замыкание вокруг выражения и асинхронно запустить его на процессе p. Вернуть Future с результатом. Принимает два аргумента, p и выражение.
Примеры
julia> addprocs(1); julia> f = @spawnat 2 myid() Future(2, 1, 3, nothing) julia> fetch(f) 2исходный код
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() 2исходный код
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).ClusterManagerдолжен послать соответствующему рабочему процессу сигнал прерывания. - со значением
:finalizeдля целей очистки.
Base.killМетод
kill(manager::ClusterManager, pid::Int, config::WorkerConfig)
Реализуется менеджерами кластера. Вызывается на главном процессе функцией rmprocs. Он должен заставить удалённый рабочий процесс, указанный в pid, выйти. kill(manager::ClusterManager.....) выполняет удалённую exit() на pid.
Sockets.connectМетод
connect(manager::ClusterManager, pid::Int, config::WorkerConfig) -> (instrm::IO, outstrm::IO)
Реализуется менеджерами кластера с использованием пользовательских транспортов. Он должен установить логическое соединение с рабочим процессом с ID pid, указанным в config, и вернуть пару объектов IO. Сообщения от pid к текущему процессу будут читаться из instrm, в то время как сообщения, которые необходимо отправить pid, будут записываться в outstrm. Реализация пользовательского транспорта должна гарантировать, что сообщения передаются и принимаются полностью и в правильном порядке. connect(manager::ClusterManager.....) устанавливает соединения TCP/IP между рабочими процессами.
Distributed.init_workerФункция
init_worker(cookie::AbstractString, manager::ClusterManager=DefaultClusterManager())
Вызывается менеджерами кластера, реализующими пользовательские транспорты. Инициализирует недавно запущенный процесс как рабочий процесс. Флаг командной строки --worker[=<cookie>] имеет эффект инициализации процесса как рабочего процесса с использованием сокетов TCP/IP для транспорта. cookie является cluster_cookie.
Distributed.start_workerФункция
start_worker([out::IO=stdout], cookie::AbstractString=readline(stdin))
start_worker — это внутренняя функция, которая является точкой входа по умолчанию для рабочих процессов, подключающихся через TCP/IP. Она настраивает процесс как рабочий процесс кластера Julia.
Информация host:port записывается в поток out (по умолчанию stdout).
Функция закрывает 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/v0.7.0/stdlib/Distributed/