Spec-Zone.ru › Julia 0.6

Задачи и параллельное вычисление

Задачи

Core.TaskТип

Task(func)

Создайте Task (то есть сопрограмму) для выполнения заданной функции (которая должна быть вызываемой без аргументов). Задача завершается, когда эта функция возвращает значение.

Пример

julia> a() = det(rand(1000, 1000));

julia> b = Task(a);

В этом примере b является выполнимой Task, которая еще не запущена.

исходный код

Base.current_taskФункция

current_task()

Получить текущую выполняющуюся Task.

исходный код

Base.istaskdoneФункция

istaskdone(t::Task) -> Bool

Определить, завершилась ли задача.

julia> a2() = det(rand(1000, 1000));

julia> b = Task(a2);

julia> istaskdone(b)
false

julia> schedule(b);

julia> yield();

julia> istaskdone(b)
true
исходный код

Base.istaskstartedФункция

istaskstarted(t::Task) -> Bool

Определить, началось ли выполнение задачи.

julia> a3() = det(rand(1000, 1000));

julia> b = Task(a3);

julia> istaskstarted(b)
false
исходный код

Base.yieldФункция

yield()

Переключиться на планировщик, чтобы разрешить выполнение другой запланированной задачи. Задача, которая вызывает эту функцию, по-прежнему выполнима и будет немедленно перезапущена, если нет других выполнимых задач.

исходный код
yield(t::Task, arg = nothing)

Быстрая, нечестная версия планирования schedule(t, arg); yield(), которая немедленно уступает t перед вызовом планировщика.

исходный код

Base.yieldtoФункция

yieldto(t::Task, arg = nothing)

Переключиться на заданную задачу. В первый раз, когда задача переключается, функция задачи вызывается без аргументов. При последующих переключениях arg возвращается из последнего вызова задачи yieldto. Это низкоуровневый вызов, который только переключает задачи, не учитывая состояния или планирование каким-либо образом. Его использование не рекомендуется.

исходный код

Base.task_local_storageМетод

task_local_storage(key)

Получить значение ключа в локальном хранилище текущей задачи.

исходный код

Base.task_local_storageМетод

task_local_storage(key, value)

Присвоить значение ключу в локальном хранилище текущей задачи.

исходный код

Base.task_local_storageМетод

task_local_storage(body, key, value)

Вызвать функцию body с измененным локальным хранилищем задачи, в котором value присваивается key; предыдущее значение key, или его отсутствие, восстанавливается после этого. Полезно для эмуляции динамического области действия.

исходный код

Base.ConditionТип

Condition()

Создайте источник событий с срабатыванием по фронту, для которого задачи могут ожидать. Задачи, которые вызывают wait на Condition, приостанавливаются и помещаются в очередь. Задачи пробуждаются, когда notify позже вызывается на Condition. Срабатывание по фронту означает, что пробуждаться могут только те задачи, которые ожидали на момент вызова notify. Для срабатывания по уровню уведомлений необходимо сохранить дополнительное состояние, чтобы отслеживать, произошло ли уведомление. Тип Channel делает это, и поэтому может использоваться для событий с триггером по уровню.

исходный код

Base.notifyФункция

notify(condition, val=nothing; all=true, error=false)

Разбудить задачи, ожидающие условия, передав им val. Если all равно true (по умолчанию), все ожидающие задачи пробуждаются, в противном случае только одна. Если error равно true, переданное значение поднимается как исключение в разбуженных задачах.

Возвращает количество разбуженных задач. Возвращает 0, если задачи не ожидают condition.

исходный код

Base.scheduleФункция

schedule(t::Task, [val]; error=false)

Добавить Task в очередь планировщика. Это заставляет задачу постоянно выполняться, когда система в противном случае простаивает, если задача не выполняет блокирующую операцию, например, wait.

Если предоставлен второй аргумент val, он будет передан задаче (через значение возврата yieldto) при ее повторном выполнении. Если error равно true, значение поднимается как исключение в разбуженной задаче.

julia> a5() = det(rand(1000, 1000));

julia> b = Task(a5);

julia> istaskstarted(b)
false

julia> schedule(b);

julia> yield();

julia> istaskstarted(b)
true

julia> istaskdone(b)
true
исходный код

Base.@scheduleМакрос

@schedule

Оборачивает выражение в Task и добавляет его в очередь планировщика локальной машины. Похоже на @async, за исключением того, что окружающий @sync НЕ ожидает задач, запущенных с @schedule.

исходный код

Base.@taskМакрос

@task

Оборачивает выражение в Task без его выполнения и возвращает Task. Это создает только задачу и не запускает ее.

julia> a1() = det(rand(1000, 1000));

julia> b = @task a1();

julia> istaskstarted(b)
false

julia> schedule(b);

julia> yield();

julia> istaskdone(b)
true
исходный код

Base.sleepФункция

sleep(seconds)

Заблокировать текущую задачу на указанное количество секунд. Минимальное время ожидания составляет 1 миллисекунду или вход 0.001.

исходный код

Base.ChannelТип

Channel{T}(sz::Int)

Создает Channel с внутренней буферизацией, которая может содержать максимальное количество sz объектов типа T. Вызовы put! на канале с полным буфером блокируются, пока объект не будет удален с помощью take!.

Channel(0) создает канал без буферизации. put! блокируется до тех пор, пока не будет вызван соответствующий take!. И наоборот.

