Spec-Zone.ru › Julia 1.5

Задачи

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.

source

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!.

source

Base.take!Метод

take!(c::Channel)

Удаляет и возвращает значение из Channel. Блокирует, пока данные не станут доступны.

Для небуферизованных каналов блокирует, пока другая задача не выполнит put!.

source

Base.isreadyМетод

isready(c::Channel)

Определяет, есть ли значение в Channel. Возвращает результат сразу, не блокирует.

Для небуферизованных каналов возвращает true если есть задачи, ожидающие put!.

source

Base.fetchМетод

fetch(c::Channel)

Ожидает и получает первый доступный элемент из канала. Элемент не удаляется. fetch не поддерживается для небуферизованных (размером 0) каналов.

source

Base.closeМетод

close(c::Channel[, excp::Exception])

Закрывает канал. Исключение (опционально заданное excp):

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

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/

Spec-Zone.ru

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