Задачи и параллельное вычисление
Задачи
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)
Закрывает канал. Исключение выбрасывается:
исходный код
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: Подключается только процесс-драйвер, т. е.pid1, к рабочим процессам. Рабочие процессы не подключаются друг к другу.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: Ожидание значения, которое будет добавлено в канал.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.
Base.Distributed.remotecall_waitМетод
remotecall_wait(f, id::Integer, args...; kwargs...)
Выполняет более быстрый wait(remotecall(...)) в одном сообщении на Worker, указанном идентификатором рабочего процесса id. Ключевые аргументы, если таковые имеются, передаются в f.
См. также wait и remotecall.
Base.Distributed.remotecall_fetchМетод
remotecall_fetch(f, id::Integer, args...; kwargs...)
Выполняет fetch(remotecall(...)) в одном сообщении. Ключевые аргументы, если таковые имеются, передаются в f. Любые исключения на удаленном узле обрабатываются в RemoteException и генерируются.
См. также fetch и remotecall.
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.
Base.put!Метод
put!(rr::RemoteChannel, args...)
Сохраняет набор значений в RemoteChannel. Если канал заполнен, блокируется до освобождения места. Возвращает свой первый аргумент.
Base.put!Метод
put!(rr::Future, v)
Сохраняет значение в Future rr. Future — это одноразовые удалённые ссылки. put! на уже установленной Future выбросит Exception. Все асинхронные удалённые вызовы возвращают Future и устанавливают значение в возвращаемое значение вызова при завершении.
Base.take!Метод
take!(rr::RemoteChannel, args...)
Извлекает значение(я) из RemoteChannel rr, удаляя значение(я) в процессе.
Base.isreadyМетод
isready(rr::RemoteChannel, args...)
Определяет, имеет ли RemoteChannel сохранённое значение. Обратите внимание, что эта функция может вызвать гонку, так как к моменту получения её результата он может быть больше не верным. Однако она может безопасно использоваться с Future, поскольку они присваиваются только один раз.
Base.isreadyМетод
isready(rr::Future)
Определяет, имеет ли Future сохранённое значение.
Если аргумент Future принадлежит другому узлу, этот вызов будет блокироваться, ожидая ответа. Рекомендуется ожидать rr в отдельной задаче вместо этого или использовать локальный Channel в качестве прокси:
c = Channel(1) @async put!(c, remotecall_fetch(long_computation, p)) isready(c) # will not blocksource
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 только один раз на каждый рабочий процесс.
Base.Distributed.default_worker_poolФункция
default_worker_pool()
WorkerPool содержащий свободные workers() — используется remote(f) и pmap (по умолчанию).
Base.Distributed.clear!Метод
clear!(pool::CachingPool) -> pool
Удаляет все кэшированные функции со всех участвующих рабочих процессов.
source
Base.Distributed.remoteФункция
remote([::AbstractWorkerPool], f) -> Function
Возвращает анонимную функцию, которая выполняет функцию f на доступном рабочем процессе с помощью remotecall_fetch.
Base.Distributed.remotecallМетод
remotecall(f, pool::AbstractWorkerPool, args...; kwargs...) -> Future
Вариант remotecall(f, pid, ....). Ожидает и получает свободный рабочий процесс из pool и выполняет remotecall на нём.
Base.Distributed.remotecall_waitМетод
remotecall_wait(f, pool::AbstractWorkerPool, args...; kwargs...) -> Future
Вариант remotecall_wait(f, pid, ....). Ожидает и получает свободный рабочий процесс из pool и выполняет remotecall_wait на нём.
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/