Другие конструкторы:

  • Channel(Inf): эквивалентно Channel{Any}(typemax(Int))

  • Channel(sz): эквивалентно Channel{Any}(sz)

исходный код

Base.put!Метод

put!(c::Channel, v)

Добавляет элемент v в канал c. Блокируется, если канал заполнен.

Для каналов без буферизации блокируется до тех пор, пока другая задача не выполнит take!.

исходный код

Base.take!Метод

take!(c::Channel)

Удаляет и возвращает значение из Channel. Блокируется до тех пор, пока данные не станут доступны.

Для каналов без буферизации блокируется до тех пор, пока другая задача не выполнит put!.

исходный код

Base.isreadyМетод

isready(c::Channel)

Определить, имеет ли Channel сохраненное значение. Возвращает немедленно, не блокируется.

Для каналов без буферизации возвращает true если задачи ожидают вызова put!.

исходный код

Base.fetchМетод

fetch(c::Channel)

Ожидает и получает первый доступный элемент из канала. Элемент не удаляется. fetch не поддерживается для небуферизованного (0-размерного) канала.

исходный код

Base.closeМетод

close(c::Channel)

Закрывает канал. Исключение выбрасывается:

  • put! при обращении к закрытому каналу.

  • take! и fetch для пустого закрытого канала.

исходный код

Base.bindМетод

bind(chnl::Channel, task::Task)

Связывает жизненный цикл chnl с задачей. Канал chnl автоматически закрывается при завершении задачи. Любое неперехваченное исключение в задаче распространяется на всех ожидающих на chnl.

Объект chnl может быть явно закрыт независимо от завершения задачи. Завершение задач не влияет на уже закрытые объекты Channel.

Когда канал связан с несколькими задачами, первая завершившаяся задача закроет канал. Когда несколько каналов связаны с одной задачей, завершение задачи закроет все связанные каналы.

julia> c = Channel(0);

julia> task = @schedule foreach(i->put!(c, i), 1:4);

julia> bind(c,task);

julia> for i in c
           @show i
       end;
i = 1
i = 2
i = 3
i = 4

julia> isopen(c)
false
julia> c = Channel(0);

julia> task = @schedule (put!(c,1);error("foo"));

julia> bind(c,task);

julia> take!(c)
1

julia> put!(c,1);
ERROR: foo
Stacktrace:
 [1] check_channel_state(::Channel{Any}) at ./channels.jl:131
 [2] put!(::Channel{Any}, ::Int64) at ./channels.jl:261
исходный код

Base.asyncmapФункция

asyncmap(f, c...; ntasks=0, batch_size=nothing)

Использует несколько параллельных задач для применения f к коллекции (или нескольким коллекциям одинаковой длины). Для нескольких аргументов-коллекций f применяется поэлементно.

ntasks задает количество задач, выполняемых одновременно. В зависимости от длины коллекций, если ntasks не указано, для одновременного отображения используется до 100 задач.

ntasks также может быть задана как функция без аргументов. В этом случае количество задач, выполняемых параллельно, проверяется перед обработкой каждого элемента, и новая задача запускается, если значение ntasks_func() меньше текущего количества задач.

Если batch_size указано, коллекция обрабатывается в пакетном режиме. f должна быть функцией, которая должна принимать Vector кортежей аргументов и возвращать вектор результатов. Вектор входных данных будет иметь длину batch_size или меньше.

Следующие примеры демонстрируют выполнение в разных задачах, возвращая object_id задач, в которых выполняется функция отображения.

Сначала, если ntasks не определено, каждый элемент обрабатывается в отдельной задаче.

julia> tskoid() = object_id(current_task());

julia> asyncmap(x->tskoid(), 1:5)
5-element Array{UInt64,1}:
 0x6e15e66c75c75853
 0x440f8819a1baa682
 0x9fb3eeadd0c83985
 0xebd3e35fe90d4050
 0x29efc93edce2b961

julia> length(unique(asyncmap(x->tskoid(), 1:5)))
5

С ntasks=2 все элементы обрабатываются в 2 задачах.

julia> asyncmap(x->tskoid(), 1:5; ntasks=2)
5-element Array{UInt64,1}:
 0x027ab1680df7ae94
 0xa23d2f80cd7cf157
 0x027ab1680df7ae94
 0xa23d2f80cd7cf157
 0x027ab1680df7ae94

julia> length(unique(asyncmap(x->tskoid(), 1:5; ntasks=2)))
2

Если batch_size определено, функция отображения должна быть изменена для принятия массива кортежей аргументов и возвращения массива результатов. map используется в измененной функции отображения для достижения этой цели.

julia> batch_func(input) = map(x->string("args_tuple: ", x, ", element_val: ", x[1], ", task: ", tskoid()), input)
batch_func (generic function with 1 method)

julia> asyncmap(batch_func, 1:5; ntasks=2, batch_size=2)
5-element Array{String,1}:
 "args_tuple: (1,), element_val: 1, task: 9118321258196414413"
 "args_tuple: (2,), element_val: 2, task: 4904288162898683522"
 "args_tuple: (3,), element_val: 3, task: 9118321258196414413"
 "args_tuple: (4,), element_val: 4, task: 4904288162898683522"
 "args_tuple: (5,), element_val: 5, task: 9118321258196414413"
Примечание

В настоящее время все задачи в Julia выполняются в одном потоке ОС кооперативно. Поэтому ayncmap полезна только тогда, когда функция отображения включает ввод-вывод — диск, сеть, вызов удаленного работника и т. д.

исходный код

