Задачи
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, чтобы убедиться, что ожидаемое условие выполнено перед продолжением.
Особое примечание для Threads.Condition:
Вызывающая сторона должна удерживать lock, которому принадлежит c, перед вызовом этого метода. Вызывающая задача будет заблокирована до тех пор, пока какая-то другая задача не разбудит её, обычно вызывая notify` для того же объекта Condition. Блокировка будет атомарно освобождена при блокировании (даже если она была заблокирована рекурсивно), и будет повторно приобретена перед возвратом.
wait(r::Future)
Ожидать, пока значение не станет доступным для указанного Future.
wait(r::RemoteChannel, args...)
Ожидать, пока значение не станет доступным в указанном RemoteChannel.
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 для потокобезопасной версии.
Base.Threads.ConditionТип
Threads.Condition([lock])
Потокобезопасная версия Base.Condition.
Для этой функциональности требуется как минимум Julia 1.2.
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() = 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) truesource
Base.EventТип
Event()
Создать источник событий с триггером по уровню. Задачи, которые вызывают wait на Event, приостанавливаются и добавляются в очередь до тех пор, пока не вызывается notify для Event. После вызова notify, Event остаётся в сигнальном состоянии, и задачи больше не будут блокироваться при ожидании его.
Для этой функциональности требуется как минимум Julia 1.1.
Base.SemaphoreТип
Semaphore(sem_size)
Создать счётный семафор, который позволяет максимум sem_size приобретений быть в использовании в любой момент. Каждое приобретение должно соответствовать освобождению.
Base.acquireФункция
acquire(s::Semaphore)
Ожидать, пока один из sem_size разрешений станет доступным, блокируясь до тех пор, пока одно не будет приобретено.
Base.releaseФункция
release(s::Semaphore)
Возвратить одно разрешение в пул, возможно, позволив другой задаче его приобрести и возобновить выполнение.
source
Base.AbstractLockТип
AbstractLock
Абстрактный супертип, описывающий типы, которые реализуют синхронизирующие примитивы: lock, trylock, unlock и islocked.
Base.lockФункция
lock(lock)
Приобрести lock при его доступности. Если блокировка уже заблокирована другой задачей/потоком, подождите, пока она не станет доступной.
Каждое lock должно быть сопоставлено с unlock.
lock(f::Function, lock)
Приобрести lock, выполнить f с удерживаемой lock, и освободить lock при возвращении f. Если блокировка уже заблокирована другой задачей/потоком, подождите, пока она не станет доступной.
По возвращении из этой функции lock была освобождена, поэтому вызывающий код не должен пытаться unlock её.
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):
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/