Spec-Zone.ru › Julia 1.6

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

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

Для решения этих задач 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 - это ожидающий FIFO-очередь, к которому могут иметь доступ для чтения и записи несколько задач.

Давайте определим задачу производителя, которая производит значения с помощью вызова 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
        @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.",: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.",: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> @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
           @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

Дополнительные операции с задачами

Операции с задачами основаны на примитиве низкого уровня, называемом 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–2021 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.6.0/manual/asynchronous-programming/

Spec-Zone.ru

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