Base.asyncmap!Функция

asyncmap!(f, results, c...; ntasks=0, batch_size=nothing)

Подобно asyncmap(), но сохраняет вывод в results вместо возвращения коллекции.

исходный код

Общая поддержка параллельных вычислений

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

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

  • 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. По умолчанию "$JULIA_HOME/julia" или "$JULIA_HOME/julia-debug" соответственно.

  • exeflags: дополнительные флаги, передаваемые рабочим процессам.

  • topology: Указывает, как рабочие процессы подключаются друг к другу. Отправка сообщения между неподключенными рабочими процессами приводит к ошибке.

    • topology=:all_to_all: Все процессы подключены друг к другу. По умолчанию.

    • topology=:master_slave: Подключается только процесс-драйвер, т. е. pid 1, к рабочим процессам. Рабочие процессы не подключаются друг к другу.

    • topology=:custom: Метод launch менеджера кластера определяет топологию подключения через поля ident и connect_idents в WorkerConfig. Рабочий процесс с идентификатором менеджера кластера ident подключится ко всем рабочим процессам, указанным в connect_idents.

Переменные среды:

Если процесс-мастер не сможет установить соединение с только что запущенным рабочим процессом в течение 60,0 секунд, рабочий процесс считает это критической ситуацией и завершит работу. Этот таймаут можно настроить с помощью переменной среды JULIA_WORKER_TIMEOUT. Значение JULIA_WORKER_TIMEOUT в процессе-мастере задаёт количество секунд, которые только что запущенный рабочий процесс ждёт установления соединения.

исходный код
addprocs(; kwargs...) -> List of process identifiers

Эквивалентно addprocs(Sys.CPU_CORES; kwargs...)

Обратите внимание, что рабочие процессы не запускают скрипт инициализации .juliarc.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, и enable_threaded_blas имеют тот же эффект, что и документировано для addprocs(machines).

исходный код

Base.Distributed.nprocsФункция

nprocs()

Получить количество доступных процессов.

исходный код

Base.Distributed.nworkersФункция

nworkers()

Получить количество доступных рабочих процессов. Это на единицу меньше, чем nprocs(). Равно nprocs() если nprocs() == 1.

исходный код

Base.Distributed.procsМетод

procs()

Возвращает список всех идентификаторов процессов.

исходный код

Base.Distributed.procsМетод

procs(pid::Integer)

Возвращает список всех идентификаторов процессов на том же физическом узле. В частности, возвращаются все рабочие процессы, привязанные к тому же IP-адресу, что и pid.

исходный код

Base.Distributed.workersФункция

workers()

Возвращает список идентификаторов всех рабочих процессов.

исходный код

Base.Distributed.rmprocsФункция

rmprocs(pids...; waitfor=typemax(Int))

Удаляет указанные рабочие процессы. Обратите внимание, что только процесс 1 может добавлять или удалять рабочие процессы.

Аргумент waitfor определяет, как долго ждать завершения работы рабочих процессов: - Если не указано, rmprocs будет ждать, пока все запрошенные pids не будут удалены. - Возникает исключение ErrorException, если все рабочие процессы не могут быть завершены до истечения запрошенных waitfor секунд. - При значении waitfor равном 0, вызов возвращает результат немедленно, а рабочие процессы планируются на удаление в другой задаче. Возвращается объект запланированного удаления Task. Пользователь должен вызвать wait для задачи перед вызовом других параллельных вычислений.

исходный код

Base.Distributed.interruptФункция

interrupt(pids::Integer...)

Прерывает текущую выполняющуюся задачу на указанных рабочих процессах. Это эквивалентно нажатию клавиш Ctrl-C на локальном компьютере. Если аргументы не указаны, прерываются все рабочие процессы.

исходный код
interrupt(pids::AbstractVector=workers())

Прерывает текущую выполняющуюся задачу на указанных рабочих процессах. Это эквивалентно нажатию клавиш Ctrl-C на локальном компьютере. Если аргументы не указаны, прерываются все рабочие процессы.

исходный код

Base.Distributed.myidФункция

myid()

Получить идентификатор текущего процесса.

исходный код

Base.Distributed.pmapФункция

pmap([::AbstractWorkerPool], f, 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))
исходный код

Base.Distributed.RemoteExceptionТип

RemoteException(captured)

Исключения при удалённых вычислениях обрабатываются и повторно генерируются локально. RemoteException оборачивает pid рабочего процесса и перехваченное исключение. CapturedException перехватывает удалённое исключение и сериализуемую форму стека вызовов, когда исключение было сгенерировано.

исходный код

Base.Distributed.FutureТип

Future(pid::Integer=myid())

Создать Future на процессе pid. По умолчанию pid — текущий процесс.

исходный код

Base.Distributed.RemoteChannelМетод

RemoteChannel(pid::Integer=myid())

Создать ссылку на Channel{Any}(1) на процессе pid. По умолчанию pid — текущий процесс.

исходный код

Base.Distributed.RemoteChannelМетод

RemoteChannel(f::Function, pid::Integer=myid())

Создать ссылки на удалённые каналы определённого размера и типа. f() — функция, которая, когда выполняется на pid, должна возвращать реализацию AbstractChannel.

Например, RemoteChannel(()->Channel{Int}(10), pid), вернёт ссылку на канал типа Int размера 10 на pid.

По умолчанию pid — текущий процесс.

исходный код

Base.waitФункция

wait([x])

