Задачи и параллельное вычисление
Задачи
-
Task(func) -
Создайте
Task(т. е. сопрограмму) для выполнения данной функции (которая должна быть вызываемой без аргументов). Задача завершается, когда эта функция возвращает значение.
-
yieldto(task, arg = nothing) -
Переключиться на заданную задачу. В первый раз при переключении на задачу функция задачи вызывается без аргументов. При последующих переключениях
argвозвращается из последнего вызова задачиyieldto. Это вызов низкого уровня, который только переключает задачи, не учитывая состояния или планирование каким-либо образом. Его использование не рекомендуется.
-
current_task() -
Получить текущую выполняющуюся
Task.
-
istaskdone(task) → Bool -
Определить, завершилась ли задача.
-
istaskstarted(task) → Bool -
Определить, начала ли задача выполнение.
-
consume(task, values...) -
Получить следующее значение, переданное
produceуказанной задачей. Дополнительные аргументы могут быть переданы, чтобы быть возвращенными из последнего вызоваproduceв производителе.
-
produce(value) -
Отправить заданное значение последнему вызову
consume, переключившись на задачу-потребителя. Если следующий вызовconsumeпередает какие-либо значения, они возвращаютсяproduce.
-
yield() -
Переключиться на планировщик, чтобы позволить другой запланированной задаче запуститься. Задача, которая вызывает эту функцию, все еще может быть выполнена и будет перезапущена немедленно, если нет других готовых задач.
-
task_local_storage(key) -
Получить значение ключа из локального хранилища текущей задачи.
-
task_local_storage(key, value) -
Присвоить значение ключу в локальном хранилище текущей задачи.
-
task_local_storage(body, key, value) -
Вызвать функцию
bodyс измененным локальным хранилищем задачи, в которомvalueприсвоеноkey; предыдущее значениеkey, или его отсутствие, восстанавливается после этого. Полезно для эмуляции динамического обладания.
-
Condition() -
Создать источник событий, срабатывающих при возникновении события, для которого задачи могут ожидать. Задачи, которые вызывают
waitнаCondition, приостанавливаются и помещаются в очередь. Задачи возобновляются, когдаnotifyпозже вызывается наCondition. Триггер по краю означает, что могут быть разбужены только те задачи, которые ожидают в момент вызоваnotify. Для триггеров по уровню необходимо хранить дополнительное состояние, чтобы отслеживать, произошло ли уведомление. ТипChannelделает это, и поэтому его можно использовать для событий с триггером по уровню.
-
notify(condition, val=nothing; all=true, error=false) -
Разбудить задачи, ожидающие события, передав им
val. Еслиallравноtrue(по умолчанию), все ожидающие задачи разбуживаются, в противном случае только одна. Еслиerrorравноtrue, переданное значение поднимается как исключение в разбуженных задачах.
-
schedule(t::Task, [val]; error=false) -
Добавить задачу в очередь планировщика. Это заставляет задачу постоянно выполняться, когда система в противном случае простаивает, если задача не выполняет блокирующую операцию, такую как
wait.Если предоставлен второй аргумент
val, он будет передан задаче (через возвращаемое значениеyieldto) при ее повторном запуске. Еслиerrorравноtrue, значение поднимается как исключение в разбуженной задаче.
-
@schedule() -
Оборачивает выражение в
Taskи добавляет его в очередь планировщика локальной машины.
-
@task() -
Оборачивает выражение в
Taskбез его выполнения и возвращаетTask. Это создает только задачу и не запускает ее.
-
sleep(seconds) -
Заблокировать текущую задачу на указанное количество секунд. Минимальное время ожидания - 1 миллисекунда или вход
0.001.
-
Channel{T}(sz::Int) -
Строит
Channel, который может содержать максимальное количествоszобъектов типаT. Вызовыput!на полном канале блокируются до тех пор, пока объект не будет удален с помощьюtake!.Другие конструкторы:
-
Channel()- эквивалентноChannel{Any}(32) -
Channel(sz::Int)эквивалентноChannel{Any}(sz)
-
Общая поддержка параллельного вычисления
-
addprocs(np::Integer; restrict=true, kwargs...) → List of process identifiers -
Запускает рабочие процессы с помощью встроенного
LocalManager, который запускает рабочие процессы только на локальном хосте. Это можно использовать для эффективного использования нескольких ядер.addprocs(4)добавит 4 процесса на локальной машине. Еслиrestrictравноtrue, привязка ограничена127.0.0.1.
-
addprocs(; kwargs...) → List of process identifiers -
Эквивалентно
addprocs(Sys.CPU_CORES; kwargs...)Обратите внимание, что рабочие процессы не выполняют скрипт запуска
.juliarc.jl, а также не синхронизируют свое глобальное состояние (такое как глобальные переменные, новые определения методов и загруженные модули) ни с одним из других запущенных процессов.
-
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()). -
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(manager::ClusterManager; kwargs...) → List of process identifiers -
Запускает рабочие процессы через указанного менеджера кластера.
Например, кластеры Beowulf поддерживаются с помощью настраиваемого менеджера кластера, реализованного в пакете
ClusterManagers.jl.Количество секунд, в течение которого недавно запущенный рабочий процесс ожидает установления соединения с главным процессом, можно указать с помощью переменной
JULIA_WORKER_TIMEOUTв среде рабочего процесса. Актуально только при использовании TCP/IP в качестве транспорта.
-
nprocs() -
Получить количество доступных процессов.
-
nworkers() -
Получить количество доступных рабочих процессов. Это на единицу меньше, чем
nprocs(). Равноnprocs(), еслиnprocs() == 1.
-
procs() -
Возвращает список всех идентификаторов процессов.
-
procs(pid::Integer) -
Возвращает список всех идентификаторов процессов на одном физическом узле. Возвращаются все рабочие процессы, привязанные к одному IP-адресу, как у
pid.
-
workers() -
Возвращает список всех идентификаторов рабочих процессов.
-
rmprocs(pids...; waitfor=0.0) -
Удаляет указанные рабочие процессы. Обратите внимание, что только процесс 1 может добавлять или удалять рабочие процессы — если другой рабочий процесс попытается вызвать
rmprocs, будет выброшена ошибка. Необязательный аргументwaitforопределяет, как долго первый процесс будет ожидать завершения работы рабочих процессов.
-
interrupt(pids::AbstractVector=workers()) -
Прервать текущую выполняемую задачу на указанных рабочих. Это эквивалентно нажатию Ctrl-C на локальной машине. Если аргументы не указаны, все рабочие прерываются.
-
interrupt(pids::Integer...) -
Прервать текущую выполняемую задачу на указанных рабочих. Это эквивалентно нажатию Ctrl-C на локальной машине. Если аргументы не указаны, все рабочие прерываются.
-
myid() -
Получить идентификатор текущего процесса.
-
asyncmap(f, c...) → collection -
Преобразовать коллекцию
c, применяя@async fк каждому элементу.Для нескольких коллекций аргументов, применять f поэлементно.
-
pmap([::AbstractWorkerPool, ]f, c...; distributed=true, batch_size=1, on_error=nothing, retry_n=0, retry_max_delay=DEFAULT_RETRY_MAX_DELAY, retry_on=DEFAULT_RETRY_ON) → collection -
Преобразовать коллекцию
c, применяяfк каждому элементу с использованием доступных рабочих и задач.Для нескольких коллекций аргументов, применять f поэлементно.
Обратите внимание, что
fдолжно быть доступно всем рабочим процессам; см. Доступность кода и загрузка пакетов для получения подробностей.Если пул рабочих не указан, используются все доступные рабочие, т.е. используется стандартный пул рабочих.
По умолчанию,
pmapраспределяет вычисление по всем указанным рабочим. Для использования только локального процесса и распределения по задачам, укажитеdistributed=false. Это эквивалентноasyncmap.pmapтакже может использовать комбинацию процессов и задач с помощью аргументаbatch_size. Для размеров пакетов больше 1 коллекция разбивается на несколько пакетов, которые распределяются по рабочим. Каждый такой пакет обрабатывается параллельно с помощью задач в каждом работнике. Указанныйbatch_sizeявляется верхним пределом, фактический размер пакетов может быть меньше и рассчитывается в зависимости от доступного количества рабочих и длины коллекции.Любая ошибка останавливает pmap от обработки оставшейся части коллекции. Чтобы переопределить это поведение, можно указать функцию обработки ошибок через аргумент
on_error, которая принимает один аргумент, т.е. исключение. Функция может остановить обработку, повторно выбросив ошибку, или, чтобы продолжить, вернуть любое значение, которое затем возвращается в строке с результатами вызывающей стороне.Неудавшее вычисление также может быть повторено с помощью
retry_on,retry_n,retry_max_delay, которые передаются вretryв качестве аргументовretry_on,nиmax_delayсоответственно. Если задана пакетная обработка и весь пакет завершается неудачей, все элементы в пакете повторяются.Следующее эквивалентно:
-
pmap(f, c; distributed=false)иasyncmap(f,c) -
pmap(f, c; retry_n=1)иasyncmap(retry(remote(f)),c) -
pmap(f, c; retry_n=1, on_error=e->e)иasyncmap(x->try retry(remote(f))(x) catch e; e end, c)
-
-
remotecall(f, id::Integer, args...; kwargs...) → Future -
Вызвать функцию
fасинхронно с заданными аргументами в указанном процессе. ВозвращаетFuture. Аргументы ключевых слов, если таковые имеются, передаются вf.
-
Base.process_messages(r_stream::IO, w_stream::IO, incoming::Bool=true) -
Вызывается менеджерами кластера, использующими пользовательские транспортные средства. Его следует вызывать, когда реализация пользовательского транспорта получает первое сообщение от удалённого рабочего. Пользовательский транспорт должен управлять логическим соединением с удалённым рабочим и предоставить два объекта
IO, один для входящих сообщений, а другой — для сообщений, адресованных удалённому рабочему. Еслиincomingравноtrue, удалённый узел инициировал подключение. Тот из пары, кто инициирует подключение, отправляет куки кластера и номер версии Julia для выполнения процесса аутентификации.
-
RemoteException(captured) -
Исключения при удалённых вычислениях捕获并重新抛出到本地。
RemoteExceptionоборачивает pid рабочего и пойманное исключение.CapturedExceptioncaptures the remote exception and a serializable form of the call stack when the exception was raised.
-
Future(pid::Integer=myid()) -
Создать
Futureв процессеpid. По умолчаниюpid— это текущий процесс.
-
RemoteChannel(pid::Integer=myid()) -
Создать ссылку на
Channel{Any}(1)в процессеpid. По умолчаниюpid— это текущий процесс.
-
RemoteChannel(f::Function, pid::Integer=myid()) -
Создать ссылки на удаленные каналы определенного размера и типа.
f()— это функция, которая, когда выполняется наpid, должна возвращать реализациюAbstractChannel.Например,
RemoteChannel(()->Channel{Int}(10), pid), вернет ссылку на канал типаIntи размера 10 наpid.По умолчанию
pid— это текущий процесс.
-
wait([x]) -
Заблокировать текущую задачу до наступления какого-либо события, в зависимости от типа аргумента:
-
RemoteChannel: Ожидать получения значения на указанном удаленном канале. -
Future: Ожидать получения значения для указанного будущего. -
Channel: Ожидать добавления значения в канал. -
Condition: Ожидатьnotifyна условии. -
Process: Ожидать завершения процесса или цепочки процессов. Полеexitcodeпроцесса может быть использовано для определения успеха или неудачи. -
Task: Ожидать завершенияTaskи возвращения его результата. Если задача завершается с исключением, исключение передаётся (перебрасывается в задачу, которая вызвалаwait). -
RawFD: Ожидать изменений в дескрипторе файла (см.poll_fdдля аргументов ключевых слов и возвращаемого кода)
Если аргумент не передан, задача блокируется на неопределённый период. Задачу можно перезапустить только с помощью явного вызова
scheduleилиyieldto.Часто
waitвызывается внутри циклаwhileдля обеспечения выполнения ожидаемого условия перед продолжением. -
-
fetch(x) -
Ожидает и извлекает значение из
xв зависимости от типаx. Не удаляет извлечённый элемент:-
Future: Ожидает и получает значение Future. Извлечённое значение кэшируется локально. Дальнейшие вызовыfetchк той же ссылке возвращают кэшированное значение. Если удалённое значение является исключением, выбрасываетRemoteException, которое捕获 удалённое исключение и трассировку стека. -
RemoteChannel: Ожидает и получает значение удалённой ссылки. Исключения обрабатываются так же, как и дляFuture. -
Channel: Ожидает и получает первый доступный элемент из канала.
-
-
remotecall_wait(f, id::Integer, args...; kwargs...) -
Выполнить более быстрый
wait(remotecall(...))в одном сообщении наWorker, указанном идентификатором рабочегоid. Аргументы ключевых слов, если таковые имеются, передаются вf.
-
remotecall_fetch(f, id::Integer, args...; kwargs...) -
Выполнить
fetch(remotecall(...))в одном сообщении. Аргументы ключевых слов, если таковые имеются, передаются вf. Любые удаленные исключения捕获 вRemoteExceptionи генерируются.
-
put!(rr::RemoteChannel, args...) -
Сохранить набор значений в
RemoteChannel. Если канал заполнен, блокируется до тех пор, пока не освободится место. Возвращает свой первый аргумент.
-
put!(rr::Future, v) -
Сохранить значение в
Futurerr.Future— это удалённые ссылки, которые можно установить только один раз.put!на уже установленнойFutureвыброситException. Все асинхронные удалённые вызовы возвращаютFutureи устанавливают значение в результат вызова при завершении.
-
put!(c::Channel, v) -
Добавляет элемент
vв каналc. Блокируется, если канал заполнен.
-
take!(rr::RemoteChannel, args...) -
Извлечь значение(я) с удаленного канала, удаляя значение(я) в процессе.
-
take!(c::Channel) -
Удаляет и возвращает значение с
Channel. Блокируется до тех пор, пока данные не станут доступны.
-
isready(c::Channel) -
Определяет, содержит ли
Channelсохранённое значение.isreadyдляChannelне блокируется.
-
isready(rr::RemoteChannel, args...) -
Определяет, содержит ли
RemoteChannelсохранённое значение. Обратите внимание, что эта функция может вызвать гонки, поскольку к моменту получения результата она может быть больше не верна. Однако её можно безопасно использовать сFuture, так как они присваиваются только один раз.
-
isready(rr::Future) -
Определяет, содержит ли
Futureсохранённое значение.Если аргумент
Futureпринадлежит другому узлу, этот вызов заблокируется, чтобы дождаться ответа. Рекомендуется ожидатьrrв отдельной задаче или использовать локальныйChannelв качестве прокси:c = Channel(1) @async put!(c, remotecall_fetch(long_computation, p)) isready(c) # will not block
-
close(c::Channel) -
Закрывает канал. Исключение генерируется при:
-
put!на закрытом канале. -
take!иfetchна пустом, закрытом канале.
-
-
WorkerPool(workers) -
Создать WorkerPool из вектора идентификаторов рабочих.
-
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будут переданы только один раз на каждый рабочий узел.
-
default_worker_pool() -
WorkerPool, содержащий свободные
workers()(используетсяremote(f)).
-
remote([::AbstractWorkerPool, ]f) → Function -
Возвращает лямбду, которая выполняет функцию
fна доступном работнике с использованиемremotecall_fetch.
-
remotecall(f, pool::AbstractWorkerPool, args...; kwargs...) -
Вызывает
f(args...; kwargs...)на одном из рабочих узлов вpool. ВозвращаетFuture.
-
remotecall_wait(f, pool::AbstractWorkerPool, args...; kwargs...) -
Вызывает
f(args...; kwargs...)на одном из рабочих узлов вpool. Ожидает завершения, возвращаетFuture.
-
remotecall_fetch(f, pool::AbstractWorkerPool, args...; kwargs...) -
Вызывает
f(args...; kwargs...)на одном из рабочих узлов вpool. Ожидает завершения и возвращает результат.
-
timedwait(testcb::Function, secs::Float64; pollint::Float64=0.1) -
Ожидает, пока
testcbвернётtrueили в течениеsecsсекунд, что произойдёт раньше.testcbопрошивается каждыеpollintсекунды.
-
@spawn() -
Создаёт замыкание вокруг выражения и выполняет его на автоматически выбранном процессе, возвращая
Futureрезультата.
-
@spawnat() -
Принимает два аргумента,
pи выражение. Замыкание создаётся вокруг выражения и выполняется асинхронно на процессеp. ВозвращаетFutureрезультата.
-
@fetch() -
Эквивалентно
fetch(@spawn expr).
-
@fetchfrom() -
Эквивалентно
fetch(@spawnat p expr).
-
@async() -
Подобно
@schedule,@asyncпомещает выражение вTaskи добавляет его в очередь планировщика локальной машины. Дополнительно добавляет задачу в набор элементов, которые ждёт ближайшее окружающее@sync.@asyncтакже помещает выражение в блокlet x=x, y=y, ...для создания новой области видимости с копиями всех переменных, на которые ссылается выражение.
-
@sync() -
Ожидает завершения всех динамически вложенных применений
@async,@spawn,@spawnatи@parallel. Все исключения, сгенерированные вложенными асинхронными операциями, собираются и выбрасываются какCompositeException.
-
@parallel() -
Параллельный цикл for вида :
@parallel [reducer] for var = range body endУказанный диапазон разбивается и выполняется локально на всех рабочих узлах. В случае указания необязательной редуцирующей функции,
@parallelвыполняет локальные редукции на каждом рабочем узле с последующей редукцией на вызывающем процессе.Обратите внимание, что без редуцирующей функции
@parallelвыполняется асинхронно, т.е. порождает независимые задачи на всех доступных рабочих узлах и возвращается немедленно, не дожидаясь завершения. Для ожидания завершения добавьте префикс@sync, как в примере:@sync @parallel for var = range body end
-
@everywhere() -
Выполняет выражение на всех процессах. Ошибки в любом из процессов собираются в
CompositeExceptionи выбрасываются. Например:@everywhere bar=1
определит
barв модулеMainна всех процессах.В отличие от
@spawnи@spawnat,@everywhereне захватывает какие-либо локальные переменные. Добавление префикса@everywhereс@evalпозволяет нам транслировать локальные переменные с помощью интерполяции:foo = 1 @eval @everywhere bar=$foo
-
clear!(pool::CachingPool) → pool -
Удаляет все кэшированные функции со всех участвующих рабочих узлов.
-
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.channel_from_id(id) → c -
Низкоуровневый API, который возвращает базовое
AbstractChannelдляid, возвращённогоBase.remoteref_id(). Вызов допустим только на том узле, где существует базовое соединение.
-
Base.worker_id_from_socket(s) → pid -
Низкоуровневый API, который, получив подключение
IOилиWorker, возвращаетpidрабочего узла, к которому оно подключено. Это полезно при написании пользовательских методовserializeдля типа, что оптимизирует выходные данные в зависимости от идентификатора получающего процесса.
-
Возвращает куки кластера.
-
Base.cluster_cookie(cookie) → cookie -
Устанавливает переданный куки в качестве куки кластера и возвращает его.
Общие массивы
-
Создаёт
SharedArrayтипа bitstypeTи размераdimsна указанных процессахpids, все из которых должны находиться на одном хосте.Если
pidsне указан, общий массив будет сопоставлен со всеми процессами на текущем хосте, включая мастер. Ноlocalindexesиindexpidsбудут ссылаться только на рабочие процессы. Это облегчает разработку кода распределения задач, позволяя использовать рабочие процессы для фактического вычисления, а мастер-процесс выполняет роль драйвера.Если указана функция
initтипаinitfn(S::SharedArray), она вызывается на всех участвующих рабочих узлах.
-
SharedArray(filename::AbstractString, T::Type, dims::NTuple, [offset=0]; mode=nothing, init=false, pids=Int[]) -
Создаёт
SharedArray, поддерживаемый файломfilename, с типом элементовT(должен бытьbitstype) и размеромdims, на процессах, указанныхpids— все из которых должны находиться на одном хосте. Этот файл отображается в памяти хоста, что влечёт за собой следующие последствия:- Данные массива должны быть представлены в двоичном формате (например, формат ASCII, как CSV, не поддерживается)
- Любые изменения, которые вы вносите в значения массива (например,
A[3] = 0), также изменят значения на диске
Если
pidsне указан, общий массив будет сопоставлен со всеми процессами на текущем хосте, включая мастер. Ноlocalindexesиindexpidsбудут ссылаться только на рабочие процессы. Это облегчает разработку кода распределения задач, позволяя использовать рабочие процессы для фактического вычисления, а мастер-процесс выполняет роль драйвера.modeдолжен быть одним из"r","r+","w+", или"a+", и по умолчанию равен"r+", если файл, указанныйfilename, уже существует, или"w+", если нет. Если указана функцияinitтипаinitfn(S::SharedArray), она вызывается на всех участвующих рабочих узлах. Вы не можете указать функциюinitесли файл не доступен для записи.offsetпозволяет пропустить указанное количество байт в начале файла.
-
procs(S::SharedArray) -
Получает вектор процессов, которые отобразили общий массив.
-
sdata(S::SharedArray) -
Возвращает фактический объект
ArraybackingS.
-
indexpids(S::SharedArray) -
Возвращает индекс текущего рабочего узла в векторе
pids, то есть список рабочих узлов, отображающих SharedArray.
-
localindexes(S::SharedArray) -
Возвращает диапазон, описывающий «стандартные» индексы, обрабатываемые текущим процессом. Этот диапазон должен интерпретироваться как линейный индекс, т.е. как поддиапазон
1:length(S). В многопроцессорных контекстах возвращает пустой диапазон в родительском процессе (или любом процессе, для которогоindexpidsвозвращает 0).Стоит подчеркнуть, что
localindexesсуществует только как удобство, и вы можете разбить работу над массивом между рабочими узлами любым способом. Для SharedArray все индексы должны быть одинаково быстрыми для каждого рабочего процесса.
Многопоточность
Этот экспериментальный интерфейс поддерживает возможности многопоточности Julia. Типы и функции, описанные здесь, могут (и, скорее всего, будут) изменены в будущем.
-
Threads.threadid() -
Получить номер идентификатора текущей потоковой нити выполнения. У главного потока ID
1.
-
Threads.nthreads() -
Получить количество потоков, доступных для процесса Julia. Это верхняя граница (включая) количества потоков, используемых для
threadid().
-
Threads.@threads() -
Макрос для распараллеливания цикла for для выполнения с несколькими потоками. Он создаёт
nthreads()потоков, распределяет пространство итераций между ними и выполняет итерации параллельно. В конце цикла ставится барьер, который ожидает завершения выполнения всех потоков, и цикл возвращает управление.
-
Threads.Atomic{T}() -
Хранит ссылку на объект типа
T, гарантируя, что к нему обращаются атомарно, то есть безопасным для многопоточности способом.Только определённые «простые» типы могут быть использованы атомарно, а именно целочисленные и типы с плавающей точкой. Это типы
Int8...``Int128``,UInt8...``UInt128``, иFloat16...``Float64``.Новые атомарные объекты могут быть созданы из неатомарных значений; если ни одно не указано, атомарный объект инициализируется нулём.
К атомарным объектам можно получить доступ с помощью обозначения
[]:x::Atomic{Int} x[] = 1 val = x[]Атомарные операции используют префикс
atomic_, такой какatomic_add!,atomic_xchg!, и т.д.
-
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не было изменено тем временем.
-
Threads.atomic_xchg!{T}(x::Atomic{T}, newval::T) -
Атомарное обменное присвоение значения в
xАтомарно меняет значение в
xнаnewval. Возвращает старое значение.Для получения более подробной информации см. инструкцию LLVM
atomicrmw xchg.
-
Threads.atomic_add!{T}(x::Atomic{T}, val::T) -
Атомарное прибавление
valкxВыполняет
x[] += valатомарно. Возвращает старое (!) значение.Для получения более подробной информации см. инструкцию LLVM
atomicrmw add.
-
Threads.atomic_sub!{T}(x::Atomic{T}, val::T) -
Атомарное вычитание
valизxВыполняет
x[] -= valатомарно. Возвращает старое (!) значение.Для получения более подробной информации см. инструкцию LLVM
atomicrmw sub.
-
Threads.atomic_and!{T}(x::Atomic{T}, val::T) -
Атомарное побитовое И
xсvalВыполняет
x[] &= valатомарно. Возвращает старое (!) значение.Для получения более подробной информации см. инструкцию LLVM
atomicrmw and.
-
Threads.atomic_nand!{T}(x::Atomic{T}, val::T) -
Атомарное побитовое НЕ-И
xсvalВыполняет
x[] = ~(x[] & val)атомарно. Возвращает старое (!) значение.Для получения более подробной информации см. инструкцию LLVM
atomicrmw nand.
-
Threads.atomic_or!{T}(x::Atomic{T}, val::T) -
Атомарное побитовое ИЛИ
xсvalВыполняет
x[] |= valатомарно. Возвращает старое (!) значение.Для получения более подробной информации см. инструкцию LLVM
atomicrmw or.
-
Threads.atomic_xor!{T}(x::Atomic{T}, val::T) -
Атомарное побитовое ИСКЛЮЧАЮЩЕЕ ИЛИ
xсvalВыполняет
x[] $= valатомарно. Возвращает старое (!) значение.Для получения более подробной информации см. инструкцию LLVM
atomicrmw xor.
-
Threads.atomic_max!{T}(x::Atomic{T}, val::T) -
Атомарно сохраняет максимальное значение из
xиvalвxВыполняет
x[] = max(x[], val)атомарно. Возвращает старое (!) значение.Для получения более подробной информации см. инструкцию LLVM
atomicrmw min.
-
Threads.atomic_min!{T}(x::Atomic{T}, val::T) -
Атомарно сохраняет минимальное значение из
xиvalвxВыполняет
x[] = min(x[], val)атомарно. Возвращает старое (!) значение.Для получения более подробной информации см. инструкцию LLVM
atomicrmw max.
-
Threads.atomic_fence() -
Вставить барьер последовательной согласованности памяти
Вставляет барьер памяти с семантикой последовательной согласованности. Существуют алгоритмы, где это необходимо, т.е. где порядок приобретения/освобождения недостаточен.
Вероятно, эта операция очень ресурсоёмкая. Учитывая, что все другие атомарные операции в Julia уже имеют семантику приобретения/освобождения, явные барьеры в большинстве случаев не нужны.
Для получения более подробной информации см. инструкцию LLVM
fence.
ccall с использованием пула потоков (Экспериментально)
-
@threadcall((cfunc, clib), rettype, (argtypes...), argvals...) -
Макрос
@threadcallвызывается так же, какccall, но выполняет работу в другом потоке. Это полезно, когда необходимо вызвать блокирующую функцию C, не блокируя основнойjuliaпоток. Конкурентность ограничена размером пула потоков libuv, по умолчанию равным 4 потокам, но может быть увеличена путём установки переменной окруженияUV_THREADPOOL_SIZEи перезапуска процессаjulia.Обратите внимание, что вызываемая функция не должна вызывать обратные вызовы в Julia.
Примитивы синхронизации
-
AbstractLock -
Абстрактный супертип, описывающий типы, реализующие примитивы синхронизации потоков:
lock,trylock,unlock, иislocked
-
lock(the_lock) -
Приобретает блокировку, когда она становится доступной. Если блокировка уже заблокирована другой задачей/потоком, она ожидает, пока она не станет доступной.
Каждое
lockдолжно быть согласовано сunlock.
-
unlock(the_lock) -
Освобождает владение блокировкой.
Если это рекурсивная блокировка, которая уже была приобретена ранее, она просто уменьшает внутренний счётчик и возвращает управление немедленно.
-
trylock(the_lock) → Success (Boolean) -
Приобретает блокировку, если она доступна, возвращает
trueв случае успеха. Если блокировка уже заблокирована другой задачей/потоком, возвращаетfalse.Каждое успешное
trylockдолжно быть согласовано сunlock.
-
islocked(the_lock) → Status (Boolean) -
Проверить, удерживается ли блокировка какой-либо задачей/потоком. Не следует использовать для синхронизации (используйте
trylock).
-
ReentrantLock() -
Создаёт рекурсивную блокировку для синхронизации задач. Одна и та же задача может приобретать блокировку столько раз, сколько нужно. Каждое
lockдолжно быть согласовано сunlock.Эта блокировка НЕ потокобезопасна. См.
Threads.Mutexдля потокобезопасной блокировки.
-
Mutex() -
Это стандартные системные мьютексы для блокировки критических участков кода.
В Windows это объект критической секции, в pthreads это
pthread_mutex_t.См. также SpinLock для более лёгкой блокировки.
-
SpinLock() -
Создаёт нерекурсивную блокировку. Рекурсивное использование приведёт к тупику. Каждое
lockдолжно быть согласовано сunlock.Блокировки test-and-test-and-set спин-блокировок наиболее быстрые до примерно 30 конкурирующих потоков. Если у вас больше конкуренции, то, возможно, блокировка — не лучший способ синхронизации.
См. также RecursiveSpinLock для версии, которая допускает рекурсию.
См. также Mutex для более эффективной версии на одном ядре или если блокировка может удерживаться в течение значительного времени.
-
RecursiveSpinLock() -
Создаёт рекурсивную блокировку. Один и тот же поток может приобретать блокировку сколько угодно раз. Каждое
lockдолжно быть согласовано сunlock.См. также SpinLock для немного более быстрой версии.
См. также Mutex для более эффективной версии на одном ядре или если блокировка может удерживаться в течение значительного времени.
-
Semaphore(sem_size) -
Создаёт семафор со счётчиком, который позволяет максимум
sem_sizeприобретений, которые могут использоваться в любой момент. Каждое приобретение должно быть согласовано с освобождением.Эта конструкция НЕ потокобезопасна.
-
acquire(s::Semaphore) -
Ожидать, пока один из
sem_sizeразрешений станет доступным, блокируя, пока не будет приобретён.
-
release(s::Semaphore) -
Возвращает одно разрешение в пул, что, возможно, позволит другой задаче приобрести его и продолжить выполнение.
Интерфейс менеджера кластера
Этот интерфейс предоставляет механизм запуска и управления рабочими процессами Julia в разных кластерных средах. В модуле Base присутствуют LocalManager для запуска дополнительных рабочих процессов на том же хосте и SSHManager для запуска на удалённых хостах через ssh. Для соединения и передачи сообщений между процессами используются сокеты TCP/IP. Кластерные менеджеры могут использовать другой транспорт.
-
launch(manager::ClusterManager, params::Dict, launched::Array, launch_ntfy::Condition) -
Реализуется кластерными менеджерами. Для каждого рабочего процесса Julia, запущенного этой функцией, она должна добавить запись
WorkerConfigвlaunchedи уведомитьlaunch_ntfy. Функция ДОЛЖНА завершиться, как только будут запущены все рабочие процессы, запрошенныеmanager.params— словарь всех ключевых аргументов, с которыми была вызвана функцияaddprocs.
-
manage(manager::ClusterManager, id::Integer, config::WorkerConfig. op::Symbol) -
Реализуется кластерными менеджерами. Она вызывается на главном процессе, в течение жизни рабочего процесса, с соответствующими
opзначениями:- с
:register/:deregisterпри добавлении/удалении рабочего процесса из пула рабочих процессов Julia. - с
:interruptпри вызовеinterrupt(workers).ClusterManagerдолжен отправить соответствующему рабочему процессу сигнал прерывания. - с
:finalizeдля целей очистки.
- с
-
kill(manager::ClusterManager, pid::Int, config::WorkerConfig) -
Реализуется кластерными менеджерами. Вызывается на главном процессе процессом
rmprocs. Она должна заставить удалённый рабочий процесс, указанныйpid, завершиться.Base.kill(manager::ClusterManager.....)выполняет удалённуюexit()наpid
-
init_worker(cookie::AbstractString, manager::ClusterManager=DefaultClusterManager()) -
Вызывается кластерными менеджерами, реализующими пользовательские транспортные средства. Она инициализирует только что запущенный процесс как рабочий процесс. Аргумент командной строки
--workerимеет эффект инициализации процесса как рабочего процесса с использованием сокетов TCP/IP для транспорта.cookie—cluster_cookie().
-
connect(manager::ClusterManager, pid::Int, config::WorkerConfig) → (instrm::IO, outstrm::IO) -
Реализуется кластерными менеджерами, использующими пользовательские транспортные средства. Она должна установить логическое соединение с рабочим процессом с идентификатором
pid, указаннымconfig, и вернуть паруIOобъектов. Сообщения отpidдо текущего процесса будут читаться изinstrm, в то время как сообщения, которые необходимо отправитьpid, будут записываться вoutstrm. Реализация пользовательского транспорта должна гарантировать, что сообщения доставляются и принимаются полностью и в правильном порядке.Base.connect(manager::ClusterManager.....)устанавливает соединения сокетов TCP/IP между рабочими процессами.
© 2009–2016 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/release-0.5/stdlib/parallel/