Задачи
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 и добавляет его в очередь планировщика локальной машины.
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(r::Future)
Ожидать доступности значения для указанного Future.
wait(r::RemoteChannel, args...)
Ожидать доступности значения на указанном RemoteChannel.
wait([x])
Заблокировать текущую задачу до возникновения события, в зависимости от типа аргумента:
-
Channel: Ожидание добавления значения в канал. -
Condition: Ожидание событияnotifyна условии. -
Process: Ожидание выхода процесса или цепочки процессов. Полеexitcodeпроцесса может использоваться для определения успеха или неудачи. -
Task: Ожидание завершенияTask. Если задача завершается с исключением, выбрасываетсяTaskFailedException(который обертывает завершившуюся задачу с ошибкой). -
RawFD: Ожидание изменений в дескрипторе файла (см. пакетFileWatching).
Если аргумент не передан, задача блокируется на неопределенное время. Задачу можно перезапустить только путем явного вызова schedule или yieldto.
Часто wait вызывается внутри цикла while для обеспечения выполнения ожидаемого условия перед продолжением.
Особое примечание для Threads.Condition:
Вызывающий код должен удерживать блокировку lock, которая владеет c перед вызовом этого метода. Вызывающая задача будет заблокирована до тех пор, пока какая-то другая задача не разбудит её, обычно вызвав notify` на том же объекте Condition. Блокировка будет атомарно освобождена при блокировке (даже если она была заблокирована рекурсивно), и будет повторно приобретена перед возвратом.
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)
Приобрести блокировку, когда она станет доступной. Если блокировка уже заблокирована другой задачей/потоком, подождать, пока она не станет доступной.
Каждое приобретение блокировки должно быть согласовано с unlock.
Base.unlockФункция
unlock(lock)
Освободить право владения lock.
Если это рекурсивная блокировка, которая уже была приобретена ранее, уменьшить внутренний счётчик и возвратиться немедленно.
source
Base.trylockФункция
trylock(lock) -> Success (Boolean)
Приобрести блокировку, если она доступна, и вернуть true в случае успеха. Если блокировка уже заблокирована другой задачей/потоком, вернуть false.
Каждое успешное приобретение блокировки должно быть согласовано с 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.3.1/base/parallel/