Задачи
Core.TaskТип
Task(func)
Создать Task (т. е. сопрограмму) для выполнения заданной функции func (которая должна быть вызываемой без аргументов). Задача завершается, когда эта функция возвращает значение. Задача будет выполняться в «возрасте мира» родительского элемента при создании, когда scheduled.
Примеры
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.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.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
Эта функция требует как минимум Julia 1.3.
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.scheduleФункция
schedule(t::Task, [val]; error=false)
Добавить Task в очередь планировщика. Это заставляет задачу постоянно выполняться, когда система в противном случае простаивает, если задача не выполняет блокирующую операцию, такую как wait.
Если предоставлен второй аргумент val, он будет передан в задачу (через значение возврата yieldto) при повторном выполнении. Если error равно true, значение поднимается как исключение в разбуженной задаче.
Неправильно использовать schedule на произвольной Task задаче, которая уже запущена. Дополнительную информацию см. в справочнике API.
Примеры
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исходный код
Синхронизация
Синхронизация
Base.errormonitorФункция
errormonitor(t::Task)
Выводит журнал ошибок в stderr, если задача t завершается неудачно.
Base.@syncМакрос
@sync
Ожидает, пока все лексически вложенные использования @async, @spawn, @spawnat и @distributed будут завершены. Все исключения, брошенные вложенными асинхронными операциями, собираются и выбрасываются как CompositeException.
Base.waitФункция
Особое примечание для Threads.Condition:
Вызывающий код должен удерживать lock, владеющий Threads.Condition, перед вызовом этого метода. Вызывающая задача будет заблокирована до тех пор, пока не будет разбужена другой задачей, как правило, вызовом notify для того же объекта Threads.Condition. Блокировка будет атомарно освобождена во время блокировки (даже если она была заблокирована рекурсивно), и будет повторно захвачена перед возвратом.
wait(r::Future)
Ожидание значения, которое станет доступным для указанной Future.
wait(r::RemoteChannel, args...)
Ожидание значения, которое станет доступным в указанном RemoteChannel.
wait([x])
Заблокировать текущую задачу до наступления какого-либо события, в зависимости от типа аргумента:
-
Channel: Ожидание добавления значения в канал. -
Condition: Ожиданиеnotifyна условии и возвращение параметраval, переданного вnotify. -
Process: Ожидание завершения процесса или цепочки процессов. Полеexitcodeпроцесса может использоваться для определения успеха или неудачи. -
Task: Ожидание завершенияTask. Если задача завершается с исключением, выбрасываетсяTaskFailedException(который оборачивает неудачную задачу). -
RawFD: Ожидание изменений в дескрипторе файла (см. пакетFileWatching).
Если аргумент не передан, задача блокируется на неопределенный период. Задачу можно перезапустить только с помощью явного вызова schedule или yieldto.
Часто wait вызывается внутри цикла while, чтобы убедиться, что ожидаемое условие выполняется перед продолжением.
Base.fetchМетод
fetch(t::Task)
Ожидание завершения задачи, а затем возвращение её значения результата. Если задача завершается с исключением, выбрасывается TaskFailedException (который оборачивает неудачную задачу).
Base.timedwaitФункция
timedwait(callback::Function, timeout::Real; pollint::Real=0.1)
Ожидает, пока callback вернёт true или пройдёт timeout секунд, в зависимости от того, что наступит раньше. callback опрошается каждые pollint секунды. Минимальное значение для timeout и pollint составляет 0.001, то есть 1 миллисекунда.
Возвращает :ok или :timed_out
исходный код
Base.ConditionТип
Condition()
Создает источник событий, на который могут ждать задачи. Задачи, которые вызывают wait на Condition, приостанавливаются и помещаются в очередь. Задачи разбуживаются, когда позже вызывается notify на Condition. Триггерирование по краю означает, что разбужены только те задачи, которые ожидали на момент вызова notify. Для триггерирования по уровню нужно хранить дополнительное состояние, чтобы отслеживать, произошло ли уведомление. Типы Channel и Threads.Event делают это и могут быть использованы для событий по уровню.
Этот объект НЕ потокобезопасен. См. Threads.Condition для потокобезопасной версии.
Base.Threads.ConditionТип
Threads.Condition([lock])
Потокобезопасная версия Base.Condition.
Для вызова wait или notify на Threads.Condition, необходимо сначала вызвать lock на нём. При вызове wait, блокировка атомарно освобождается во время блокировки, и будет повторно захвачена перед тем, как wait вернёт управление. Поэтому типичное использование Threads.Condition c выглядит следующим образом:
lock(c)
try
while !thing_we_are_waiting_for
wait(c)
end
finally
unlock(c)
end
Для этой функциональности требуется как минимум Julia 1.2.
Base.EventТип
Event([autoreset=false])
Создаёт источник событий с триггерированием по уровню. Задачи, которые вызывают wait на Event, приостанавливаются и помещаются в очередь до вызова notify на Event. После вызова notify, Event остаётся в состоянии сигнализации и задачи больше не будут блокироваться при ожидании его, пока не будет вызван reset. Если autoreset равно true, не более одной задачи будет освобождена из очереди wait для каждого вызова notify.
Это обеспечивает упорядоченность памяти при получении и освобождении уведомлений.
Для этой функциональности требуется как минимум Julia 1.1.
Функциональность и гарантии порядка памяти autoreset требуют как минимум Julia 1.8.
Base.notifyФункция
notify(condition, val=nothing; all=true, error=false)
Разбудить задачи, ожидающие условия, передавая им val. Если all равно true (по умолчанию), все ожидающие задачи разбуживаются, в противном случае только одна. Если error равно true, переданное значение поднимается как исключение в разбуженных задачах.
Возвращает количество разбуженных задач. Возвращает 0, если на condition нет ожидающих задач.
Base.resetМетод
reset(::Event)
Сбросить событие в состояние не установлено. Последующие вызовы wait будут блокироваться до повторного вызова notify.
Base.SemaphoreТип
Semaphore(sem_size)
Создаёт счётный семафор, который позволяет не более sem_size приобретений быть активными одновременно. Каждое приобретение должно быть согласовано с освобождением.
Это обеспечивает согласование памяти при вызовах acquire/release.
исходный код
Base.acquireФункция
acquire(s::Semaphore)
Ожидать освобождения одного из sem_size разрешений, блокируя до тех пор, пока одно не будет получено.
acquire(f, s::Semaphore)
Выполнить f после получения из семафора s, и release по завершении или ошибке.
Например, форма do-блока, которая гарантирует, что только 2 вызова foo будут активными одновременно:
s = Base.Semaphore(2)
@sync for _ in 1:100
Threads.@spawn begin
Base.acquire(s) do
foo()
end
end
end
Этот метод требует как минимум Julia 1.8.
Base.releaseФункция
release(s::Semaphore)
Возвратить одно разрешение в пул, возможно, позволяя другой задаче его получить и возобновить выполнение.
исходный код
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 её.
Использование Channel в качестве второго аргумента требует Julia 1.7 или более поздней версии.
Base.unlockФункция
unlock(lock)
Освобождает владение lock.
Если это рекурсивный замок, который был захвачен ранее, уменьшите внутренний счётчик и вернитесь сразу же.
исходный код
Base.trylockФункция
trylock(lock) -> Success (Boolean)
Захватите блокировку, если она доступна, и верните true, если успешно. Если блокировка уже захвачена другой задачей/потоком, верните false.
Каждый успешный trylock должен быть сопоставлен с unlock.
Base.islockedФункция
islocked(lock) -> Status (Boolean)
Проверьте, удерживается ли lock какой-либо задачей/потоком. Это не следует использовать для синхронизации (см. вместо этого trylock).
Base.ReentrantLockТип
ReentrantLock()
Создаёт рекурсивный замок для синхронизации Taskов. Такая же задача может захватывать замок столько раз, сколько нужно. Каждый lock должен быть сопоставлен с unlock.
Вызов «замка» также запретит запуск финализаторов в этом потоке до соответствующего «разблокирования». Использование стандартной модели блокировки, проиллюстрированной ниже, должно поддерживаться, но будьте осторожны, изменяя порядок try/lock или полностью опуская try-блок (например, пытаясь вернуться с захваченной блокировкой):
Это обеспечивает порядок приобретения/освобождения памяти при вызовах lock/unlock.
lock(l)
try
<atomic work>
finally
unlock(l)
end
исходный кодКаналы
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}(1) (2 items available)
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: TaskFailedException
Stacktrace:
[...]
nested task error: foo
[...]
исходный код
Синхронизация низкого уровня с использованием schedule и wait
Самое простое и корректное использование schedule — на Task , которое ещё не запущено (не запланировано). Однако, schedule и wait можно использовать как очень низкоуровневый строительный блок для создания интерфейсов синхронизации. Критическое предварительное условие вызова schedule(task) — вызывающая сторона должна "владеть" task; т.е., она должна знать, что вызов wait в данной task происходит в местах, известных коду, вызывающему schedule(task). Один из способов обеспечения такого предварительного условия — использование атомарных операций, как показано в следующем примере:
@enum OWEState begin
OWE_EMPTY
OWE_WAITING
OWE_NOTIFYING
end
mutable struct OneWayEvent
@atomic state::OWEState
task::Task
OneWayEvent() = new(OWE_EMPTY)
end
function Base.notify(ev::OneWayEvent)
state = @atomic ev.state
while state !== OWE_NOTIFYING
# Spin until we successfully update the state to OWE_NOTIFYING:
state, ok = @atomicreplace(ev.state, state => OWE_NOTIFYING)
if ok
if state == OWE_WAITING
# OWE_WAITING -> OWE_NOTIFYING transition means that the waiter task is
# already waiting or about to call `wait`. The notifier task must wake up
# the waiter task.
schedule(ev.task)
else
@assert state == OWE_EMPTY
# Since we are assuming that there is only one notifier task (for
# simplicity), we know that the other possible case here is OWE_EMPTY.
# We do not need to do anything because we know that the waiter task has
# not called `wait(ev::OneWayEvent)` yet.
end
break
end
end
return
end
function Base.wait(ev::OneWayEvent)
ev.task = current_task()
state, ok = @atomicreplace(ev.state, OWE_EMPTY => OWE_WAITING)
if ok
# OWE_EMPTY -> OWE_WAITING transition means that the notifier task is guaranteed to
# invoke OWE_WAITING -> OWE_NOTIFYING transition. The waiter task must call
# `wait()` immediately. In particular, it MUST NOT invoke any function that may
# yield to the scheduler at this point in code.
wait()
else
@assert state == OWE_NOTIFYING
# Otherwise, the `state` must have already been moved to OWE_NOTIFYING by the
# notifier task.
end
return
end
ev = OneWayEvent()
@sync begin
@async begin
wait(ev)
println("done")
end
println("notifying...")
notify(ev)
end
# output
notifying...
done
OneWayEvent позволяет одной задаче wait для другой задачи's notify. Это ограниченный интерфейс связи, так как wait может быть использован только один раз из одной задачи (обратите внимание на неатомарную операцию присвоения ev.task).
В этом примере, notify(ev::OneWayEvent) разрешено вызвать schedule(ev.task) только если оно изменяет состояние с OWE_WAITING на OWE_NOTIFYING. Это позволяет нам знать, что задача, выполняющая wait(ev::OneWayEvent), теперь находится в ветке ok, и что не может быть других задач, которые пытаются schedule(ev.task), так как их @atomicreplace(ev.state, state => OWE_NOTIFYING) завершатся неудачно.
© 2009–2022 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.8/base/parallel/