Spec-Zone.ru › Julia 1.9

Задачи

Core.TaskТип

Task(func)

Создайте Task (т.е. сопрограмму) для выполнения заданной функции func (которая должна быть вызываемой без аргументов). Задача завершается, когда эта функция возвращает значение. Задача будет выполняться в «возрасте мира» родительского контекста при создании, когда schedule.

Примеры

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 с помощью $, что копирует значение напрямую в созданный закрытый фрагмент кода. Это позволяет вставить значение переменной, изолируя асинхронный код от изменений значения переменной в текущей задаче.

Сильно рекомендуется всегда отдавать предпочтение 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"

В настоящее время все задачи в 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, значение повышается как исключение в разбуженной задаче.

Неправильно использовать schedule для произвольной Task, которая уже запущена. Для получения дополнительной информации см. справочник по API.

Примеры

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([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
исходный код
wait(r::Future)

Подождите, пока значение станет доступным для указанного Future.

wait(r::RemoteChannel, args...)

Подождите, пока значение станет доступным в указанном RemoteChannel.

Base.fetchМетод

fetch(t::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)

Сбросить событие в состояние «не установлено». Любые последующие вызовы wait будут блокироваться до тех пор, пока notify не будет вызван снова.

исходный код

Base.SemaphoreТип

Semaphore(sem_size)

Создаёт счётный семафор, который позволяет максимум sem_size приобретений одновременно. Каждое приобретение должно быть сопоставлено с освобождением.

Это обеспечивает упорядочение памяти при приобретении и освобождении при вызовах acquire/release.

исходный код

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)

Приобретение lock при его доступности. Если блокировка уже захвачена другой задачей/потоком, ожидать её освобождения.

Каждое lock должно быть сопоставлено с unlock.

исходный код
lock(f::Function, lock)

Приобретение lock, выполнение f с захваченной lock, и освобождение lock при возвращении f. Если блокировка уже захвачена другой задачей/потоком, ожидать её освобождения.

По завершении выполнения этой функции блокировка lock была освобождена, поэтому вызывающая сторона не должна пытаться её захватить unlock.

Использование Channel в качестве второго аргумента требует Julia 1.7 или более поздней версии.

исходный код

Base.unlockФункция

unlock(lock)

Освобождает владение lock.

Если это рекурсивный замок, который был приобретён ранее, уменьшается внутренний счётчик и сразу возвращается.

исходный код

Base.trylockФункция

trylock(lock) -> Success (Boolean)

Приобретает блокировку, если она доступна, и возвращает true при успехе. Если блокировка уже захвачена другой задачей/потоком, возвращает false.

Каждое успешное trylock должно быть сопоставлено с 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. Одна и та же задача может приобретать блокировку несколько раз. Каждый вызов 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)

Создайте новую задачу из 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}(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):

  • put! на закрытом канале.
  • take! и fetch на пустом, закрытом канале.
исходный код

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 — для задачи, которая ещё не запущена (запланирована). Однако можно использовать 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–2023 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.9/base/parallel/

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API