Spec-Zone.ru › Julia 1.7

Задачи

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

исходный код
END_OF_DOCUMENT_MARKER

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) генерируется:

  • 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.7.0/base/parallel/

Spec-Zone.ru

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