Spec-Zone.ru › Julia 1.9

Асинхронное программирование

Когда программе нужно взаимодействовать с внешним миром, например, обмениваться сообщениями с другой машиной через интернет, операции в программе могут выполняться в непредсказуемом порядке. Представьте, что вашей программе нужно загрузить файл. Мы хотим инициировать операцию загрузки, выполнить другие операции, пока ждем ее завершения, а затем возобновить код, которому нужен загруженный файл, когда он будет доступен. Такая ситуация относится к области асинхронного программирования, иногда также называемого параллельным программированием (поскольку, концептуально, несколько вещей происходят одновременно).

Для решения этих задач Julia предоставляет Task (также известные под другими названиями, такими как симметричные сопрограммы, легкие потоки, кооперативное многозадачность или одноразовые продолжения). Когда часть вычислительной работы (на практике выполнение определенной функции) обозначена как Task, становится возможным прервать ее, переключившись на другую Task. Исходную Task можно позже возобновить, в этот момент она продолжит работу с того места, где остановилась. На первый взгляд, это может показаться похожим на вызов функции. Однако есть два ключевых отличия. Во-первых, переключение задач не использует памяти, поэтому любое количество переключений задач может происходить без потребления стека вызовов. Во-вторых, переключение между задачами может происходить в любом порядке, в отличие от вызовов функций, где вызываемая функция должна завершить выполнение, прежде чем управление вернется к вызывающей функции.

Основные операции с Task

Вы можете рассматривать Task как указатель на единицу вычислительной работы, которая должна быть выполнена. Она имеет жизненный цикл: создание-запуск-выполнение-завершение. Задачи создаются путем вызова конструктора Task для выполнения функции с 0 аргументами или с помощью макроса @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, принимающий функцию с 1 аргументом в качестве аргумента, может использоваться для запуска задачи, связанной с каналом. Затем мы можем 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 ожидает функцию с 0 аргументами, метод Channel, создающий привязанный к задаче канал, ожидает функцию, которая принимает один аргумент типа Channel. Общая схема заключается в том, что производитель параметризован, в этом случае необходима частичная применённая функция для создания функции с 0 или 1 аргументом анонимной функции.

Для объектов 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–2023 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.9/manual/asynchronous-programming/

Spec-Zone.ru

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