Задачи
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.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на условии. -
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.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 её.
Использование 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.
Вызов '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, задача Task, созданная для 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.7.0/base/parallel/