Задачи
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)
Ожидает завершения задачи, затем возвращает её значение результата. Если задача завершается с исключением, исключение передаётся (перевыбрасывается в задачу, которая вызвала fetch).
исходный код
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.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Функция
Особое примечание для Threads.Condition:
Вызывающий должен удерживать lock объекта, которому принадлежит c перед вызовом этого метода. Вызывающая задача заблокируется до тех пор, пока другая задача не разбудит её, обычно вызвав notify на том же объекте Condition. Блокировка будет атомарно освобождена во время блокировки (даже если она была заблокирована рекурсивно), и будет повторно получена перед возвратом.
wait([x])
Заблокировать текущую задачу до тех пор, пока не произойдёт событие, в зависимости от типа аргумента:
-
Channel: Ожидание добавления значения в канал. -
Condition: Ожиданиеnotifyна условии. -
Process: Ожидание завершения процесса или цепочки процессов. Полеexitcodeпроцесса можно использовать для определения успеха или неудачи. -
Task: Ожидание завершенияTaskзадачи. Если задача завершается с исключением, исключение передаётся (перевыбрасывается в задачу, которая вызвалаwait). -
RawFD: Ожидание изменений в дескрипторе файла (см. пакетFileWatching).
Если аргумент не передан, задача блокируется на неопределённый период. Задачу можно перезапустить только с помощью явного вызова schedule или yieldto.
Часто wait вызывается внутри цикла while, чтобы убедиться, что ожидаемое условие выполняется, прежде чем продолжить.
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 секунды.
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) trueисходный код
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)
Возвращает одно разрешение в пул, возможно, позволяя другой задаче получить его и возобновить выполнение.
исходный код
Base.AbstractLockТип
AbstractLock
Абстрактный супертип, описывающий типы, которые реализуют примитивы синхронизации: lock, trylock, unlock и islocked.
Base.lockФункция
lock(lock)
Получить lock, когда он станет доступным. Если блокировка уже заблокирована другой задачей/потоком, подождите, пока она не освободится.
Каждое lock должно быть согласовано с unlock.
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()
Создаёт рекурсивную блокировку для синхронизации Tasks. Одна и та же задача может получать блокировку любое количество раз. Каждое lock должно быть согласовано с Channel.
Base.ChannelТип
Channel{T}(sz::Int)
Создаёт Channel с внутренним буфером, который может содержать максимальное количество sz объектов типа T. Вызовы put! на полном канале блокируются до тех пор, пока объект не будет удалён с помощью take!.
Channel(0) создаёт небуферизованный канал. put! блокируется до тех пор, пока не вызовется соответствующее take!. И наоборот.
Другие конструкторы:
-
Channel(Inf): эквивалентноChannel{Any}(typemax(Int)) -
Channel(sz): эквивалентноChannel{Any}(sz)
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–2019 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.2.0/base/parallel/