Асинхронное программирование
Когда программе нужно взаимодействовать с внешним миром, например, обмениваться данными с другим компьютером через интернет, операции в программе могут выполняться в непредсказуемом порядке. Предположим, вашей программе нужно загрузить файл. Мы хотели бы инициировать операцию загрузки, выполнить другие операции, ожидая ее завершения, а затем возобновить код, которому нужен загруженный файл, когда он станет доступен. Такая ситуация относится к области асинхронного программирования, иногда также называемого конкуретным программированием (поскольку, по сути, несколько вещей происходят одновременно).
Для решения таких задач Julia предоставляет Task (также известные под другими названиями, такими как симметричные сопрограммы, легкие потоки, кооперативное многозадачность или одноразовые продолжения). Когда часть вычислительной работы (на практике, выполнение определенной функции) обозначена как Task, появляется возможность прервать ее, переключившись на другую Task. Исходную Task можно позже возобновить, и в этот момент она продолжится с того места, где остановилась. Сначала это может показаться похожим на вызов функции. Однако есть два ключевых отличия. Во-первых, переключение задач не использует никакого пространства, поэтому любое количество переключений задач может произойти без потребления стека вызовов. Во-вторых, переключение между задачами может происходить в любом порядке, в отличие от вызовов функций, где вызываемая функция должна завершить выполнение, прежде чем управление вернется к вызывающей функции.
Основные операции с Task
Можно представить себе Task как дескриптор блока вычислительной работы, подлежащей выполнению. Она имеет жизненный цикл: создание, запуск, выполнение, завершение. Задачи создаются путем вызова конструктора Task для функции, которая будет выполняться, или с помощью макроса @task:
julia> t = @task begin; sleep(5); println("done"); end
Task (runnable) @0x00007f13a40c0eb0
@task x эквивалентно Task(()->x).
Эта задача будет ожидать пять секунд, а затем выводить done. Однако она еще не начала выполняться. Мы можем запустить ее, когда захотим, вызвав schedule:
julia> schedule(t);
Если вы попробуете это в REPL, вы увидите, что schedule возвращается немедленно. Это потому, что она просто добавляет t в внутреннюю очередь задач для выполнения. Затем REPL выведет следующую подсказку и будет ожидать дополнительного ввода. Ожидание ввода с клавиатуры предоставляет возможность выполнения других задач, поэтому в этот момент t начнется. t вызывает sleep, который устанавливает таймер и останавливает выполнение. Если были запланированы другие задачи, они могут быть выполнены в этот момент. Через пять секунд таймер срабатывает и перезапускает t, и вы увидите, что выводится done. t тогда завершается.
Функция wait блокирует вызывающую задачу до тех пор, пока какая-то другая задача не завершит работу. Таким образом, если вы напечатаете
julia> schedule(t); wait(t)
вместо того, чтобы только вызвать schedule, вы увидите пятисекундную паузу, прежде чем появится следующая подсказка для ввода. Это потому, что REPL ожидает завершения t перед продолжением.
Часто возникает необходимость создать задачу и сразу же запланировать ее, поэтому для этой цели предоставлен макрос @async — @async x эквивалентно schedule(@task x).
Взаимодействие с каналами
В некоторых задачах различные части необходимой работы не связаны естественным образом вызовами функций; среди работ, которые нужно выполнить, нет очевидного «отправителя» или «получателя». Примером является задача «производитель-потребитель», где одна сложная процедура генерирует значения, а другая сложная процедура их потребляет. Потребитель не может просто вызвать функцию производителя, чтобы получить значение, поскольку производитель может иметь больше значений для генерации, поэтому он может еще не быть готов вернуть значение. С помощью задач производитель и потребитель могут оба работать столько, сколько им нужно, передавая значения друг другу по мере необходимости.
Julia предоставляет механизм Channel для решения этой проблемы. Channel — это очередь «первым вошел, первым вышел», для ожидания, к которой могут обращаться несколько задач для чтения и записи.
Давайте определим задачу производителя, которая производит значения с помощью вызова put!. Чтобы потреблять значения, нам нужно запланировать выполнение производителя в новой задаче. Можно использовать специальный конструктор Channel, который принимает функцию с одним аргументом в качестве аргумента, чтобы запустить задачу, связанную с каналом. Затем мы можем получать значения из объекта канала многократно с помощью take!:
julia> function producer(c::Channel)
put!(c, "start")
for n=1:4
put!(c, 2n)
end
put!(c, "stop")
end;
julia> chnl = Channel(producer);
julia> take!(chnl)
"start"
julia> take!(chnl)
2
julia> take!(chnl)
4
julia> take!(chnl)
6
julia> take!(chnl)
8
julia> take!(chnl)
"stop"
Один из способов понять это поведение заключается в том, что producer смог возвращать значения многократно. Между вызовами put!, выполнение производителя приостанавливается, и потребитель получает контроль.
Возвращаемый Channel может быть использован как итерируемый объект в цикле for, в этом случае переменная цикла принимает все произведенные значения. Цикл завершается, когда канал закрывается.
julia> for x in Channel(producer)
println(x)
end
start
2
4
6
8
stop
Обратите внимание, что нам не пришлось явно закрывать канал у производителя. Это потому, что при связывании Channel с Task открытое время жизни канала связывается с временем жизни связанной задачи. Объект канала закрывается автоматически при завершении задачи. К одной задаче может быть привязано несколько каналов, и наоборот.
Хотя конструктор Task ожидает функцию без аргументов, метод Channel, который создает связанный с задачей канал, ожидает функцию, которая принимает один аргумент типа Channel. Типичная модель для параметризованного производителя заключается в применении частичной функции для создания функции с нулем или одним аргументом анонимной функции.
Для объектов Task это можно сделать напрямую или с помощью удобного макроса:
function mytask(myarg)
...
end
taskHdl = Task(() -> mytask(7))
# or, equivalently
taskHdl = @task mytask(7)
Для организации более сложных шаблонов распределения работы можно использовать bind и schedule в сочетании с Task и Channel конструкторами, чтобы явно связать набор каналов с набором задач производителя/потребителя.
Подробнее о каналах
Канал можно представить как трубу, то есть у него есть вход для записи и вход для чтения:
Несколько пишущих задач в разных задачах могут одновременно записывать в один и тот же канал с помощью вызовов
put!.Несколько читающих задач в разных задачах могут одновременно считывать данные с помощью вызовов
take!.-
Пример:
# Given Channels c1 and c2, c1 = Channel(32) c2 = Channel(32) # and a function `foo` which reads items from c1, processes the item read # and writes a result to c2, function foo() while true data = take!(c1) [...] # process data put!(c2, result) # write out result end end # we can schedule `n` instances of `foo` to be active concurrently. for _ in 1:n errormonitor(@async foo()) end Каналы создаются с помощью конструктора
Channel{T}(sz). Канал будет содержать только объекты типаT. Если тип не указан, канал может содержать объекты любого типа.szотносится к максимальному числу элементов, которые могут храниться в канале в любой момент времени. Например,Channel(32)создает канал, который может содержать максимум 32 объекта любого типа. КаналChannel{MyType}(64)может содержать до 64 объектов типаMyTypeв любой момент времени.Если
Channelпуст, читающие задачи (в вызовеtake!) будут блокироваться до тех пор, пока данные не появятся.Если
Channelполон, пишущие задачи (в вызовеput!) будут блокироваться до тех пор, пока место не освободится.isreadyпроверяет наличие каких-либо объектов в канале, аwaitожидает появления объекта.-
Канал
Channelизначально находится в открытом состоянии. Это означает, что к нему можно свободно обращаться для чтения и записи с помощью вызововtake!иput!.closeзакрываетChannel. В закрытомChannelвызовput!завершится ошибкой. Например:julia> c = Channel(2); julia> put!(c, 1) # `put!` on an open channel succeeds 1 julia> close(c); julia> put!(c, 2) # `put!` on a closed channel throws an exception. ERROR: InvalidStateException: Channel is closed. Stacktrace: [...]
-
take!иfetch(которая извлекает, но не удаляет значение) для закрытого канала успешно возвращают любые существующие значения, пока он не опустеет. Продолжая приведенный выше пример:julia> fetch(c) # Any number of `fetch` calls succeed. 1 julia> fetch(c) 1 julia> take!(c) # The first `take!` removes the value. 1 julia> take!(c) # No more data available on a closed channel. ERROR: InvalidStateException: Channel is closed. Stacktrace: [...]
Рассмотрим простой пример использования каналов для межзадачного обмена данными. Мы запускаем 4 задачи для обработки данных из одного канала jobs . Задачи, идентифицированные по идентификатору (job_id), записываются в канал. Каждая задача в этом имитировании считывает job_id, ждет случайное время и записывает назад кортеж из job_id и имитированного времени в канал результатов. Наконец, выводятся все results.
julia> const jobs = Channel{Int}(32);
julia> const results = Channel{Tuple}(32);
julia> function do_work()
for job_id in jobs
exec_time = rand()
sleep(exec_time) # simulates elapsed time doing actual work
# typically performed externally.
put!(results, (job_id, exec_time))
end
end;
julia> function make_jobs(n)
for i in 1:n
put!(jobs, i)
end
end;
julia> n = 12;
julia> errormonitor(@async make_jobs(n)); # feed the jobs channel with "n" jobs
julia> for i in 1:4 # start 4 tasks to process requests in parallel
errormonitor(@async do_work())
end
julia> @elapsed while n > 0 # print out results
job_id, exec_time = take!(results)
println("$job_id finished in $(round(exec_time; digits=2)) seconds")
global n = n - 1
end
4 finished in 0.22 seconds
3 finished in 0.45 seconds
1 finished in 0.5 seconds
7 finished in 0.14 seconds
2 finished in 0.78 seconds
5 finished in 0.9 seconds
9 finished in 0.36 seconds
6 finished in 0.87 seconds
8 finished in 0.79 seconds
10 finished in 0.64 seconds
12 finished in 0.5 seconds
11 finished in 0.97 seconds
0.029772311
Вместо errormonitor(t), более надежное решение может заключаться в использовании bind(results, t), так как это позволит не только регистрировать любые непредвиденные ошибки, но и заставить связанные ресурсы закрыться и распространить исключение повсюду.
Дополнительные операции с задачами
Операции с задачами основаны на примитиве низкого уровня, называемом yieldto. yieldto(task, value) приостанавливает текущую задачу, переключается на указанную task, и вызывает последний вызов yieldto этой задачи для возвращения указанного value. Обратите внимание, что yieldto является единственной операцией, необходимой для использования управления потоком задач; вместо вызова и возвращения мы всегда просто переключаемся на другую задачу. Именно поэтому эта функция также называется «симметричными сопрограммами»; каждая задача переключается туда и обратно с помощью одного и того же механизма.
yieldto мощная, но большинство применений задач не вызывают её напрямую. Подумайте, почему это может быть так. Если вы переключаетесь с текущей задачи, вы, вероятно, захотите вернуться к ней в какой-то момент, но знание того, когда вернуться, и знание того, какая задача несет ответственность за возврат, может потребовать значительной координации. Например, put! и take! являются блокирующими операциями, которые, когда используются в контексте каналов, сохраняют состояние, чтобы запомнить, кто являются потребителями. Не нужно вручную отслеживать задачу-потребителя, что делает put! проще в использовании, чем низкоуровневый yieldto.
Помимо yieldto, для эффективного использования задач необходимы ещё несколько основных функций.
-
current_taskполучает ссылку на текущую выполняемую задачу. -
istaskdoneпроверяет, завершилась ли задача. -
istaskstartedпроверяет, была ли задача запущена. -
task_local_storageманипулирует хранилищем пар «ключ-значение», специфичным для текущей задачи.
Задачи и события
Большинство переключений задач происходит в результате ожидания событий, таких как запросы ввода-вывода, и выполняются планировщиком, включённым в Julia Base. Планировщик поддерживает очередь выполняемых задач и выполняет цикл обработки событий, который перезапускает задачи на основе внешних событий, таких как приход сообщения.
Основная функция для ожидания события — wait. Несколько объектов реализуют wait; например, в случае объекта Process, wait будет ждать его завершения. wait часто подразумевается; например, wait может произойти внутри вызова read для ожидания появления данных.
Во всех этих случаях wait в конечном итоге работает с объектом Condition, который отвечает за очереди и перезапуск задач. Когда задача вызывает wait на объекте Condition, задача помечается как невыполняемая, добавляется в очередь условия и переключается на планировщик. Затем планировщик выберет другую задачу для выполнения или заблокируется в ожидании внешних событий. При успешном выполнении обработчик событий в конечном итоге вызовет notify на условии, что снова сделает ожидающие это условие задачи выполнимыми.
Задача, созданная явно с помощью вызова Task, изначально неизвестна планировщику. Это позволяет вам управлять задачами вручную с помощью yieldto, если вы этого хотите. Однако, когда такая задача ожидает события, она всё равно автоматически перезапускается при его появлении, как ожидалось.
© 2009–2024 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.10/manual/asynchronous-programming/