Задачи
Core.TaskТип
Task(func)
Создайте Task (то есть сопрограмму) для выполнения заданной функции func (которая должна быть вызываема без аргументов). Задача завершается, когда эта функция возвращает значение. Задача будет выполняться в «возрасте мира» родительского элемента при создании, когда schedule.
По умолчанию задачи будут иметь установленный флаг sticky в значение true t.sticky. Это соответствует историческому значению по умолчанию для @async. Задачи со значением sticky могут выполняться только на потоке-работнике, на котором они были впервые запланированы, и при планировании делают задачу, из которой они были запланированы, sticky. Чтобы получить поведение Threads.@spawn, установите флаг sticky вручную в значение false.
Примеры
julia> a() = sum(i for i in 1:1000); julia> b = Task(a);
В этом примере b является исполняемым Task , который еще не начал работу.
Base.@taskМакрос
@task
Оборачивает выражение в Task без его выполнения и возвращает Task. Это создаёт задачу, но не запускает её.
По умолчанию задачи будут иметь установленный флаг sticky в значение true t.sticky. Это соответствует историческому значению по умолчанию для @async. Задачи со значением sticky могут выполняться только на потоке-работнике, на котором они были впервые запланированы, и при планировании делают задачу, из которой они были запланированы, sticky. Чтобы получить поведение Threads.@spawn, установите флаг sticky вручную в значение false.
Примеры
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 с помощью $, что копирует значение напрямую в построенное внутреннее закрытие. Это позволяет вставлять значение переменной, изолируя асинхронный код от изменений значения переменной в текущей задаче.
Сильно рекомендуется отдавать предпочтение Threads.@spawn перед @async всегда, даже когда параллелизм не требуется, особенно в публично распространяемых библиотеках. Это связано с тем, что использование @async отключает миграцию задачи родителя через потоки-рабочие в текущей реализации Julia. Таким образом, казалось бы, безобидное использование @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"
исходный код
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, значение повышается как исключение в разбуженной задаче.
Неправильно использовать schedule для произвольной Task , которая уже запущена. Дополнительную информацию см. в справочнике API.
По умолчанию у задач установлен флаг sticky со значением true t.sticky. Это соответствует историческому значению по умолчанию для @async. Задачи с флагом sticky могут выполняться только на том рабочем потоке, на котором они были впервые запланированы, и при планировании сделают задачу, из которой они были запланированы, sticky. Чтобы получить поведение Threads.@spawn, вручную установите флаг sticky в false.
Примеры
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 завершается с ошибкой.
Примеры
julia> Base._wait(errormonitor(Threads.@spawn error("task failed")))
Unhandled Task ERROR: task failed
Stacktrace:
[...]
исходный код
Base.@syncМакрос
@sync
Подождать, пока все лексически вложенные использования @async, @spawn, @spawnat и @distributed будут завершены. Все исключения, сгенерированные вложенными асинхронными операциями, собираются и генерируются как CompositeException.
Примеры
julia> Threads.nthreads()
4
julia> @sync begin
Threads.@spawn println("Thread-id $(Threads.threadid()), task 1")
Threads.@spawn println("Thread-id $(Threads.threadid()), task 2")
end;
Thread-id 3, task 1
Thread-id 1, task 2
исходный код
Base.waitФункция
Особое примечание для Threads.Condition:
Вызывающая сторона должна удерживать lock, которая принадлежит Threads.Condition, перед вызовом этого метода. Вызывающая задача будет заблокирована до тех пор, пока какая-то другая задача не разбудит её, обычно вызывая notify для того же объекта Threads.Condition. Блокировка будет атомарно освобождена при блокировке (даже если она была заблокирована рекурсивно), и будет повторно приобретена перед возвратом.
wait(r::Future)
Ожидание значения, доступного для указанного Future.
wait(r::RemoteChannel, args...)
Ожидание значения, доступного в указанном RemoteChannel.
wait([x])
Блокирует текущую задачу до тех пор, пока не произойдёт какое-либо событие, в зависимости от типа аргумента:
-
Channel: Ожидание добавления значения в канал. -
Condition: Ожиданиеnotifyна условии и возврат параметраval, переданного вnotify. Ожидание на условии дополнительно позволяет передатьfirst=true, что приводит к тому, что ожидающий становится первым в очереди для пробуждения поnotify, вместо обычного поведения FIFO. -
Process: Ожидание завершения процесса или цепочки процессов. Полеexitcodeпроцесса может использоваться для определения успеха или неудачи. -
Task: Ожидание завершенияTask. Если задача завершается с исключением, генерируетсяTaskFailedException(который оборачивает завершившуюся задачу). -
RawFD: Ожидание изменений в дескрипторе файла (см. пакетFileWatching).
Если аргумент не передан, задача блокируется на неопределённый срок. Задачу можно перезапустить только путём явного вызова schedule или yieldto.
Часто wait вызывается в цикле while для обеспечения выполнения ожидаемого условия до продолжения.
wait(c::Channel)
Блокировка до тех пор, пока Channel isready.
julia> c = Channel(1); julia> isready(c) false julia> task = Task(() -> wait(c)); julia> schedule(task); julia> istaskdone(task) # task is blocked because channel is not ready false julia> put!(c, 1); julia> istaskdone(task) # task is now unblocked trueисходный код
Base.fetchМетод
fetch(t::Task)
Ожидание завершения Task, затем возвращение его значения результата. Если задача завершается с исключением, генерируется TaskFailedException (который оборачивает завершившуюся задачу).
Base.fetchМетод
fetch(x::Any)
Возвращение x.
Base.timedwaitФункция
timedwait(testcb, timeout::Real; pollint::Real=0.1)
Ожидание, пока testcb() вернёт true или пройдёт timeout секунд, в зависимости от того, что произойдёт раньше. Проверка функции выполняется каждые pollint секунды. Минимальное значение для pollint составляет 0,001 секунды, то есть 1 миллисекунда.
Возвращение :ok или :timed_out.
Base.ConditionТип
Condition()
Создание источника событий с реагированием на изменения, для которого задачи могут ожидать. Задачи, которые вызывают wait на Condition, приостанавливаются и помещаются в очередь. Задачи активируются, когда notify позднее вызывается для Condition. Реагирование на изменения означает, что только те задачи, которые ожидают в момент вызова notify, могут быть активированы. Для уведомлений с реагированием на изменения требуется дополнительное состояние для отслеживания того, произошло ли уведомление. Типы Channel и Threads.Event это делают и могут использоваться для событий с реагированием на изменения.
Этот объект НЕ потокобезопасен. См. Threads.Condition для потокобезопасной версии.
Base.Threads.ConditionТип
Threads.Condition([lock])
Потокобезопасная версия Base.Condition.
Для вызова wait или notify на Threads.Condition, необходимо сначала вызвать lock на нём. Когда wait вызывается, блокировка атомарно освобождается во время блокирования и будет повторно приобретена перед wait возвратом. Поэтому типичный способ использования Threads.Condition c выглядит следующим образом:
lock(c)
try
while !thing_we_are_waiting_for
wait(c)
end
finally
unlock(c)
end
Для этой функциональности требуется по крайней мере Julia 1.2.
Base.EventТип
Event([autoreset=false])
Создайте источник событий с триггером уровня. Задачи, которые вызывают wait на Event, приостанавливаются и помещаются в очередь до тех пор, пока не будет вызвано notify на Event. После вызова notify, Event остаётся в сигнализированном состоянии, и задачи больше не будут блокироваться при ожидании его, пока не будет вызвано reset.
Если autoreset имеет значение true, не более одной задачи будет освобождена из wait для каждого вызова notify.
Это обеспечивает порядок памяти «приобретение и освобождение» при уведомлении/ожидании.
Для этой функции требуется по крайней мере Julia 1.1.
Функциональность autoreset и гарантия порядка памяти требуют по крайней мере Julia 1.8.
Base.notifyФункция
notify(condition, val=nothing; all=true, error=false)
Разбудить задачи, ожидающие условия, передав им val. Если all имеет значение true (по умолчанию), все ожидающие задачи разбуживаются, в противном случае только одна. Если error равно true, переданное значение поднимается в качестве исключения в разбуженных задачах.
Возвращает количество разбуженных задач. Возвращает 0, если на condition нет ожидающих задач.
Base.resetМетод
reset(::Event)
Сбросить Event в состояние «не установлено». Затем любые последующие вызовы wait будут блокироваться, пока снова не будет вызвано notify.
Base.SemaphoreТип
Semaphore(sem_size)
Создать счётный семафор, который позволяет не более sem_size приобретениям быть активными одновременно. Каждое приобретение должно быть согласовано с освобождением.
Это обеспечивает порядок памяти «приобретение и освобождение» при вызовах «приобрести/освободить».
исходный код
Base.acquireФункция
acquire(s::Semaphore)
Подождать, пока один из sem_size разрешений станет доступным, блокируясь до тех пор, пока одно не будет приобретено.
acquire(f, s::Semaphore)
Выполнить f после приобретения из семафора s, и release по завершении или ошибке.
Например, блок do, который гарантирует, что только 2 вызова foo будут активны одновременно:
s = Base.Semaphore(2)
@sync for _ in 1:100
Threads.@spawn begin
Base.acquire(s) do
foo()
end
end
end
Этот метод требует по крайней мере Julia 1.8.
Base.releaseФункция
release(s::Semaphore)
Возвратить одно разрешение в пул, возможно, разрешая другой задаче его приобрести и возобновить выполнение.
исходный код
Base.AbstractLockТип
AbstractLock
Абстрактный супертип, описывающий типы, которые реализуют синхронизирующие примитивы: lock, trylock, unlock и islocked.
Base.lockФункция
lock(lock)
Приобрести блокировку, когда она станет доступна. Если блокировка уже заблокирована другой задачей/потоком, подождать, пока она не станет доступной.
Каждое «приобретение» должно быть согласовано с «освобождением» unlock.
lock(f::Function, lock)
Приобрести блокировку, выполнить f с удерживаемой блокировкой и освободить блокировку, когда f вернётся. Если блокировка уже заблокирована другой задачей/потоком, подождать, пока она не станет доступной.
Когда эта функция возвращается, блокировка была освобождена, поэтому вызывающий не должен пытаться unlock её.
Использование Channel в качестве второго аргумента требует Julia 1.7 или более поздней версии.
Base.unlockФункция
unlock(lock)
Освобождает владение блокировкой lock.
Если это рекурсивная блокировка, которая была приобретена ранее, уменьшите внутренний счётчик и вернитесь сразу.
исходный код
Base.trylockФункция
trylock(lock) -> Success (Boolean)
Приобрести блокировку, если она доступна, и вернуть true в случае успеха. Если блокировка уже заблокирована другой задачей/потоком, вернуть false.
Каждое успешное «приобретение» должно быть согласовано с «освобождением» unlock.
Функция trylock в сочетании с islocked может использоваться для написания алгоритмов «проверка-и-проверка-и-установка» или экспоненциального отката, *если это поддерживается реализацией typeof(lock)* (читайте её документацию).
Base.islockedФункция
islocked(lock) -> Status (Boolean)
Проверить, удерживается ли блокировка lock какой-либо задачей/потоком. Эта функция сама по себе не должна использоваться для синхронизации. Однако, islocked в сочетании с trylock может использоваться для написания алгоритмов «проверка-и-проверка-и-установка» или экспоненциального отката, *если это поддерживается реализацией typeof(lock)* (читайте её документацию).
Дополнительная помощь
Например, экспоненциальный откат можно реализовать следующим образом, если реализация lock удовлетворяет описанным ниже свойствам.
nspins = 0
while true
while islocked(lock)
GC.safepoint()
nspins += 1
nspins > LIMIT && error("timeout")
end
trylock(lock) && break
backoff()
end
Реализация
Реализации блокировок рекомендуется определить islocked со следующими свойствами и отметить это в своей документации.
-
islocked(lock)не содержит гонок данных. - Если
islocked(lock)возвращаетfalse, непосредственный вызовtrylock(lock)должен быть успешным (возвращаетtrue), если нет вмешательства от других задач.
Base.ReentrantLockТип
ReentrantLock()
Создаёт рекурсивную блокировку для синхронизации задач Task. Одна и та же задача может приобретать блокировку столько раз, сколько потребуется (отсюда и часть «Reentrant» в имени). Каждое приобретение lock должно быть согласовано с освобождением unlock.
Вызов «lock» также запретит выполнение финализаторов в данном потоке до соответствующего «unlock». Использование стандартной схемы блокировки, проиллюстрированной ниже, должно быть естественным образом поддержано, но будьте осторожны при инвертировании порядка try/lock или полном отсутствии блока try (например, попытке вернуться с удерживаемой блокировкой):
Это обеспечивает порядок памяти «приобретение/освобождение» при вызовах «заблокировать/разблокировать».
lock(l)
try
<atomic work>
finally
unlock(l)
end
Если !islocked(lck::ReentrantLock) удерживается, trylock(lck) успешно выполняется, если нет других задач, пытающихся удержать блокировку «одновременно».
Каналы
Base.AbstractChannelТип
AbstractChannel{T}
Представление канала, передающего объекты типа T.
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, threadpool=nothing)
Создаёт новую задачу из func, привязывает её к новому каналу типа T и размера size, и планирует задачу, всё в одном вызове. Канал автоматически закрывается при завершении задачи.
func должен принимать привязанный канал в качестве единственного аргумента.
Если вам нужна ссылка на созданную задачу, передайте объект Ref{Task} через аргумент taskref.
Если spawn=true, то Task для func может быть запланирована в другом потоке параллельно, что эквивалентно созданию задачи с помощью Threads.@spawn.
Если spawn=true и аргумент threadpool не задан, он по умолчанию равен :default.
Если аргумент threadpool задан (в значении :default или :interactive ), это подразумевает, что spawn=true и новая задача запускаются в указанном пуле потоков.
Возвращает 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, но эти конструкторы устарели.
Аргумент threadpool= был добавлен в Julia 1.9.
julia> chnl = Channel{Char}(1, spawn=true) do ch
for c in "hello world"
put!(ch, c)
end
end
Channel{Char}(1) (2 items 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!.
Примеры
Буферизованный канал:
julia> c = Channel(1); julia> put!(c, 1); julia> take!(c) 1
Небуферизованный канал:
julia> c = Channel(0); julia> task = Task(() -> put!(c, 1)); julia> schedule(task); julia> take!(c) 1исходный код
Base.isreadyМетод
isready(c::Channel)
Определяет, содержит ли Channel сохранённое значение. Возвращает значение немедленно, не блокируется.
Для небуферизованных каналов возвращает true если задачи ожидают вызова put!.
Примеры
Буферизованный канал:
julia> c = Channel(1); julia> isready(c) false julia> put!(c, 1); julia> isready(c) true
Небуферизованный канал:
julia> c = Channel(); julia> isready(c) # no tasks waiting to put! false julia> task = Task(() -> put!(c, 1)); julia> schedule(task); # schedule a put! task julia> isready(c) trueисходный код
Base.fetchМетод
fetch(c::Channel)
Ожидает и возвращает (не удаляя) первый доступный элемент из Channel. Примечание: fetch не поддерживается в небуферизованном (0-размерном) Channel.
Примеры
Буферизованный канал:
julia> c = Channel(3) do ch
foreach(i -> put!(ch, i), 1:3)
end;
julia> fetch(c)
1
julia> collect(c) # item is not removed
3-element Vector{Any}:
1
2
3
исходный код
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
[...]
исходный код
Низкоуровневая синхронизация с помощью schedule и wait
Самое простое правильное использование schedule — это на Task который ещё не запущен (не запланирован). Однако, можно использовать schedule и wait в качестве очень низкоуровневого строительного блока для создания интерфейсов синхронизации. Ключевым предварительным условием вызова schedule(task) является то, что вызывающий должен «владеть» task; т. е., он должен знать, что вызов wait в данной task происходит в местах, известных коду, вызывающему schedule(task) Один из способов обеспечения такого предварительного условия — использование атомарных операций, как показано в следующем примере:
@enum OWEState begin
OWE_EMPTY
OWE_WAITING
OWE_NOTIFYING
end
mutable struct OneWayEvent
@atomic state::OWEState
task::Task
OneWayEvent() = new(OWE_EMPTY)
end
function Base.notify(ev::OneWayEvent)
state = @atomic ev.state
while state !== OWE_NOTIFYING
# Spin until we successfully update the state to OWE_NOTIFYING:
state, ok = @atomicreplace(ev.state, state => OWE_NOTIFYING)
if ok
if state == OWE_WAITING
# OWE_WAITING -> OWE_NOTIFYING transition means that the waiter task is
# already waiting or about to call `wait`. The notifier task must wake up
# the waiter task.
schedule(ev.task)
else
@assert state == OWE_EMPTY
# Since we are assuming that there is only one notifier task (for
# simplicity), we know that the other possible case here is OWE_EMPTY.
# We do not need to do anything because we know that the waiter task has
# not called `wait(ev::OneWayEvent)` yet.
end
break
end
end
return
end
function Base.wait(ev::OneWayEvent)
ev.task = current_task()
state, ok = @atomicreplace(ev.state, OWE_EMPTY => OWE_WAITING)
if ok
# OWE_EMPTY -> OWE_WAITING transition means that the notifier task is guaranteed to
# invoke OWE_WAITING -> OWE_NOTIFYING transition. The waiter task must call
# `wait()` immediately. In particular, it MUST NOT invoke any function that may
# yield to the scheduler at this point in code.
wait()
else
@assert state == OWE_NOTIFYING
# Otherwise, the `state` must have already been moved to OWE_NOTIFYING by the
# notifier task.
end
return
end
ev = OneWayEvent()
@sync begin
@async begin
wait(ev)
println("done")
end
println("notifying...")
notify(ev)
end
# output
notifying...
done
OneWayEvent позволяет одной задаче wait для другой задачи notify. Это ограниченный интерфейс связи, поскольку wait может быть использован только один раз одной задачей (обратите внимание на неатомарную операцию присваивания ev.task).
В этом примере notify(ev::OneWayEvent) разрешено вызвать schedule(ev.task) только тогда, когда *он* изменит состояние с OWE_WAITING на OWE_NOTIFYING. Это позволяет нам узнать, что задача, выполняющая wait(ev::OneWayEvent), теперь находится в ветке ok, и что не может быть других задач, которые пытаются schedule(ev.task), так как их @atomicreplace(ev.state, state => OWE_NOTIFYING) завершится неудачей.
© 2009–2024 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.10/base/parallel/