Заблокировать текущую задачу до тех пор, пока не произойдёт какое-либо событие, в зависимости от типа аргумента:

  • RemoteChannel: Ожидание значения, которое станет доступным в указанном удалённом канале.

  • Future: Ожидание значения, которое станет доступным для указанной задачи.

  • Channel: Ожидание значения, которое будет добавлено в канал.

  • Condition: Ожидание вызова notify для условия.

  • Process: Ожидание выхода процесса или цепочки процессов. Поле exitcode процесса может быть использовано для определения успеха или неудачи.

  • Task: Ожидание завершения Task и возвращение её результирующего значения. Если задача завершается с исключением, исключение распространяется (перебрасывается в задачу, которая вызвала wait).

  • RawFD: Ожидание изменений в дескрипторе файла (см. poll_fd для параметров и кода возврата).

Если аргумент не передан, задача блокируется на неопределённое время. Задачу можно перезапустить только с помощью явного вызова schedule или yieldto.

Часто wait вызывается внутри цикла while, чтобы гарантировать, что ожидаемое условие выполняется перед продолжением.

исходный код

Base.fetchМетод

fetch(x)

Ожидает и извлекает значение из x в зависимости от типа x:

  • Future: Ожидание и получение значения Future. Извлеченное значение кешируется локально. Дальнейшие вызовы fetch к той же ссылке возвращают значение из кэша. Если удалённое значение представляет собой исключение, генерирует RemoteException, который перехватывает удалённое исключение и стек вызовов.

  • RemoteChannel: Ожидание и получение значения удалённой ссылки. Исключение обрабатывается так же, как и для Future .

Извлеченный элемент не удаляется.

source

Base.Distributed.remotecallМетод

remotecall(f, id::Integer, args...; kwargs...) -> Future

Вызывает функцию f асинхронно на заданных аргументах на указанном процессе. Возвращает Future. В случае наличия ключевых аргументов, они передаются в f.

source

Base.Distributed.remotecall_waitМетод

remotecall_wait(f, id::Integer, args...; kwargs...)

Выполняет более быстрый wait(remotecall(...)) в одном сообщении на Worker, указанном идентификатором рабочего процесса id. Ключевые аргументы, если таковые имеются, передаются в f.

См. также wait и remotecall.

source

Base.Distributed.remotecall_fetchМетод

remotecall_fetch(f, id::Integer, args...; kwargs...)

Выполняет fetch(remotecall(...)) в одном сообщении. Ключевые аргументы, если таковые имеются, передаются в f. Любые исключения на удаленном узле обрабатываются в RemoteException и генерируются.

См. также fetch и remotecall.

source

Base.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.

source

Base.put!Метод

put!(rr::RemoteChannel, args...)

Сохраняет набор значений в RemoteChannel. Если канал заполнен, блокируется до освобождения места. Возвращает свой первый аргумент.

source

Base.put!Метод

put!(rr::Future, v)

Сохраняет значение в Future rr. Future — это одноразовые удалённые ссылки. put! на уже установленной Future выбросит Exception. Все асинхронные удалённые вызовы возвращают Future и устанавливают значение в возвращаемое значение вызова при завершении.

source

Base.take!Метод

take!(rr::RemoteChannel, args...)

Извлекает значение(я) из RemoteChannel rr, удаляя значение(я) в процессе.

source

Base.isreadyМетод

isready(rr::RemoteChannel, args...)

Определяет, имеет ли RemoteChannel сохранённое значение. Обратите внимание, что эта функция может вызвать гонку, так как к моменту получения её результата он может быть больше не верным. Однако она может безопасно использоваться с Future, поскольку они присваиваются только один раз.

source

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
source

Base.Distributed.WorkerPoolТип

WorkerPool(workers::Vector{Int})

Создаёт WorkerPool из вектора идентификаторов рабочих процессов.

source

Base.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 только один раз на каждый рабочий процесс.

source

Base.Distributed.default_worker_poolФункция

default_worker_pool()

WorkerPool содержащий свободные workers() — используется remote(f) и pmap (по умолчанию).

source

Base.Distributed.clear!Метод

clear!(pool::CachingPool) -> pool

Удаляет все кэшированные функции со всех участвующих рабочих процессов.

source

Base.Distributed.remoteФункция

remote([::AbstractWorkerPool], f) -> Function

Возвращает анонимную функцию, которая выполняет функцию f на доступном рабочем процессе с помощью remotecall_fetch.

source

Base.Distributed.remotecallМетод

remotecall(f, pool::AbstractWorkerPool, args...; kwargs...) -> Future

Вариант remotecall(f, pid, ....). Ожидает и получает свободный рабочий процесс из pool и выполняет remotecall на нём.

source

Base.Distributed.remotecall_waitМетод

remotecall_wait(f, pool::AbstractWorkerPool, args...; kwargs...) -> Future

Вариант remotecall_wait(f, pid, ....). Ожидает и получает свободный рабочий процесс из pool и выполняет remotecall_wait на нём.

source

Base.Distributed.remotecall_fetchМетод

remotecall_fetch(f, pool::AbstractWorkerPool, args...; kwargs...) -> result

WorkerPool вариант remotecall_fetch(f, pid, ....). Ожидает и забирает свободного работника из pool и выполняет remotecall_fetch над ним.

исходный код

Base.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 секунд.

исходный код

Base.Distributed.@spawnМакрос

@spawn

Создаёт замыкание вокруг выражения и выполняет его на автоматически выбранном процессе, возвращая Future результата.

исходный код

