Spec-Zone.ru › Julia 1.5

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

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. Требуется, чтобы 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.

  • multiplex: если true , то для SSH-туннелирования будет использоваться 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: Только процесс драйвера, т.е. pid 1, подключается к рабочим процессам. Рабочие процессы не подключаются друг к другу.

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

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, вызвана из задачи с id 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

Вариант 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.9995177101692958
исходный код

Distributed.remote_doМетод

remote_do(f, pool::AbstractWorkerPool, args...; kwargs...) -> nothing

Вариант remote_do(f, pid, ....) с WorkerPool. Ожидает и берёт свободную задачу из 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).

исходный код

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.ClusterManagerТип

ClusterManager

Базовый тип для менеджеров кластеров, которые контролируют рабочие процессы как кластер. Менеджеры кластеров реализуют, как рабочие процессы могут добавляться, удаляться и взаимодействовать. SSHManager и LocalManager являются подтипами этого.

исходный код

Distributed.WorkerConfigТип

WorkerConfig

Тип, используемый ClusterManager для управления рабочими процессами, добавленными в кластер. Некоторые поля используются всеми менеджерами кластеров для доступа к хосту:

  • io – соединение, используемое для доступа к рабочему процессу (подтип IO или Nothing)
  • host – адрес хоста (либо AbstractString, либо 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)

Реализуется менеджерами кластера с пользовательскими транспортом. Она должна установить логическое соединение с рабочим процессом с 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); close_stdin::Bool=true, stderr_to_stdout::Bool=true)

start_worker — внутренняя функция, являющаяся точкой входа по умолчанию для рабочих процессов, подключающихся через TCP/IP. Она настраивает процесс как рабочий процесс кластера Julia.

Информация о хосте:порту записывается в поток out (по умолчанию stdout).

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

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

источник

Distributed.process_messagesФункция

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

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

См. также cluster_cookie.

источник

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

Spec-Zone.ru

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