Spec-Zone.ru › Julia 1.6

Задачи

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.@syncМакрос

@sync

Подождите, пока все лексически вложенные использования @async, @spawn, @spawnat и @distributed будут завершены. Все исключения, брошенные вложенными асинхронными операциями, собираются и выбрасываются как CompositeException.

источник

Base.waitФункция

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, чтобы убедиться, что ожидаемое условие выполняется перед продолжением.

источник

Особое примечание для Threads.Condition:

Вызывающий должен удерживать lock, который владеет Threads.Condition, перед вызовом этого метода. Вызывающая задача будет заблокирована до тех пор, пока какая-то другая задача не разбудит ее, обычно вызывая notify для того же объекта Threads.Condition. Блокировка будет атомарно освобождена при блокировке (даже если она была заблокирована рекурсивно), и будет повторно приобретена перед возвратом.

источник

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)

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

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

источник
lock(f::Function, lock)

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

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

источник

Base.unlockФункция

unlock(lock)

Освобождает владение блокировкой.

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

источник

Base.trylockФункция

trylock(lock) -> Success (Boolean)

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

Каждый успешный захват должен быть сопоставлен с 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, задача, созданная для 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):

  • 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
[...]
исходный код

© 2009–2021 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.6.0/base/parallel/

Spec-Zone.ru

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