Base.Distributed.@spawnatМакрос

@spawnat

Принимает два аргумента, p и выражение. Создаёт замыкание вокруг выражения и выполняет его асинхронно на процессе p. Возвращает Future результата.

исходный код

Base.Distributed.@fetchМакрос

@fetch

Эквивалентно fetch(@spawn expr). См. fetch и @spawn.

исходный код

Base.Distributed.@fetchfromМакрос

@fetchfrom

Эквивалентно fetch(@spawnat p expr). См. fetch и @spawnat.

исходный код

Base.@asyncМакрос

@async

Как @schedule, @async оборачивает выражение в Task и добавляет его в очередь планировщика локальной машины. Кроме того, добавляет задачу в набор элементов, которые ожидает ближайший окружающий @sync.

исходный код

Base.@syncМакрос

@sync

Ожидает завершения всех динамически вложенных вызовов @async, @spawn, @spawnat и @parallel. Все исключения, брошенные вложенными асинхронными операциями, собираются и бросаются как CompositeException.

исходный код

Base.Distributed.@parallelМакрос

@parallel

Параллельный цикл for следующего вида:

@parallel [reducer] for var = range
    body
end

Указанный диапазон разбивается и выполняется локально на всех работниках. В случае наличия необязательной редукторной функции @parallel выполняет локальные редукции на каждом работнике с последующей окончательной редукцией на вызывающем процессе.

Обратите внимание, что без редукторной функции @parallel выполняется асинхронно, т.е. генерирует независимые задачи на всех доступных работниках и возвращается немедленно без ожидания завершения. Для ожидания завершения добавьте вызов с @sync, например:

@sync @parallel for var = range
    body
end
исходный код

Base.Distributed.@everywhereМакрос

@everywhere expr

Выполнить выражение под Main на всех участках. Эквивалентно вызову eval(Main, expr) на всех процессах. Ошибки на любом из процессов собираются в CompositeException и выбрасываются. Например:

@everywhere bar=1

определит Main.bar на всех процессах.

В отличие от @spawn и @spawnat, @everywhere не захватывает локальные переменные. Добавление префикса @everywhere с @eval позволяет нам транслировать локальные переменные с помощью интерполяции:

foo = 1
@eval @everywhere bar=$foo

Выражение оценивается под Main независимо от того, откуда @everywhere вызывается. Например:

module FooBar
    foo() = @everywhere bar()=myid()
end
FooBar.foo()

приведёт к тому, что Main.bar будет определено на всех процессах, а не FooBar.bar.

исходный код

Base.Distributed.clear!Метод

clear!(syms, pids=workers(); mod=Main)

Очищает глобальные связи в модулях, инициализируя их значением nothing. syms должно быть типа Symbol или коллекцией Symbol. pids и mod идентифицируют процессы и модуль, в которых необходимо переинициализировать глобальные переменные. Очищаются только те имена, которые найдены как определённые под mod.

Возникает исключение, если запрос на очистку глобальной константы.

исходный код

Base.Distributed.remoteref_idФункция

Base.remoteref_id(r::AbstractRemoteRef) -> RRID

Future и RemoteChannel идентифицируются полями:

  • where - указывает узел, на котором фактически находится подлежащий объекту/хранилищу объект, к которому относится ссылка.

  • whence - указывает узел, с которого была создана удалённая ссылка. Обратите внимание, что это отличается от узла, на котором фактически находится подлежащий объекту объект. Например, вызов RemoteChannel(2) с главного процесса приведёт к значению where равным 2, а whence — 1.

  • id уникален для всех ссылок, созданных с работником, указанным whence.

Вместе whence и id однозначно идентифицируют ссылку на всех работниках.

Base.remoteref_id — это низкоуровневый API, который возвращает объект Base.RRID, который оборачивает значения whence и id удалённой ссылки.

исходный код

Base.Distributed.channel_from_idФункция

Base.channel_from_id(id) -> c

Низкоуровневый API, который возвращает базовое AbstractChannel для id , возвращённого remoteref_id. Вызов допустим только на том узле, где существует базовая очередь.

исходный код

Base.Distributed.worker_id_from_socketФункция

Base.worker_id_from_socket(s) -> pid

Низкоуровневый API, который, получая соединение IO или Worker, возвращает идентификатор pid подключённого к нему работника. Это полезно при написании пользовательских методов serialize для типа, оптимизирующих вывод данных в зависимости от идентификатора процесса, который их получает.

исходный код

Base.Distributed.cluster_cookieМетод

Base.cluster_cookie() -> cookie

Возвращает куки кластера.

исходный код

Base.Distributed.cluster_cookieМетод

Base.cluster_cookie(cookie) -> cookie

Устанавливает переданный куки в качестве куки кластера, затем возвращает его.

исходный код

Удалённые массивы

Base.SharedArrayТип

SharedArray{T}(dims::NTuple; init=false, pids=Int[])
SharedArray{T,N}(...)

Создайте SharedArray типа bits T и размера dims по всем процессам, указанным в pids, — все они должны находиться на одном хосте. Если N задано вызовом SharedArray{T,N}(dims), то N должно совпадать по длине с dims.

Если pids не указано, общий массив будет отображён по всем процессам на текущем хосте, включая мастер. Однако, localindexes и indexpids будут относиться только к рабочим процессам. Это упрощает код распределения работы, позволяя использовать рабочие процессы для фактических вычислений, а мастер-процесс выполняет роль драйвера.

