Spec-Zone.ru › Julia 1.4

Задачи

Core.TaskТип

Task(func)

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

Примеры

julia> a() = sum(i for i in 1:1000);

julia> b = Task(a);

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

исходный код

Base.@taskМакрос

@task

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

Примеры

julia> a1() = sum(i for i in 1:1000);

julia> b = @task a1();

julia> istaskstarted(b)
false

julia> schedule(b);

julia> yield();

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

Base.@asyncМакрос

@async

Оборачивает выражение в Task и добавляет его в очередь планировщика локальной машины.

Значения могут быть интерполированы в @async с помощью $, что копирует значение напрямую в построенное закрытие. Это позволяет вставлять значение переменной, изолируя асинхронный код от изменений значения переменной в текущей задаче.

Интерполяция значений с помощью $ доступна начиная с Julia 1.4.

исходный код

Base.@syncМакрос

@sync

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

исходный код

Base.asyncmapФункция

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

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

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

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

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

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

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

julia> tskoid() = objectid(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 выполняются в одном потоке ОС кооперативно. Следовательно, asyncmap полезно только тогда, когда функция отображения включает в себя ввод-вывод - диск, сеть, вызов удалённого работника и т.д.

исходный код

Base.asyncmap!Функция

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

Аналогично asyncmap, но сохраняет вывод в results вместо возврата коллекции.

исходный код

Base.fetchМетод

fetch(t::Task)

Ожидает завершения задачи, затем возвращает её значение. Если задача завершается с исключением, выбрасывается TaskFailedException (который оборачивает завершившуюся с ошибкой задачу).

исходный код

Base.current_taskФункция

current_task()

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

исходный код

Base.istaskdoneФункция

istaskdone(t::Task) -> Bool

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

Примеры

julia> a2() = sum(i for i in 1: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() = sum(i for i in 1:1000);

julia> b = Task(a3);

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

Base.istaskfailedФункция

istaskfailed(t::Task) -> Bool

Определяет, завершилась ли задача из-за выброшенного исключения.

Примеры

julia> a4() = error("task failed");

julia> b = Task(a4);

julia> istaskfailed(b)
false

julia> schedule(b);

julia> yield();

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

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.yieldФункция

yield()

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

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

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

исходный код

Base.yieldtoФункция

yieldto(t::Task, arg = nothing)

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

исходный код

Base.sleepФункция

sleep(seconds)

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

исходный код

Base.waitФункция

wait([x])

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

  • Channel: Ожидать добавления значения в канал.
  • Condition: Ожидать notify на условии.
  • Process: Ожидать завершения процесса или цепочки процессов. Поле exitcode процесса может использоваться для определения успеха или неудачи.
  • Task: Ожидать завершения Task. Если задача завершается с исключением, выбрасывается TaskFailedException (который оборачивает завершившуюся с ошибкой задачу).
  • RawFD: Ожидать изменений в дескрипторе файла (см. пакет FileWatching).

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

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

source

Особое примечание для Threads.Condition:

Вызывающая сторона должна удерживать lock, которому принадлежит c, перед вызовом этого метода. Вызывающая задача будет заблокирована до тех пор, пока какая-то другая задача не разбудит её, обычно вызывая notify` для того же объекта Condition. Блокировка будет атомарно освобождена при блокировании (даже если она была заблокирована рекурсивно), и будет повторно приобретена перед возвратом.

source
wait(r::Future)

Ожидать, пока значение не станет доступным для указанного Future.

source
wait(r::RemoteChannel, args...)

Ожидать, пока значение не станет доступным в указанном RemoteChannel.

source

Base.timedwaitФункция

timedwait(testcb::Function, secs::Float64; pollint::Float64=0.1)

Ожидает, пока testcb вернёт true или в течение secs секунд, что наступит раньше. testcb опрошается каждые pollint секунд.

Возвращает :ok, :timed_out или :error

source

Base.ConditionТип

Condition()

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

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

source

Base.Threads.ConditionТип

Threads.Condition([lock])

Потокобезопасная версия Base.Condition.

Для этой функциональности требуется как минимум Julia 1.2.

source

Base.notifyФункция

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

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

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

source

Base.scheduleФункция

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

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

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

Примеры

julia> a5() = sum(i for i in 1:1000);

julia> b = Task(a5);

julia> istaskstarted(b)
false

julia> schedule(b);

julia> yield();

julia> istaskstarted(b)
true

julia> istaskdone(b)
true
source

Base.EventТип

Event()

Создать источник событий с триггером по уровню. Задачи, которые вызывают wait на Event, приостанавливаются и добавляются в очередь до тех пор, пока не вызывается notify для Event. После вызова notify, Event остаётся в сигнальном состоянии, и задачи больше не будут блокироваться при ожидании его.

Для этой функциональности требуется как минимум Julia 1.1.

source

Base.SemaphoreТип

Semaphore(sem_size)

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

source

Base.acquireФункция

acquire(s::Semaphore)

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

source

Base.releaseФункция

release(s::Semaphore)

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

source

Base.AbstractLockТип

AbstractLock

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

source

Base.lockФункция

lock(lock)

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

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

source
lock(f::Function, lock)

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

По возвращении из этой функции lock была освобождена, поэтому вызывающий код не должен пытаться unlock её.

source

Base.unlockФункция

unlock(lock)

Освобождает владение lock.

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

source

Base.trylockФункция

trylock(lock) -> Success (Boolean)

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

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

исходный код

Base.islockedФункция

islocked(lock) -> Status (Boolean)

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

исходный код

Base.ReentrantLockТип

ReentrantLock()

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

исходный код

Base.ChannelТип

Channel{T=Any}(size::Int=0)

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

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

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

  • Channel(): конструктор по умолчанию, эквивалентный Channel{Any}(0)
  • Channel(Inf): эквивалентный Channel{Any}(typemax(Int))
  • Channel(sz): эквивалентный Channel{Any}(sz)

Конструктор по умолчанию Channel() и по умолчанию size=0 были добавлены в Julia 1.3.

исходный код

Base.ChannelМетод

Channel{T=Any}(func::Function, size=0; taskref=nothing, spawn=false)

Создаёт новую задачу из func, привязывает её к новому каналу типа T и размера size, и планирует задачу, всё в одном вызове.

func должен принимать связанный канал в качестве единственного аргумента.

Если вам нужен ссылка на созданную задачу, передайте объект Ref{Task} через ключевой аргумент taskref.

Если spawn = true, задача, созданная для func может быть запланирована в другом потоке параллельно, что эквивалентно созданию задачи с помощью Threads.@spawn.

Возвращает Channel.

Примеры

julia> chnl = Channel() do ch
           foreach(i -> put!(ch, i), 1:4)
       end;

julia> typeof(chnl)
Channel{Any}

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

Ссылка на созданную задачу:

julia> taskref = Ref{Task}();

julia> chnl = Channel(taskref=taskref) do ch
           println(take!(ch))
       end;

julia> istaskdone(taskref[])
false

julia> put!(chnl, "Hello");
Hello

julia> istaskdone(taskref[])
true

Параметр spawn= был добавлен в Julia 1.3. Этот конструктор был добавлен в Julia 1.3. В более ранних версиях Julia Channel использовал ключевые аргументы для установки size и T, но эти конструкторы устарели.

julia> chnl = Channel{Char}(1, spawn=true) do ch
           for c in "hello world"
               put!(ch, c)
           end
       end
Channel{Char}(sz_max:1,sz_curr:1)

julia> String(collect(chnl))
"hello world"
исходный код

Base.put!Метод

put!(c::Channel, v)

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

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

v теперь преобразуется в тип канала с помощью convert, как только put! будет вызван.

исходный код

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[, excp::Exception])

Закрывает канал. Исключение (указанное необязательно excp):

  • put! на закрытом канале.
  • take! и fetch на пустом, закрытом канале.
исходный код

Base.bindМетод

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

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

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

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

Примеры

julia> c = Channel(0);

julia> task = @async 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 = @async (put!(c,1);error("foo"));

julia> bind(c,task);

julia> take!(c)
1

julia> put!(c,1);
ERROR: foo
Stacktrace:
[...]
исходный код

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

Spec-Zone.ru

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