Spec-Zone.ru › Julia 0.5

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

Задачи

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

END_OF_DOCUMENT_MARKER
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 рабочего и пойманное исключение. CapturedException captures 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)

Сохранить значение в Future rr. 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

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

Base.cluster_cookie(cookie) → cookie

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

Общие массивы

SharedArray(T::Type, dims::NTuple; init=false, pids=Int[])

Создаёт SharedArray типа bitstype T и размера 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)

Возвращает фактический объект Array backing S.

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/

Spec-Zone.ru

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