Если задана функция init типа initfn(S::SharedArray), она вызывается на всех участвующих рабочих процессах.

Общий массив остаётся валидным до тех пор, пока ссылка на объект SharedArray существует на узле, который создал отображение.

SharedArray{T}(filename::AbstractString, dims::NTuple, [offset=0]; mode=nothing, init=false, pids=Int[])
SharedArray{T,N}(...)

Создайте SharedArray на основе файла filename, с типом элементов T (должен быть типа bits) и размером dims, по всем процессам, указанным в pids, — все они должны находиться на одном хосте. Этот файл отображается в памяти хоста, что влечёт следующие последствия:

  • Данные массива должны быть представлены в двоичном формате (например, формат ASCII, такой как CSV, не поддерживается)

  • Любые изменения, которые вы вносите в значения массива (например, A[3] = 0), также изменят значения на диске

Если pids не указано, общий массив будет отображён по всем процессам на текущем хосте, включая мастер. Однако, localindexes и indexpids будут относиться только к рабочим процессам. Это упрощает код распределения работы, позволяя использовать рабочие процессы для фактических вычислений, а мастер-процесс выполняет роль драйвера.

mode должно быть одним из "r", "r+", "w+", или "a+", и по умолчанию равно "r+", если файл, указанный в filename, уже существует, или "w+", если нет. Если задана функция init типа initfn(S::SharedArray), она вызывается на всех участвующих рабочих процессах. Вы не можете указать функцию init, если файл не является записываемым.

offset позволяет пропустить указанное количество байтов в начале файла.

источник

Base.Distributed.procsМетод

procs(S::SharedArray)

Получить вектор процессов, отображающих общий массив.

источник

Base.sdataФункция

sdata(S::SharedArray)

Возвращает фактический объект Array, поддерживающий S.

источник

Base.indexpidsФункция

indexpids(S::SharedArray)

Возвращает индекс текущего рабочего процесса в списке рабочих процессов, отображающих SharedArray (то есть в том же списке, что возвращается procs(S) ), или 0, если SharedArray не отображён локально.

источник

Base.localindexesФункция

localindexes(S::SharedArray)

Возвращает диапазон, описывающий "стандартные" индексы, которые будет обрабатывать текущий процесс. Этот диапазон следует интерпретировать в смысле линейной индексации, то есть как поддиапазон 1:length(S). В многопроцессорных контекстах возвращает пустой диапазон в родительском процессе (или в любом процессе, для которого indexpids возвращает 0).

Стоит подчеркнуть, что localindexes существует только как удобство, и вы можете распределять работу по массиву среди рабочих процессов как угодно. Для SharedArray, все индексы должны быть одинаково быстрыми для каждого рабочего процесса.

источник

Многопоточность

Этот экспериментальный интерфейс поддерживает многопоточность Julia. Типы и функции, описанные здесь, могут (и, скорее всего, будут) меняться в будущем.

Base.Threads.threadidФункция

Threads.threadid()

Получить номер идентификатора текущей потоковой нити выполнения. Главный поток имеет идентификатор 1.

источник

Base.Threads.nthreadsФункция

Threads.nthreads()

Получить количество потоков, доступных процессу Julia. Это верхняя граница (включительно) для threadid().

источник

Base.Threads.@threadsМакрос

Threads.@threads

Макрос для распараллеливания цикла for для выполнения с несколькими потоками. Он порождает nthreads() потоков, распределяет пространство итераций между ними и выполняет итерации параллельно. В конце цикла устанавливается барьер, который ожидает завершения всех потоков, и цикл возвращает результат.

источник

Base.Threads.AtomicТип

Threads.Atomic{T}

Содержит ссылку на объект типа T, гарантируя, что к нему обращаются только атомарно, т.е. безопасным для потоков способом.

Атомарно могут использоваться только некоторые "простые" типы, а именно примитивные целочисленные и плавающие типы. Это Int8...Int128, UInt8...UInt128, и Float16...Float64.

Новые атомарные объекты могут быть созданы из неатомарных значений; если значение не указано, атомарный объект инициализируется нулём.

К атомарным объектам можно обратиться, используя обозначение []:

julia> x = Threads.Atomic{Int}(3)
Base.Threads.Atomic{Int64}(3)

julia> x[] = 1
1

julia> x[]
1

Атомарные операции используют префикс atomic_, такой как atomic_add!, atomic_xchg!, и т.д.

источник

Base.Threads.atomic_cas!Функция

Threads.atomic_cas!{T}(x::Atomic{T}, cmp::T, newval::T)

Атомарно сравнить и установить x

Атомарно сравнивает значение в x с cmp. Если они равны, записывает newval в x. В противном случае, оставляет x неизменным. Возвращает старое значение в x. Сравнивая возвращённое значение с cmp (через === ), можно узнать, было ли x изменено и теперь содержит новое значение newval.

Для получения дополнительной информации, см. инструкцию LLVM cmpxchg.

Эта функция может использоваться для реализации транзакционных семантик. До транзакции записывается значение в x. После транзакции новое значение сохраняется только в том случае, если x не было изменено в промежутке времени.

julia> x = Threads.Atomic{Int}(3)
Base.Threads.Atomic{Int64}(3)

julia> Threads.atomic_cas!(x, 4, 2);

julia> x
Base.Threads.Atomic{Int64}(3)

julia> Threads.atomic_cas!(x, 3, 2);

julia> x
Base.Threads.Atomic{Int64}(2)
источник

Base.Threads.atomic_xchg!Функция

