Задачи
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.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, значение повышается как исключение в разбуженной задаче.
Примеры
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.@syncМакрос
@sync
Подождите, пока все лексически вложенные использования @async, @spawn, @spawnat и @distributed будут завершены. Все исключения, брошенные вложенными асинхронными операциями, собираются и выбрасываются как CompositeException.
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, который владеет Threads.Condition, перед вызовом этого метода. Вызывающая задача будет заблокирована до тех пор, пока какая-то другая задача не разбудит ее, обычно вызывая notify для того же объекта Threads.Condition. Блокировка будет атомарно освобождена при блокировке (даже если она была заблокирована рекурсивно), и будет повторно приобретена перед возвратом.
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.notifyФункция
notify(condition, val=nothing; all=true, error=false)
Разбудить задачи, ожидающие условия, передав им val. Если all равно true (по умолчанию), все ожидающие задачи разбуживаются, в противном случае только одна. Если error равно true, переданное значение поднимается как исключение в разбуженных задачах.
Возвращает количество разбуженных задач. Возвращает 0, если на condition нет ожидающих задач.
Base.SemaphoreТип
Semaphore(sem_size)
Создаёт счётный семафор, который позволяет максимум sem_size присвоений быть используемыми в любой момент. Каждый захват должен быть сопоставлен с освобождением.
Base.acquireФункция
acquire(s::Semaphore)
Подождите, пока один из sem_size разрешений станет доступным, блокируя выполнение, пока одно не будет приобретено.
Base.releaseФункция
release(s::Semaphore)
Возвращает одно разрешение в пул, что, возможно, позволит другой задаче приобрести его и возобновить выполнение.
источник
Base.AbstractLockТип
AbstractLock
Абстрактный супертип, описывающий типы, которые реализуют синхронизационные примитивы: lock, trylock, unlock и islocked.
Base.lockФункция
lock(lock)
Приобрести блокировку, когда она станет доступной. Если блокировка уже заблокирована другой задачей/потоком, подождите, пока она не станет доступной.
Каждый захват должен быть сопоставлен с unlock.
lock(f::Function, lock)
Приобрести блокировку, выполнить f с удерживаемой блокировкой и освободить блокировку, когда f вернётся. Если блокировка уже заблокирована другой задачей/потоком, подождите, пока она не станет доступной.
После возврата из этой функции блокировка была освобождена, поэтому вызывающий не должен пытаться её захватить.
источник
Base.unlockФункция
unlock(lock)
Освобождает владение блокировкой.
Если это рекурсивная блокировка, которая была приобретена ранее, декрементирует внутренний счётчик и возвращает немедленно.
источник
Base.trylockФункция
trylock(lock) -> Success (Boolean)
Приобрести блокировку, если она доступна, и вернуть true в случае успеха. Если блокировка уже заблокирована другой задачей/потоком, вернуть false.
Каждый успешный захват должен быть сопоставлен с unlock.
Base.islockedФункция
islocked(lock) -> Status (Boolean)
Проверить, удерживается ли lock какой-либо задачей/потоком. Это не следует использовать для синхронизации (см. вместо этого trylock).
Base.ReentrantLockТип
ReentrantLock()
Создаёт повторно входящую блокировку для синхронизации Task. Один и тот же процесс может приобретать блокировку столько раз, сколько необходимо. Каждая lock должна быть согласована с unlock.
Вызов 'lock' также запретит выполнение финализаторов в этом потоке до соответствующего 'unlock'. Использование стандартной схемы блокировки, показанной ниже, должно поддерживаться естественным образом, но будьте осторожны, изменяя порядок try/lock или опускай блок try полностью (например, пытаясь вернуться с блокировкой, которая всё ещё удерживается):
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) (1 item 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
[...]
исходный код
© 2009–2021 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.6.0/base/parallel/