Задачи
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
исходный код
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Функция
Особое примечание для Threads.Condition:
Вызывающая сторона должна удерживать lock, которая владеет c перед вызовом этого метода. Вызывающая задача будет заблокирована до тех пор, пока какая-то другая задача не разбудит её, обычно вызвав notify для того же объекта Condition. Блокировка будет атомарно освобождена при блокировке (даже если она была заблокирована рекурсивно), и будет повторно приобретена перед возвратом.
wait([x])
Заблокировать текущую задачу до тех пор, пока не произойдёт какое-то событие, в зависимости от типа аргумента:
-
Channel: Ожидать добавления значения в канал. -
Condition: Ожидатьnotifyдля условия. -
Process: Ожидать завершения процесса или цепочки процессов. Полеexitcodeпроцесса можно использовать для определения успеха или неудачи. -
Task: Ожидать завершенияTask. Если задача завершится с исключением, будет выброшеноTaskFailedException(которое оборачивает завершившуюся с ошибкой задачу). -
RawFD: Ожидать изменений в дескрипторе файла (см. пакетFileWatching).
Если аргумент не передан, задача блокируется на неопределённый период. Задачу можно перезапустить только с помощью явного вызова schedule или yieldto.
Часто wait вызывается внутри цикла while для обеспечения выполнения ожидаемого условия перед продолжением.
wait(r::Future)
Ожидание значения, которое станет доступным для указанного Future.
wait(r::RemoteChannel, args...)
Ожидание значения, которое станет доступным в указанном RemoteChannel.
Base.fetchМетод
fetch(t::Task)
Ожидать завершения задачи, а затем вернуть её значение результата. Если задача завершится с исключением, будет выброшено TaskFailedException (которое оборачивает завершившуюся с ошибкой задачу).
Base.timedwaitФункция
timedwait(testcb::Function, timeout::Real; pollint::Real=0.1)
Ожидает, пока testcb вернёт true или в течение timeout секунд, в зависимости от того, что наступит раньше. testcb опрашивается каждые pollint секунды. Минимальная продолжительность для timeout и pollint составляет 1 миллисекунду или 0.001.
Возвращает :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)
Приобрести lock при его появлении. Если блокировка уже захвачена другой задачей/потоком, подождать, пока она не станет доступна.
Каждое lock должно быть согласовано с unlock.
lock(f::Function, lock)
Приобрести lock, выполнить f с удерживаемой lock, и освободить lock при возвращении f. Если блокировка уже захвачена другой задачей/потоком, подождать, пока она не станет доступна.
Когда эта функция возвращает, блокировка освобождена, поэтому вызывающая сторона не должна пытаться 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()
Создаёт рекурсивный замок для синхронизации 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 каналы использовали ключевые аргументы для установки 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"
source
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:
foo
Stacktrace:
[...]
source
© 2009–2020 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.5.3/base/parallel/