Threads.atomic_xchg!{T}(x::Atomic{T}, newval::T)

Атомарно обменять значение в x

Атомарно обменивает значение в x со значением newval. Возвращает старое значение.

Для получения дополнительной информации, см. инструкцию LLVM atomicrmw xchg.

julia> x = Threads.Atomic{Int}(3)
Base.Threads.Atomic{Int64}(3)

julia> Threads.atomic_xchg!(x, 2)
3

julia> x[]
2
источник

Base.Threads.atomic_add!Функция

Threads.atomic_add!{T}(x::Atomic{T}, val::T)

Атомарно прибавить val к x

Выполняет x[] += val атомарно. Возвращает старое значение.

Для получения дополнительной информации, см. инструкцию LLVM atomicrmw add.

julia> x = Threads.Atomic{Int}(3)
Base.Threads.Atomic{Int64}(3)

julia> Threads.atomic_add!(x, 2)
3

julia> x[]
5
источник

Base.Threads.atomic_sub!Функция

Threads.atomic_sub!{T}(x::Atomic{T}, val::T)

Атомарно вычесть val из x

Выполняет x[] -= val атомарно. Возвращает старое значение.

Для получения дополнительной информации, см. инструкцию LLVM atomicrmw sub.

julia> x = Threads.Atomic{Int}(3)
Base.Threads.Atomic{Int64}(3)

julia> Threads.atomic_sub!(x, 2)
3

julia> x[]
1
источник

Base.Threads.atomic_and!Функция

Threads.atomic_and!{T}(x::Atomic{T}, val::T)

Атомарно выполнить побитовое И x с val

Выполняет x[] &= val атомарно. Возвращает старое значение.

Для получения дополнительной информации, см. инструкцию LLVM atomicrmw and.

julia> x = Threads.Atomic{Int}(3)
Base.Threads.Atomic{Int64}(3)

julia> Threads.atomic_and!(x, 2)
3

julia> x[]
2
источник

Base.Threads.atomic_nand!Функция

Threads.atomic_nand!{T}(x::Atomic{T}, val::T)

Атомарно выполнить побитовое И НЕ x с val

Выполняет x[] = ~(x[] & val) атомарно. Возвращает значение старое значение.

Дополнительные сведения см. в инструкции LLVM atomicrmw nand.

julia> x = Threads.Atomic{Int}(3)
Base.Threads.Atomic{Int64}(3)

julia> Threads.atomic_nand!(x, 2)
3

julia> x[]
-3
исходный код

Base.Threads.atomic_or!Функция

Threads.atomic_or!{T}(x::Atomic{T}, val::T)

Атомарно выполняет побитовое ИЛИ x с val

Выполняет x[] |= val атомарно. Возвращает значение старое значение.

Дополнительные сведения см. в инструкции LLVM atomicrmw or.

julia> x = Threads.Atomic{Int}(5)
Base.Threads.Atomic{Int64}(5)

julia> Threads.atomic_or!(x, 7)
5

julia> x[]
7
исходный код

Base.Threads.atomic_xor!Функция

Threads.atomic_xor!{T}(x::Atomic{T}, val::T)

Атомарно выполняет побитовое XOR (исключающее ИЛИ) x с val

Выполняет x[] $= val атомарно. Возвращает значение старое значение.

Дополнительные сведения см. в инструкции LLVM atomicrmw xor.

julia> x = Threads.Atomic{Int}(5)
Base.Threads.Atomic{Int64}(5)

julia> Threads.atomic_xor!(x, 7)
5

julia> x[]
2
исходный код

Base.Threads.atomic_max!Функция

Threads.atomic_max!{T}(x::Atomic{T}, val::T)

Атомарно сохраняет максимальное значение из x и val в x

Выполняет x[] = max(x[], val) атомарно. Возвращает значение старое значение.

Дополнительные сведения см. в инструкции LLVM atomicrmw max.

julia> x = Threads.Atomic{Int}(5)
Base.Threads.Atomic{Int64}(5)

julia> Threads.atomic_max!(x, 7)
5

julia> x[]
7
исходный код

Base.Threads.atomic_min!Функция

Threads.atomic_min!{T}(x::Atomic{T}, val::T)

Атомарно сохраняет минимальное значение из x и val в x

Выполняет x[] = min(x[], val) атомарно. Возвращает значение старое значение.

Дополнительные сведения см. в инструкции LLVM atomicrmw min.

julia> x = Threads.Atomic{Int}(7)
Base.Threads.Atomic{Int64}(7)

julia> Threads.atomic_min!(x, 5)
7

julia> x[]
5
исходный код

Base.Threads.atomic_fenceФункция

Threads.atomic_fence()

Вставка барьера последовательной согласованности памяти

Вставляет барьер памяти с семантикой последовательной согласованности. В некоторых алгоритмах это необходимо, т.е. когда порядка приобретения/выпуска недостаточно.

Вероятно, это очень дорогостоящая операция. Учитывая, что все атомарные операции в Julia уже имеют семантику приобретения/выпуска, явные барьеры в большинстве случаев не нужны.

Дополнительные сведения см. в инструкции LLVM fence.

исходный код

Вызов ccall с использованием пула потоков (Экспериментально)

Base.@threadcallМакрос

@threadcall

Макрос @threadcall вызывается так же, как и ccall, но выполняет работу в другом потоке. Это полезно, когда вы хотите вызвать блокирующую функцию C, не заставляя основной julia поток заблокироваться. Конкурентность ограничена размером пула потоков libuv, по умолчанию равным 4 потокам, но может быть увеличена путем установки переменной среды UV_THREADPOOL_SIZE и перезапуска процесса julia.

Обратите внимание, что вызываемая функция не должна вызывать обратные вызовы в Julia.

исходный код

Примитивы синхронизации

Base.Threads.AbstractLockТип

AbstractLock

Абстрактный супертип, описывающий типы, реализующие потокобезопасные примитивы синхронизации: lock, trylock, unlock, и islocked

исходный код

Base.lockФункция

lock(the_lock)

Приобретает блокировку, когда она становится доступной. Если блокировка уже заблокирована другой задачей/потоком, она ожидает, пока она не станет доступной.

Каждый lock должен быть сопоставлен с unlock.

исходный код

Base.unlockФункция

unlock(the_lock)

Освобождает владение блокировкой.

Если это рекурсивная блокировка, которая уже была приобретена ранее, она просто уменьшает внутренний счетчик и возвращается немедленно.

исходный код

Base.trylockФункция

trylock(the_lock) -> Success (Boolean)

Приобретает блокировку, если она доступна, возвращая true в случае успеха. Если блокировка уже заблокирована другой задачей/потоком, возвращает false.

Каждый успешный trylock должен быть сопоставлен с unlock.

исходный код

Base.islockedФункция

islocked(the_lock) -> Status (Boolean)

Проверка, удерживается ли блокировка какой-либо задачей/потоком. Не следует использовать для синхронизации (см. вместо этого trylock).

исходный код

Base.ReentrantLockТип

ReentrantLock()

Создает рекурсивную блокировку для синхронизации задач. Одна и та же задача может приобретать блокировку столько раз, сколько требуется. Каждый lock должен быть сопоставлен с unlock.

Эта блокировка НЕ потокобезопасна. См. Threads.Mutex для потокобезопасной блокировки.

исходный код

Base.Threads.MutexТип

Mutex()

Это стандартные системные мьютексы для блокировки критических разделов логики.

В Windows это объект критической секции, в pthreads это pthread_mutex_t.

См. также SpinLock для более легкой блокировки.

исходный код

Base.Threads.SpinLockТип

SpinLock()

Создаёт нерекурсивную блокировку. Рекурсивное использование приведёт к тупиковой ситуации. Каждый lock должен быть сопоставлен с unlock.

Спин-блокировки «проверить-и-проверить-и-установить» самые быстрые до примерно 30-ти конкурирующих потоков. Если у вас больше конкуренции, то, возможно, блокировка не является правильным способом синхронизации.

См. также RecursiveSpinLock для версии, допускающей рекурсию.

См. также Mutex для более эффективной версии на одном ядре или если блокировка может удерживаться в течение значительного времени.

исходный код

Base.Threads.RecursiveSpinLockТип

RecursiveSpinLock()

Создаёт рекурсивную блокировку. Один и тот же поток может приобретать блокировку столько раз, сколько требуется. Каждый lock должен быть сопоставлен с unlock.

См. также SpinLock для немного более быстрой версии.

См. также Mutex для более эффективной версии на одном ядре или если блокировка может удерживаться в течение значительного времени.

исходный код

Base.SemaphoreТип

Semaphore(sem_size)

Создаёт счётный семафор, который позволяет максимум sem_size приобретений быть в использовании в любой момент. Каждое приобретение должно быть сопоставлено с выпуском.

Этот конструкт НЕ потокобезопасен.

исходный код

Base.acquireФункция

acquire(s::Semaphore)

Ожидает, пока один из sem_size разрешений станет доступным, блокируясь, пока одно не будет приобретено.

исходный код

Base.releaseФункция

release(s::Semaphore)

Возвращает одно разрешение в пул, возможно, позволяя другой задаче приобрести его и возобновить выполнение.

исходный код

Интерфейс менеджера кластера

Этот интерфейс предоставляет механизм для запуска и управления рабочими узлами Julia в различных средах кластеров. В Base есть два типа менеджеров: LocalManager, для запуска дополнительных рабочих узлов на одном хосте, и SSHManager, для запуска на удалённых хостах через ssh. Для соединения и передачи сообщений между процессами используются сокеты TCP/IP. Менеджеры кластеров могут предоставить другой транспорт.

Base.Distributed.launchФункция

launch(manager::ClusterManager, params::Dict, launched::Array, launch_ntfy::Condition)

Реализуется менеджерами кластеров. Для каждого запущенного Julia-воркера этой функцией, она должна добавить запись WorkerConfig в launched и уведомить launch_ntfy. Функция ДОЛЖНА завершиться, как только будут запущены все воркеры, запрошенные manager. params — словарь всех ключевых аргументов, с которыми вызывалась addprocs.

исходный код

Base.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.

исходный код

Base.Distributed.init_workerФункция

init_worker(cookie::AbstractString, manager::ClusterManager=DefaultClusterManager())

Вызывается менеджерами кластеров, реализующими пользовательские транспорты. Инициализирует только что запущенный процесс как воркер. Аргумент командной строки --worker имеет эффект инициализации процесса как воркера, используя TCP/IP сокеты для транспорта. cookie — cluster_cookie.

исходный код

Base.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 сокет-соединения между воркерами.

исходный код

Base.Distributed.process_messagesФункция

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

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

См. также cluster_cookie.

исходный код

© 2009–2016 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/release-0.6/stdlib/parallel/

Spec-Zone.ru

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