PYTHON ASYNCIO: цикл событий и FUTURE, TASK
Цикл событий – это цикл, который может: регистрировать задачи для выполнения, выполняет их, задерживает или даже отменяет их и обрабатывать различные события, связанные с этими операциями. Обычно мы планируем несколько асинхронных функций в цикле событий. Цикл запускает одну функцию, в то время как эта функция ожидает ввода-вывода, приостанавливает ее и запускает другую. Когда первая функция завершает ввод-вывод, она возобновляется. Таким образом, две или более функции могут работать совместно. Это главная цель цикла событий.
Цикл событий также может передавать ресурсоемкие функции в пул потоков для обработки. Внутренние элементы цикла событий довольно сложны, и нам не нужно беспокоиться об этом. Нам просто нужно помнить, что цикл обработки событий – это механизм, с помощью которого мы можем планировать наши асинхронные функции и запускать их.
Futures и Tasks (Задачи)
Future – это объект, который должен иметь результат в будущем. Это такое себе обещание на получения результата выполнения, при этом может быть получен не только успешный результат, но и ошибка возникшая в результате выполнения.
Задача (Task) – это подкласс Future, который служит оберткой для корутины. Когда корутина завершается, результат Задания становится доступным.
Корутина – это способ приостановить функцию и периодически возвращать серию значений. Функция созданная как корутина может приостановить свое выполнение, используя в выражении ключевые слова yield, yield from или await (python 3.5+). Функция приостанавливается до тех пор, пока оператор yield не получит значение.
Объединение цикла событий и задач
Все предельно просто. Нам нужен цикл событий, и нам нужно зарегистрировать наши объекты Future / Task в цикле событий. Цикл будет планировать и запускать их. Мы можем добавить обратные вызовы к нашим объектам Future / Task, чтобы мы могли получать уведомления, когда появятся результаты выполнения.
Типовым примером можн считать использование корутин, которые оборачиваются Future и далее создается объект Task. Когда корутина выполняет оператор yield, она приостанавливается. Когда ей передается значение, она возобновляет свою работу. Когда кортина завершается return , задача завершена и получает результирующие значение. Любой связанный обратный вызов выполняется. Если корутина вызывает исключение, задача завершается с ошибкой.
import asyncio
@ asyncio . coroutine
def slow_operation ( ) :
# yield from suspends execution until
# there’s some result from asyncio.sleep
yield from asyncio . sleep ( 1 )
# our task is done, here’s the result
return ‘Future is done!’
def got_result ( future ) :
print ( future . result ( ) )
# Our main event loop
loop = asyncio . get_event_loop ( )
# We create a task from a coroutine
task = loop . create_task ( slow_operation ( ) )
# Please notify us when the task is complete
task . add_done_callback ( got_result )
# The loop will close when the task has resolved
loop . run_until_complete ( task )
В примере, декоратор @asyncio.coroutine – преобразует генератор в корутину. Операция loop.create_task(slow_operation()) – создает задачу из корутины созданной вызовом slow_operation. Операция task.add_done_callback(got_result) – регистрирует обратный вызов для целей обработки результатов. Метод loop.run_until_complete(task) – запускает цикл событий, регистрирует в нем задачу и ждет ее завершение. После завершения задачи цикл событий завершается.
Функция run_until_complete отличный метод для управления циклом событий. Конечно возможен и такой вариант:
import asyncio
async def slow_operation ( ) :
await asyncio . sleep ( 1 )
return ‘Future is done!’
def got_result ( future ) :
print ( future . result ( ) )
# We have result, so let’s stop
loop . stop ( )
loop = asyncio . get_event_loop ( )
task = loop . create_task ( slow_operation ( ) )
task . add_done_callback ( got_result )
# We run forever
loop . run_forever ( )
Цикл событий запускается в бесконечном режиме, но в обработчике результатов мы завершаем цикл.
Важно понимание, что в asyncio нет никакой магии при реализации неблокирующих задач. Во время реализации asyncio стоял отдельно в стандартной библиотеке, т.к. остальные модули предоставляли только блокирующую функциональность. Вы можете использовать модуль concurrent.futures для оборачивания блокирующих задач в потоки или процессы и получения Future для использования в asyncio. Несколько таких примеров доступны на GitHub.
Это, наверно, главный недостаток сейчас при использовании asyncio, однако уже есть несколько библиотек, помогающих решить эту проблему.
Вам не нужна функция ensure_future если Вам не нужен результат и вы не хотите получать исключения возникшие во время выполнения.
Метод run_in_executor позволяет передать Executor и таким образом контролировать количество рабочих процессов.
Сравнение Promise из JavaScript и Future/Tast из Python
Стоит сразу заметить, что Promise в JS это не то же самое, что Future в Python. Куда более эквивалентны Promise в JS и Task в Python. При этом стоит понимать, что задача (Task) в python, состоит из Future и Corutine. Другими словами задача состоит из обещания вернуть какое-то значение и кода который ее выполнит. При этом главное отличие Future от Promise из JS, состоит в том, что Promise автоматически выполнится асинхронно, в то время как Future не выполняется вообще это просто контейнер в который значение будет помещено некоторым внешним кодом через метод set_result(), при этом Future сразу же вызовет зарегистрированные обратные вызовы (вызов регистрируются через метод add_done_callback())
В то время как Corutine основанная на тех же принципах, что и генераторы. Корутина вызывается из цикла обработки событий и возвращает объекты Future и при этом ожидает когда это событие будет завершено. После завершения в корутину может быть передано значение и получено из нее новый Future. (по анологии с циклом for и генератором который выбрасывет через yield значения).
Выражение await практически идентично yield from, поэтому, ожидая другую сопрограмму, вы останавливаетесь, пока у этой сопрограммы не будет разрешен Future, и вы не получаете возвращаемое значение. Future является повторяемым в один такт, и его итератор возвращает фактическое Future – это примерно означает, что await future равно yield from future и yield future.
Task это Future, который фактически был запущен для вычисления и привязан к циклу событий. Так что это особый вид Future (класс Task является производным от класса Future), который связан с некоторым циклом событий и имеет некоторую корутину (сопрограмму), которая служит исполнителем Task.
Задача обычно создается объектом цикла события: вы предоставляете корутину (сопрограмму) циклу, она создает объект задачи и начинает выполнять итерацию по этой корутине, как описано выше. Как только корутина закончена, Future задачи разрешается (получает результат) с помощью возвращаемого значения корутины.
Видите ли, задача очень похожа на JS Promise – она инкапсулирует фоновое задание и его результат.
Coroutine func – это фабрика корутин (сопрограмм), подобная функции генератора для генераторов. Обратите внимание на разницу между функцией корутин Python и асинхронной функцией Javascript – при вызове асинхронная функция JS создает Promise, и ее внутренний генератор немедленно начинает выполнятся, в то время как корутина Python ничего не делает, пока на ней не будет создана задача.
В Python весь код которому нужна какая-либо функция asyncio, должен всегда работает цикл обработки событий.
Однако, смешивать синхронный и асинхронный код довольно сложно – вся ваша программа должна быть асинхронной (но вы можете запускать куски синхронного кода в отдельных потоках через API asyncio threadpool API)
asyncio — параллелизм в Python
Параллелизм в Python — одна из самых сложных тем для понимания, не говоря уже о реализации. Не помогает и то, что существует множество способов создания параллельных программ. Возникает куча вопросов. Нужно ли запускать несколько потоков? Использовать несколько процессов? Использовать асинхронное программирование?
Что ж, ответ здесь один — использовать тот способ, который лучше всего подходит для вашего случая. Но если вы сомневаетесь, то используйте асинхронный ввод-вывод, когда это возможно, и потоковое программирование, когда это необходимо.
В этой статье мы рассмотрим асинхронные программы как в старых версиях Python (на случай, если вы имеете дело с устаревшим кодом), так и в новых.
Что такое asyncio?
Asyncio означает «асинхронный ввод-вывод» и относится к парадигме программирования, позволяющей достичь высокого параллелизма с помощью одного потока или цикла событий. Эта модель не является исключительно «питоничной» и реализуется также в других языках и фреймворках, наиболее известным из которых является NodeJS на JavaScript.
Понимание asyncio на примере
Чтобы понять концепцию asyncio, рассмотрим ресторан с одним официантом. Внезапно появляются три клиента, Кохли, Амир и Джон. После того как они получили от официанта меню, им требуется разное количество времени для принятия решения о том, что они будут есть.
Предположим, что Кохли требуется 5 минут, Амиру — 10 минут, а Джону — 1 минута. Если один официант начинает с Амира, то он принимает его заказ через 10 минут. Затем он обслуживает Кохли и тратит 5 минут на то, чтобы записать его заказ. Наконец, официант тратит еще 1 минуту на то, чтобы узнать, что хочет съесть Джон. Таким образом, в общей сложности он тратит 10 + 5 + 1 = 16 минут на то, чтобы записать их заказы. Однако обратите внимание, что в этой последовательности событий Джон ждет 15 минут, пока официант доберется до него, Кохли ждет 10 минут, а Амир ждет 0 минут.
Теперь подумаем, знает ли официант время, которое потребуется каждому клиенту для принятия решения. Он может начать с Джона, затем перейти к Амиру и, наконец, к Кохли. Таким образом, каждый клиент будет ждать 0 минут. Создается иллюзия трех официантов, по одному на каждого клиента, хотя на самом деле есть только один. Наконец, общее время, необходимое официанту для принятия всех трех заказов, составит 10 минут, что гораздо меньше, чем 16 минут при первом сценарии.
Тем, кто знаком с JavaScript, asyncio покажется очень похожим на работу NodeJS. В NodeJS под капотом находится однопоточный цикл событий, который обслуживает все входящие запросы.
Зачем использовать asyncio вместо многопоточности в Python?
Во-первых, очень сложно написать код, безопасный для потоков. В асинхронном коде вы точно знаете, где код будет переходить от одной задачи к другой, и возникновение состояния гонки маловероятнее.
Во-вторых, потоки потребляют достаточно много данных, поскольку каждый поток должен иметь свой стек. В асинхронном коде весь код использует один и тот же стек, и его размер остается небольшим за счет постоянного разматывания между задачами.
В-третьих, потоки являются структурами ОС и поэтому требуют большего объема памяти для поддержки платформой. С асинхронными задачами такой проблемы нет.
Как создавать асинхронные программы в старой версии Python
Итак, вы приступили к новой работе и обнаружили, что кодовая база изобилует устаревшим кодом на языке Python. В этом разделе мы познакомим вас со старым способом создания асинхронных программ.
Здесь есть о чем рассказать, поэтому давайте просто погрузимся в тему. Первое понятие, с которым вам стоит познакомиться, — это итерируемые объекты и итераторы. Они служат основой для генераторов, которые открыли двери для асинхронного программирования.
От редакции Pythonist: также предлагаем почитать статью «Итераторы и генераторы в Python».
Итерируемые объекты и итераторы
В Python итерируемый объект — это объект, элементы которого можно перебирать с помощью цикла for . Итератор способен возвращать свои члены по одному, при этом наиболее распространенным типом итераторов в Python являются последовательности, включающие списки, строки и кортежи.
Когда мы хотим получить элемент по индексу, вызывается метод __getitem__() . Помните, что не каждый тип в Python является последовательностью. Словари, множества, файловые объекты и генераторы не могут быть проиндексированы, но они являются итерируемыми. Python также позволяет создавать бесконечные итераторы, называемые генераторами.
Для того чтобы считаться итерируемым, объект должен определять один из двух методов:
- __iter__()
- __getitem__()
Итератор — это объект, который может использоваться для последовательного доступа к элементам итерируемого объекта. В Python 3 итератор предоставляет метод __next__() , а в Python 2 — метод next() . Оба метода извлекают следующий элемент из последовательности итерируемого объекта. Примечание: итератор должен поддерживать следующие методы:
Объект итератора возвращает самого себя для метода __iter__() . Это позволяет нам использовать итератор и итерируемый объект в цикле for .
Перебрав элементы объекта до конца, функция next() выбрасывает исключение StopIteration . В совокупности эти правила называются протоколом итератора. Метод __iter__() для контейнера может также возвращать так называемый генератор, который тоже является итератором.
Оператор yield()
Используемый в Python оператор yield может как выдавать значения, так и принимать параметры. Это становится особенно важным при создании функций-генераторов.
От редакции Pythonist: также предлагаем почитать «Сравнение операторов yield и return в Python (с примерами)».
Рассмотрим следующий код, который возвращает строку:
def keep_learning_synchronous(): return "Educative" if __name__ == "__main__": str = keep_learning_synchronous() print(str)
Заменив return на yield , вы заметите, что возвращаемый объект — это объект-генератор. Фактически наш метод keep_learning_asynchronous() теперь является генераторной функцией.
Функции-генераторы называются генераторами, поскольку они генерируют значения. Для того чтобы объект-генератор выдал строку из приведенного выше фрагмента кода, можно вызвать для него функцию next() .
Мы можем использовать yield в функции в виде yield . Оператор yield позволяет функции возвращать значение и приостанавливать состояние функции до тех пор, пока не будет вызвана функция next() на связанном с ней объекте-генераторе.
Операторы генераторов
Функции, содержащие выражение yield , компилируются как генераторы. Использование выражения yield в теле функции приводит к тому, что эта функция становится генератором. Эти функции возвращают объект, поддерживающий методы протокола итерации.
Созданный объект-генератор автоматически получает метод __next__() . Давайте вернемся к примеру из предыдущего раздела. Вместо использования метода next() мы можем вызвать __next__ непосредственно на объекте-генераторе:
def keep_learning_asynchronous(): yield "Educative" if __name__ == "__main__": gen = keep_learning_asynchronous() str = gen.__next__() print(str)
О генераторах нужно помнить следующие факты:
- Функции-генераторы позволяют откладывать получение сложновычисляемых значений. Следующее значение вычисляется только при необходимости. То есть генераторы не сохраняют в памяти длинные последовательности и не выполняют все дорогостоящие вычисления заранее. Это делает генераторы эффективными с точки зрения памяти и вычислений.
- Генераторы, будучи приостановленными, сохраняют местоположение кода, в котором был выполнен последний оператор yield, и всю свою локальную область видимости. Это позволяет им возобновить выполнение с того места, на котором они остановились.
- Объекты-генераторы — это не что иное, как итераторы.
- Следует различать функцию-генератор и связанный с ней объект-генератор. Функция-генератор при вызове возвращает объект-генератор, а функция next() вызывается на объекте-генераторе для выполнения кода внутри функции-генератора.
Состояния генератора
Генератор проходит через следующие состояния:
- GEN_CREATED, когда объект генератора был возвращен впервые из функции генератора и итерация еще не началась.
- GEN_RUNNING, когда для объекта генератора был вызван next , и он выполняется интерпретатором Python.
- GEN_SUSPENDED, когда генератор приостановлен на выходе.
- GEN_CLOSED, когда генератор завершил выполнение или был закрыт.
Методы для объектов генераторов
Объект генератора предоставляет различные методы, которые могут быть вызваны для работы с генератором. Например, методы throw() , send() и close() .
Корутины на основе генератора
Python сделал различие между генераторами Python и генераторами, которые предназначены для использования как корутины (асинхронные функции). Такие корутины называются генераторными и требуют добавления декоратора @asynio.coroutine в определение функции, хотя это не является строгим требованием.
В генераторных корутинах вместо синтаксиса yield используется синтаксис yield from .
От редакции Pythonist: также предлагаем почитать «Конструкция yield from».
- выходить из другой корутины
- возвращать выражение
- вызывать исключение
Корутины в Python позволяют реализовать кооперативную многозадачность. Кооперативная многозадачность — это подход, при котором выполняющаяся задача добровольно уступает поток другим процессам. Задача может сделать это, когда она логически заблокирована, например, в ожидании пользовательского ввода или когда она инициировала сетевой запрос и будет простаивать некоторое время.
Корутину можно определить как специальную функцию, которая может передавать управление вызывающему ее процессу без потери своего состояния.
В чем же разница между короутинами и генераторами?
Генераторы — это, по сути, итераторы, хотя внешне они похожи на функции. В общем случае различие между генераторами и корутинами заключается в следующем:
- Генераторы возвращают значение, в то время как корутина передает управление другой корутине и может возобновить выполнение с того момента, когда она передала управление.
- Генератор не может принимать аргументы после запуска, в то время как корутина может.
- Генераторы используются в основном для упрощения написания итераторов. Они являются разновидностью корутин и иногда также называются семикорутинами.
Пример корутины на основе генератора
Простейшая генераторная корутина, которую мы можем написать, выглядит следующим образом:
@asyncio.coroutine def do_something_important(): yield from asyncio.sleep(1)
Корутина находится в состоянии сна в течение одной секунды. Обратите внимание на декоратор и использование yield from . Без них вы не смогли бы использовать корутину с asyncio.
Оператор yield from передает управление обратно циклу событий и возобновляет выполнение после завершения работы корутины asyncio.sleep() .
Заметим, что asyncio.sleep() сама является корутиной. Модифицируем эту корутину так, чтобы она вызывала другую корутину, выполняющую сон. Изменения показаны ниже:
@asyncio.coroutine def go_to_sleep(sleep): print("sleeping for " + str(sleep) + " seconds") yield from asyncio.sleep(sleep) @asyncio.coroutine def do_something_important(sleep): # what is more important than getting # enough sleep! yield from go_to_sleep(sleep)
Теперь представьте, что вы трижды последовательно вызываете корутину do_something_important() со значениями 1, 2 и 3 соответственно.
Без использования потоков или мультипроцессинга последовательный код будет выполнен за 1 + 2 + 3 = 6 секунд. Но если использовать asyncio, то тот же самый код может быть выполнен примерно за 3 секунды, несмотря на то, что все вызовы выполняются в одном потоке.
Интуитивно понятно, что при возникновении блокирующей операции управление передается обратно в цикл событий, и выполнение возобновляется только после завершения блокирующей операции.
В случае Python генераторы используются как производители данных, а корутины — как их потребители. Раньше, до появления в Python 3.5 поддержки собственных корутин, корутины реализовывались с помощью генераторов.
Объекты и тех, и других имеют тип generator. Однако, начиная с версии 3.5, Python делает различие между корутинами и генераторами.
Как создавать асинхронные программы на Python 3
Существует три основных элемента для создания асинхронных программ на Python: нативные корутины, циклы событий (event loops) и фьючерсы (futures). Давайте углубимся в каждый из них и рассмотрим подробнее.
Нативные корутины
В Python 3.5 в языке появилась поддержка нативных корутин. Под «нативными» подразумевается, что в языке появился синтаксис для специального определения корутин. Нативные корутины могут быть определены с помощью синтаксиса async/await .
Прежде чем перейти к более подробному описанию, приведем пример очень простой нативной корутины:
async def coro(): await asyncio.sleep(1)
Эта корутина может быть запущена с помощью цикла событий следующим образом:
loop = asyncio.get_event_loop() loop.run_until_complete(coro())
Async
Создать нативную корутину можно с помощью async def . Метод с префиксом async def автоматически становится корутиной.
async def useless_native_coroutine(): pass
Метод inspect.iscoroutine() вернет True для объекта coroutine, возвращенного из приведенной выше функции coroutine. Заметьте, что yield или yield from не могут появляться в теле асинхронной функции, иначе это будет отмечено как синтаксическая ошибка.
import inspect import asyncio async def useless_native_coroutine(): pass if __name__ == "__main__": coro = useless_native_coroutine() print(inspect.iscoroutine(coro)) //Returns True
Await
Функция await может использоваться для получения результата выполнения объекта coroutine. Вы используете команду await следующим образом: await , где должен быть awaitable объектом.
Awaitable объекты должны реализовывать метод __await__() , который должен возвращать итератор. Если вы помните, yield from также ожидает, что его аргументом будет итератор, из которого можно получить итератор. Под капотом await заимствует реализацию yield from с дополнительной проверкой, действительно ли его аргумент является awaitable.
Следующие объекты принадлежат к типу awaitable объектов:
- объект нативной корутины, возвращаемый при вызове функции нативной корутины
- объект coroutine на основе генератора, возвращаемый из генератора, декорированного @types.coroutine или @asyncio.coroutine . Декорированные генераторные корутины являются awaitable, даже если у них нет метода __await__() .
- объекты Future
- объекты Task (Task является подклассом Future).
- объекты, определенные с помощью CPython C API, имеют функцию tp_as_async.am_await() , возвращающую итератор (аналогично методу __await__() ).
Кроме того, await должен появляться внутри async-определенного метода, иначе это синтаксическая ошибка. В настоящее время:
- генераторы используются для обозначения функций, которые только производят значения,
- ванильные корутины только получают значения,
- корутины на основе генераторов идентифицируются по наличию yield from в теле метода,
- нативные корутины определяются с использованием синтаксиса async/await .
Подытожим данный раздел:
- Для возврата значений в генераторах используется yield
- Генераторы, которые могут получать значения извне, являются корутинами
- Генераторы с yield from в теле функции являются генераторными корутинами, а методы, определенные с помощью async def , — нативными корутинами.
- Используйте декораторы @asyncio.coroutine или @types.coroutine для генераторных корутинов, чтобы сделать их совместимыми с нативными корутинами.
Циклы событий
Цикл событий (event loop) — это программная конструкция, которая ожидает возникновения событий, а затем передает их обработчику событий.
Событием может быть клик пользователя на кнопку интерфейса пользователя, запуск процесса загрузки файла или получение ответа от стороннего сервера.
В основе асинхронного программирования лежат циклы событий.
Эта концепция не является чем-то новым. На самом деле многие языки программирования обеспечивают асинхронное программирование с помощью циклов событий. В Python циклы событий запускают асинхронные задачи и колбэки, выполняют сетевые операции ввода-вывода, запускают подпроцессы и делегируют дорогостоящие вызовы функций пулу потоков.
Один из наиболее распространенных примеров использования этой концепции — это веб-серверы, реализованные с использованием асинхронного дизайна. Веб-сервер ожидает получения HTTP-запроса и возвращает соответствующий ресурс. Те, кто знаком с JavaScript, могут вспомнить, что NodeJS работает по тому же принципу: это веб-сервер, который запускает цикл событий для получения веб-запросов в одном потоке.
Некоторые другие веб-серверы для обработки каждого запроса создают новый поток или форк нового процесса. В некоторых бенчмарках асинхронные веб-серверы на основе цикла событий превосходили многопоточные, что может показаться нелогичным на первый взгляд.
Запуск цикла событий
В Python 3.7+ предпочтительным способом запуска цикла событий является использование метода asyncio.run() . Метод представляет собой блокирующий вызов, который блокирует до тех пор, пока не завершится переданная корутина. Пример программы:
async def do_something_important(): await asyncio.sleep(10) if __name__ == "__main__": asyncio.run(do_something_important())
Примечание. Если вы работаете с Python 3.5, то API asyncio.run() недоступен. В этом случае необходимо явно получить цикл событий с помощью asyncio.new_event_loop() и запустить желаемую корутину с помощью run_until_complete() , определенной для объекта цикла.
Запуск нескольких циклов событий
Вам никогда не придется самостоятельно запускать цикл событий. Вместо этого следует использовать API более высокого уровня для отправки корутин. В учебных целях мы продемонстрируем запуск цикла событий в одном потоке.
В приведенном ниже примере кода используется API asyncio.new_event_loop() для получения нового цикла событий и последующего использования его для запуска другой корутины.
import asyncio, random from threading import Thread from threading import current_thread async def do_something_important(sleep_for): print("Is event loop running in thread = \n".format(current_thread().getName(), asyncio.get_event_loop().is_running())) await asyncio.sleep(sleep_for) def launch_event_loops(): # get a new event loop loop = asyncio.new_event_loop() # set the event loop for the current thread asyncio.set_event_loop(loop) # run a coroutine on the event loop loop.run_until_complete(do_something_important(random.randint(1, 5))) # remember to close the loop loop.close() if __name__ == "__main__": t1 = Thread(target=launch_event_loops) t2 = Thread(target=launch_event_loops) t1.start() t2.start() print("Is event loop running in thread = \n".format(current_thread().getName(), asyncio.get_event_loop().is_running())) t1.join() t2.join()
Попробуйте запустить код, изучите вывод, и вы поймете, что каждый порожденный поток выполняет свой собственный цикл событий.
Типы циклов событий
Существует два типа циклов событий:
- SelectorEventLoop
- ProactorEventLoop
SelectorEventLoop основан на модуле selectors и является циклом по умолчанию на всех платформах. Модуль selectors содержит API-интерфейсы poll() и select() .
ProactorEventLoop , напротив, использует порты завершения ввода/вывода Windows и поддерживается только на Windows. Мы не будем вдаваться в более тонкие детали реализации обоих типов, но закончим здесь замечанием, что как тип, так и связанная с ним политика управления циклом контролируют поведение цикла событий.
Фьючерсы и задачи
Futures
Future представляет вычисления, которые либо уже выполняются, либо будут запланированы в будущем. Это специальный низкоуровневый ожидаемый объект, который представляет собой конечный результат асинхронной операции.
Не путайте threading.Future и asyncio.Future . Первый является частью модуля потоков и для него не определен метод __iter__() .
asyncio.Future является awaitable и может быть использован с оператором yield from . В общем случае вам не придется иметь дело с фьючерсами напрямую. Обычно они предоставляются библиотеками или API asyncio.
Для наглядности мы покажем пример, в котором создается фьючерс, ожидаемый корутиной. Изучите приведенный ниже фрагмент:
import asyncio from asyncio import Future async def bar(future): print("bar will sleep for 3 seconds") await asyncio.sleep(3) print("bar resolving the future") future.done() future.set_result("future is resolved") async def foo(future): print("foo will await the future") await future print("foo finds the future resolved") async def main(): future = Future() results = await asyncio.gather(foo(future), bar(future)) if __name__ == "__main__": asyncio.run(main()) print("main exiting")
Обе корутины получают в качестве аргумента объект future . Корутина foo() ожидает разрешения future , в то время как корутина bar() разрешает future через три секунды.
Задачи
Задачи (tasks) похожи на фьючерсы. Фактически, Task является подклассом Future и может быть создан с помощью следующих методов:
- asyncio.create_task() появился в Python 3.7 и является предпочтительным способом создания задач. Метод принимает корутины и оборачивает их как задачи.
- loop.create_task() принимает только корутины.
- asyncio.ensure_future() принимает futures, coroutines и любые ожидающие объекты.
Задачи оборачивают корутины и запускают их в циклах событий. Если корутина ожидает Future, задача приостанавливает выполнение корутины и ждет завершения Future. После завершения Future выполнение обернутой программы возобновляется.
В циклах событий используется кооперативное планирование, то есть цикл событий выполняет одну задачу за раз. Пока задача ожидает завершения Future, цикл событий запускает другие задачи, обратные вызовы или выполняет операции ввода-вывода. Задачи также могут быть отменены.
Мы переписали пример фьючерса с использованием задач следующим образом:
import asyncio from asyncio import Future async def bar(future): print("bar will sleep for 3 seconds") await asyncio.sleep(3) print("bar resolving the future") future.done() future.set_result("future is resolved") async def foo(future): print("foo will await the future") await future print("foo finds the future resolved") async def main(): future = Future() loop = asyncio.get_event_loop() t1 = loop.create_task(bar(future)) t2 = loop.create_task(foo(future)) await t2, t1 if __name__ == "__main__": loop = asyncio.get_event_loop() loop.run_until_complete(main()) print("main exiting")
Цепочка корутин: старый подход против нового
Старый способ создания цепочки корутин
Одно из наиболее распространенных применений корутин — их объединение в цепочки для конвейерной обработки данных. Вы можете выстроить цепочку корутин, подобно тому, как вы передаете команды Unix в shell.
Идея заключается в том, что входные данные проходят через первую корутину, которая может выполнить некоторые действия с входными данными, а затем передает измененные данные второй корутине, которая может выполнить дополнительные операции с входными данными.
Входные данные проходят через цепочку корутин, каждая из которых выполняет некоторую операцию над входными данными, пока входные данные не достигнут последней корутины, где они возвращаются обратно вызывающей стороне (оригинальному вызывающему коду).
Рассмотрим следующий пример, в котором вычисляются значения выражения x**2 + 3 для первых ста натуральных чисел.
Вы вручную работаете с конвейером данных, используя метод next() , поэтому можете настроить цепочку, не заботясь об изменениях, необходимых для ее работы с циклом событий asyncio. Настройка осуществляется следующим образом:
- Первая корутина генерирует натуральные числа, начиная с 1.
- Вторая корутина вычисляет квадрат каждого переданного на вход числа.
- Последняя функция является генератором и прибавляет к переданному ей значению 3 и возвращает результат.
def coro3(k): yield (k + 3) def coro2(j): j = j * j yield from coro3(j) def coro1(): i = 0 while True: yield from coro2(i) i += 1 if __name__ == "__main__": # The first 100 natural numbers evaluated for the following expression # x^2 + 3 cr = coro1() for v in range(100): print("f() = ".format(v, next(cr)))
В приведенном примере конец цепочки состоит из генератора, однако эта цепочка не будет работать с циклом событий asyncio, поскольку он не работает с генераторами. Один из способов исправить это — заменить последний генератор на обычную функцию, возвращающую фьючерс с итоговым результатом.
Метод coro3() будет выглядеть так:
def coro3(k): f = Future() f.set_result(k + 3) f.done() return f
Еще один способ — добавить @asyncio.coroutine к coro3() и return вместо yield. Изменения будут выглядеть следующим образом:
@asyncio.coroutine def coro3(k): return k + 3
Важной оговоркой является то, что если бы мы использовали декоратор @types.coroutine , то программа завершилась бы неудачей. Это связано с тем, что @asyncio.coroutine может преобразовать обычную функцию в coroutine, а @types.coroutine — нет.
Обратите внимание, что в предыдущих примерах мы не декорировали функции coro1() и coro2() с помощью @asyncio.coroutine .
Обе функции являются генераторными coroutine-функциями, поскольку в их телах присутствуют yield from . Кроме того, появление декоратора не является строго обязательным, но если использовать декораторы, то программа все равно будет работать корректно.
Новый способ организации цепочки нативных корутин
Аналогично генераторам и корутинам, основанным на генераторах, мы также можем объединять в цепочки нативные корутины.
import asyncio async def coro3(k): return k + 3 async def coro2(j): j = j * j res = await coro3(j) return res async def coro1(): i = 0 while i < 100: res = await coro2(i) print("f() = ".format(i, res)) i += 1 if __name__ == "__main__": # The first 100 natural numbers evaluated for the following expression # x^2 + 3 cr = coro1() loop = asyncio.get_event_loop() loop.run_until_complete(cr)
Применение asyncio на практике
Задача состоит в том, чтобы реализовать собственную корутину, выполняющую асинхронный сон. Сигнатура корутины выглядит следующим образом.
# Implement the following coroutine where # sleep_for is defined in seconds async def asleep(sleep_for): pass
Решение
Первая мысль, которая приходит в голову, — использовать API time.sleep() для ожидания запрошенного времени сна. Однако этот API является блокирующим и будет блокировать выполняющий его поток. Очевидно, что это исключает обращение к API из главного потока. Но это не исключает возможности выполнения этого API в другом потоке.
Это наводит нас на мысль о возможном решении. Мы можем создать объект Future и ожидать его в корутине asleep() . Единственное требование теперь — чтобы другой поток разрешал фьючерсы по истечении sleep_for секунд. Частичное решение выглядит следующим образом:
async def asleep(sleep_for): future = Future() Thread(target=sync_sleep, args=(sleep_for, future)).start() await future def sync_sleep(sleep_for, future): # sleep synchronously time.sleep(sleep_for) # resolve the future future.set_result(None
Добавим остальное и посмотрим, что получится:
from threading import Thread from threading import current_thread from asyncio import Future import asyncio import time async def asleep(sleep_for): future = Future() Thread(target=sync_sleep, args=(sleep_for, future)).start() await future def sync_sleep(sleep_for, future): # sleep synchronously time.sleep(sleep_for) # resolve the future future.set_result(None) print("Sleeping completed in ".format(current_thread().getName()), flush=True) if __name__ == "__main__": start = time.time() work = list() work.append(asleep(1)) loop = asyncio.get_event_loop() loop.run_until_complete(asyncio.wait(work, return_when=asyncio.ALL_COMPLETED)) print("main program exiting after running for ".format(time.time() - start))
Удивительно, но приведенная выше программа зависает и не завершается, хотя сообщение из метода sync_sleep() выводится. Почему-то корутина asleep() так и не возобновляется после завершения ожидаемого ею фьючерса.
Причина в том, что Future не является потокобезопасным. К счастью, asyncio предоставляет метод для потокобезопасного выполнения корутины в заданном цикле. В качестве API используется run_coroutine_threadsafe() .
Итак, у нас есть способ потокобезопасного выполнения фьючерса, однако делать это нужно в другой корутине, поскольку API run_coroutine_threadsafe() принимает только корутины. Для этого необходимо немного модифицировать метод sync_sleep() следующим образом:
def sync_sleep(sleep_for, future, loop): # sleep synchronously time.sleep(sleep_for) # define a nested coroutine to resolve the future async def sleep_future_resolver(): # resolve the future future.set_result(None) asyncio.run_coroutine_threadsafe(sleep_future_resolver(), loop)
Мы определяем вложенную корутину sleep_future_resolver , которая разрешает объект Future .
Также отметим, что sync_sleepnow принимает в качестве параметра цикл событий. Это должен быть тот же цикл событий, который изначально выполнял корутину asleep() .
Изменения в корутине asleep() показаны ниже:
async def asleep(sleep_for): future = Future() # get the current event loop current_loop = asyncio.get_running_loop() Thread(target=sync_sleep, args=(sleep_for, future, current_loop)).start() await future
Вот что мы имеем в итоге:
from threading import Thread from threading import current_thread from asyncio import Future import asyncio import time async def asleep(sleep_for): future = Future() current_loop = asyncio.get_event_loop() Thread(target=sync_sleep, args=(sleep_for, future, current_loop)).start() await future def sync_sleep(sleep_for, future, loop): # sleep synchronously time.sleep(sleep_for) # define a nested coroutine to resolve the future async def sleep_future_resolver(): # resolve the future future.set_result(None) asyncio.run_coroutine_threadsafe(sleep_future_resolver(), loop) print("Sleeping completed in \n".format(current_thread().getName()), flush=True) if __name__ == "__main__": start = time.time() work = list() work.append(asleep(5)) work.append(asleep(5)) work.append(asleep(5)) work.append(asleep(5)) work.append(asleep(5)) loop = asyncio.get_event_loop() loop.run_until_complete(asyncio.wait(work, return_when=asyncio.ALL_COMPLETED)) print("main program exiting after running for ".format(time.time() - start))
Вывод показывает, что сон происходит в порожденных нами потоках, а не в главном потоке. Более того, несмотря на то, что мы пять раз вызываем корутину asleep() для сна на пять секунд каждый раз, общее время выполнения программы составляет примерно пять секунд, как и должно быть, если мы правильно реализовали решение.
В качестве упражнения подумайте, что произойдет, если мы создадим пять потоков и в каждом из них вызовем time.sleep() . В этом случае программа выполнится за пять или двадцать пять секунд? Попробуйте это сделать и понаблюдайте за временем выполнения программы.
from threading import Thread from threading import current_thread import time def sync_sleep(sleep_for): time.sleep(sleep_for) print("Sleeping completed in ".format(current_thread().getName())) if __name__ == "__main__": start = time.time() threads = list() for _ in range(0, 5): threads.append(Thread(target=sync_sleep, args=(5,))) for thread in threads: thread.start() for thread in threads: thread.join() print("main program exiting after running for ".format(time.time() - start))
Тест синхронного сна по-прежнему занимает пять секунд! Вы можете задаться вопросом, в чем разница между нашими программами асинхронного и синхронного сна? Ответ заключается в том, что вызов асинхронного сна является неблокирующим, в отличие от вызова синхронного.
Однако внутри планировщик, увидев, что поток собирается блокировать вызов сна в течение пяти секунд, переключает его на другой поток и возобновляет выполнение только по истечении не менее пяти секунд.
Асинхронный Python: различные формы конкурентности

Это перевод статьи Абу Ашраф Маснуна «Async Python: The Different Forms of Concurrency».
С появлением Python 3 довольно много шума об «асинхронности» и «параллелизме», можно полагать, что Python недавно представил эти возможности/концепции. Но это не так. Мы много раз использовали эти операции. Кроме того, новички могут подумать, что asyncio является единственным или лучшим способом воссоздать и использовать асинхронные/параллельные операции. В этой статье мы рассмотрим различные способы достижения параллелизма, их преимущества и недостатки.
Определение терминов:
Прежде чем мы углубимся в технические аспекты, важно иметь некоторое базовое понимание терминов, часто используемых в этом контексте.
Синхронный и асинхронный:
В синхронных операциях задачи выполняются друг за другом. В асинхронных — задачи могут запускаться и завершаться независимо друг от друга. Одна асинхронная задача может запускаться и продолжать выполняться, пока выполнение переходит к новой задаче. Асинхронные задачи не блокируют (не заставляют ждать завершения выполнения задачи) операции и обычно выполняются в фоновом режиме.
Например, вы должны обратиться в туристическое агентство, чтобы спланировать свой следующий отпуск. Вам нужно отправить письмо своему руководителю, прежде чем улететь. В синхронном режиме, вы сначала позвоните в туристическое агентство, и если вас попросят подождать, то вы будете ждать, пока вам не ответят. Затем вы начнёте писать письмо руководителю. Таким образом, вы выполняете задачи последовательно, одна за другой.
Но, если вы умны, то пока вас попросили подождать, вы начнёте писать письмо, и когда с вами снова заговорят, вы приостановите написание, поговорите, а затем допишете письмо. Вы также можете попросить друга позвонить в агентство, а сами написать письмо. Это асинхронность, задачи не блокируют друг друга.
Конкурентность и параллелизм:
Конкурентность подразумевает, что две задачи выполняются совместно. В нашем предыдущем примере, когда мы рассматривали асинхронный пример, мы постепенно продвигались то в написании письма, то в разговоре с турагентством. Это конкурентность.
Когда мы попросили позвонить друга, а сами писали письмо, то задачи выполнялись параллельно.
Параллелизм по сути является формой конкурентности. Но параллелизм зависит от оборудования. Например, если в CPU только одно ядро, то две задачи не могут выполняться параллельно. Они просто делят процессорное время между собой. Тогда это конкурентность, но не параллелизм. Но когда у нас есть несколько ядер, мы можем выполнять несколько операций (в зависимости от количества ядер) одновременно.
- Синхронность: блокирует операции (блокирующие)
- Асинхронность: не блокирует операции (неблокирующие)
- Конкурентность: совместный прогресс (совместные)
- Параллелизм: параллельный прогресс (параллельные)
Параллелизм подразумевает конкурентность. Но конкурентность не всегда подразумевает параллелизм.
Потоки и процессы
Python поддерживает потоки уже очень давно. Потоки позволяют выполнять операции конкурентно. Но есть проблема, связанная с Global Interpreter Lock (GIL) из-за которой потоки не могли обеспечить настоящий параллелизм. И тем не менее, с появлением multiprocessing можно использовать несколько ядер с помощью Python.
Потоки (Threads)
Рассмотрим небольшой пример. В нижеследующем коде функция worker будет выполняться в нескольких потоках асинхронно и одновременно.
import threading import time import random def worker(number): sleep = random.randrange(1, 10) time.sleep(sleep) print("I am Worker <>, I slept for <> seconds".format(number, sleep)) for i in range(5): t = threading.Thread(target=worker, args=(i,)) t.start() print("All Threads are queued, let's see when they finish!")
А вот пример выходных данных:
$ python thread_test.py All Threads are queued, let's see when they finish! I am Worker 1, I slept for 1 seconds I am Worker 3, I slept for 4 seconds I am Worker 4, I slept for 5 seconds I am Worker 2, I slept for 7 seconds I am Worker 0, I slept for 9 seconds
Таким образом мы запустили 5 потоков для совместной работы, и после их старта (т.е. после запуска функции worker) операция не ждёт завершения работы потоков прежде чем перейти к следующему оператору print. Это асинхронная операция.
В нашем примере мы передали функцию в конструктор Thread. Если бы мы хотели, то могли бы реализовать подкласс с методом (ООП стиль).
Global Interpreter Lock (GIL)
GIL нужен, чтобы сделать обработку памяти CPython проще и обеспечить наилучшую интеграцию с C.
GIL — это механизм блокировки, когда интерпретатор Python запускает в работу только один поток за раз. Это значит, только один поток может исполняться в байт-коде Python единовременно. GIL следит за тем, чтобы несколько потоков не выполнялись параллельно.
Краткие сведения о GIL:
- Одновременно может выполняться один поток.
- Интерпретатор Python переключается между потоками для достижения конкурентности.
- GIL применим к CPython (стандартной реализации). Но, например, Jython и IronPython не имеют GIL.
- GIL делает однопоточные программы быстрыми.
- Операциям ввода/вывода GIL обычно не мешает.
- GIL позволяет легко интегрировать непотокобезопасные библиотеки на C, благодаря GIL у нас есть много высокопроизводительных расширений/модулей, написанных на C.
- Для CPU-зависимых задач интерпретатор делает проверку каждые N тиков и переключает потоки. Таким образом один поток не блокирует другие.
Многие видят в GIL слабость. Я же рассматриваю это как благо, ведь были созданы такие библиотеки, как NumPy, SciPy, которые занимают особое, уникальное положение в научном обществе.
Процессы (Processes)
Чтобы достичь параллелизма, в Python был добавлен модуль multiprocessing, который предоставляет API и выглядит очень похожим, если вы использовали threading раньше.
Давайте просто пойдем и изменим предыдущий пример. Теперь модифицированная версия использует Процесс вместо Потока.
import multiprocessing import time import random def worker(number): sleep = random.randrange(1, 10) time.sleep(sleep) print("I am Worker <>, I slept for <> seconds".format(number, sleep)) for i in range(5): t = multiprocessing.Process(target=worker, args=(i,)) t.start() print("All Processes are queued, let's see when they finish!")
Что же изменилось? Я просто импортировал модуль multiprocessing вместо threading. А затем, вместо потока я использовал процесс. Вот и всё! Теперь вместо множества потоков мы используем процессы, которые запускаются на разных ядрах CPU (если, конечно, у вашего процессора несколько ядер).
С помощью класса Pool мы также можем распределить выполнение одной функции между несколькими процессами для разных входных значений.
Пример из официальных документов:
from multiprocessing import Pool def f(x): return x*x if __name__ == '__main__': p = Pool(5) print(p.map(f, [1, 2, 3]))
Здесь вместо того, чтобы перебирать список значений и вызывать функцию f по одному, мы фактически запускаем функцию в разных процессах.
Один процесс выполняет f(1), другой-f(2), а другой-f (3). Наконец, результаты снова объединяются в список. Это позволяет нам разбить тяжелые вычисления на более мелкие части и запускать их параллельно для более быстрого расчета.
Модуль concurrent.futures
Модуль concurrent.futures большой и позволяет писать асинхронный код очень легко. Мои любимчики — ThreadPoolExecutor и ProcessPoolExecutor. Эти исполнители поддерживают пул потоков или процессов. Мы отправляем наши задачи в пул, и он запускает задачи в доступном потоке / процессе. Возвращается объект Future, который можно использовать для запроса и получения результата по завершении задачи.
А вот пример ThreadPoolExecutor:
from concurrent.futures import ThreadPoolExecutor from time import sleep def return_after_5_secs(message): sleep(5) return message pool = ThreadPoolExecutor(3) future = pool.submit(return_after_5_secs, ("hello")) print(future.done()) sleep(5) print(future.done()) print(future.result())
Asyncio — что, как и почему
У вас, вероятно, есть вопрос, который есть у многих людей в сообществе Python — что asyncio приносит нового? Зачем нужен был еще один способ асинхронного ввода-вывода? Разве у нас уже не было потоков и процессов?
Зачем нам нужен asyncio?
Процессы очень дорогостоящие и требуют много ресурсов для создания. Поэтому для операций ввода/вывода в основном выбираются потоки.
Мы знаем, что ввод-вывод зависит от внешних вещей — медленные диски или неприятные сетевые лаги делают ввод-вывод часто непредсказуемым. Теперь предположим, что мы используем потоки для операций ввода-вывода. 3 потока выполняют различные задачи ввода-вывода. Интерпретатор должен был бы переключаться между конкурентными потоками и давать каждому из них некоторое время по очереди.
Назовем потоки — T1, T2 и T3. Три потока начали свою операцию ввода-вывода. T3 завершает его первым. T2 и T1 все еще ожидают ввода-вывода. Интерпретатор Python переключается на T1, но он все еще ждет. Хорошо, интерпретатор перемещается в T2, а тот все еще ждет, а затем перемещается в T3, который готов и выполняет код. Вы видите в этом проблему?
T3 был готов, но интерпретатор сначала переключился между T2 и T1 — это понесло расходы на переключение, которых мы могли бы избежать, если бы интерпретатор сначала переключился на T3, верно?
Что есть asynio?
Asyncio предоставляет нам цикл событий наряду с другими крутыми вещами. Цикл событий (event loop) отслеживает события ввода/вывода и переключает задачи, которые готовы и ждут операции ввода/вывода.
Идея очень проста. Есть цикл обработки событий. И у нас есть функции, которые выполняют асинхронные операции ввода-вывода. Мы передаем свои функции циклу событий и просим его запустить их. Цикл событий возвращает нам объект Future, словно обещание, что в будущем мы что-то получим. Мы держимся за обещание, время от времени проверяем, имеет ли оно значение, и, наконец, когда значение получено, мы используем его в некоторых других операциях.
Как использовать asyncio?
Прежде чем мы начнём, давайте взглянем на пример:
import asyncio import datetime import random async def my_sleep_func(): await asyncio.sleep(random.randint(0, 5)) async def display_date(num, loop): end_time = loop.time() + 50.0 while True: print("Loop: <> Time: <>".format(num, datetime.datetime.now())) if (loop.time() + 1.0) >= end_time: break await my_sleep_func() loop = asyncio.get_event_loop() asyncio.ensure_future(display_date(1, loop)) asyncio.ensure_future(display_date(2, loop)) loop.run_forever()
Обратите внимание, что синтаксис async/await предназначен только для Python 3.5 и выше. Пройдёмся по коду:
- У нас есть асинхронная функция display_date, которая принимает число (в качестве идентификатора) и цикл обработки событий в качестве параметров.
- Функция имеет бесконечный цикл, который прерывается через 50 секунд. Но за этот период она неоднократно печатает время и делает паузу. Функция await может ожидать завершения выполнения других асинхронных функций (корутин).
- Передаем функцию в цикл обработки событий (используя метод ensure_future).
- Запускаем цикл событий.
Всякий раз, когда происходит вызов await, asyncio понимает, что функции, вероятно, потребуется некоторое время. Таким образом, он приостанавливает выполнение, начинает мониторинг любого связанного с ним события ввода-вывода и позволяет запускать задачи. Когда asyncio замечает, что приостановленный ввод-вывод функции готов, он возобновляет функцию.
Делаем правильный выбор
Только что мы прошлись по самым популярным формам конкурентности. Но остаётся вопрос — что следует выбрать?
Это зависит от вариантов использования. Из моего опыта я склонен следовать этому псевдо-коду:
if io_bound: if io_very_slow: print("Use Asyncio") else: print("Use Threads") else: print("Multi Processing")
Такие сложные материи, как асинхронность, мы проходим на обучении Рython
AsyncIO для практикующего python-разработчика
Я помню тот момент, когда подумал «Как же медленно всё работает, что если я распараллелю вызовы?», а спустя 3 дня, взглянув на код, ничего не мог понять в жуткой каше из потоков, синхронизаторов и функций обратного вызова.
Тогда я познакомился с asyncio, и всё изменилось.
Если кто не знает, asyncio — новый модуль для организации конкурентного программирования, который появился в Python 3.4. Он предназначен для упрощения использования корутин и футур в асинхронном коде — чтобы код выглядел как синхронный, без коллбэков.
Я помню, в то время было несколько похожих инструментов, и один из них выделялся — это библиотека gevent. Я советую всем прочитать прекрасное руководство gevent для практикующего python-разработчика, в котором описана не только работа с ней, но и что такое конкурентность в общем понимании. Мне настолько понравилось та статья, что я решил использовать её как шаблон для написания введения в asyncio.
Небольшой дисклеймер — это статья не gevent vs asyncio. Nathan Road уже сделал это за меня в своей заметке. Все примеры вы можете найти на GitHub.
Я знаю, вам уже не терпится писать код, но для начала я бы хотел рассмотреть несколько концепций, которые нам пригодятся в дальнейшем.
Потоки, циклы событий, корутины и футуры
Потоки — наиболее распространённый инструмент. Думаю, вы слышали о нём и ранее, однако asyncio оперирует несколько другими понятиями: циклы событий, корутины и футуры.
- цикл событий (event loop) по большей части всего лишь управляет выполнением различных задач: регистрирует поступление и запускает в подходящий момент
- корутины — специальные функции, похожие на генераторы python, от которых ожидают (await), что они будут отдавать управление обратно в цикл событий. Необходимо, чтобы они были запущены именно через цикл событий
- футуры — объекты, в которых хранится текущий результат выполнения какой-либо задачи. Это может быть информация о том, что задача ещё не обработана или уже полученный результат; а может быть вообще исключение
Синхронное и асинхронное выполнение
В видео "Конкурентность — это не параллелизм, это лучше" Роб Пайк обращает ваше внимание на ключевую вещь. Разбиение задач на конкурентные подзадачи возможно только при таком параллелизме, когда он же и управляет этими подзадачами.
Asyncio делает тоже самое — вы можете разбивать ваш код на процедуры, которые определять как корутины, что даёт возможность управлять ими как пожелаете, включая и одновременное выполнение. Корутины содержат операторы yield, с помощью которых мы определяем места, где можно переключиться на другие ожидающие выполнения задачи.
За переключение контекста в asyncio отвечает yield, который передаёт управление обратно в event loop, а тот в свою очередь — к другой корутине. Рассмотрим базовый пример:
import asyncio async def foo(): print('Running in foo') await asyncio.sleep(0) print('Explicit context switch to foo again') async def bar(): print('Explicit context to bar') await asyncio.sleep(0) print('Implicit context switch back to bar') ioloop = asyncio.get_event_loop() tasks = [ioloop.create_task(foo()), ioloop.create_task(bar())] wait_tasks = asyncio.wait(tasks) ioloop.run_until_complete(wait_tasks) ioloop.close()
$ python3 1-sync-async-execution-asyncio-await.py Running in foo Explicit context to bar Explicit context switch to foo again Implicit context switch back to bar
* Сначала мы объявили пару простейших корутин, которые притворяются неблокирующими, используя sleep из asyncio
* Корутины могут быть запущены только из другой корутины, или обёрнуты в задачу с помощью create_task
* После того, как у нас оказались 2 задачи, объединим их, используя wait
* И, наконец, отправим на выполнение в цикл событий через run_until_complete
Используя await в какой-либо корутине, мы таким образом объявляем, что корутина может отдавать управление обратно в event loop, который, в свою очередь, запустит какую-либо следующую задачу: bar. В bar произойдёт тоже самое: на await asyncio.sleep управление будет передано обратно в цикл событий, который в нужное время вернётся к выполнению foo.
Представим 2 блокирующие задачи: gr1 и gr2, как будто они обращаются к неким сторонним сервисам, и, пока они ждут ответа, третья функция может работать асинхронно.
import time import asyncio start = time.time() def tic(): return 'at %1.1f seconds' % (time.time() - start) async def gr1(): # Busy waits for a second, but we don't want to stick around. print('gr1 started work: <>'.format(tic())) await asyncio.sleep(2) print('gr1 ended work: <>'.format(tic())) async def gr2(): # Busy waits for a second, but we don't want to stick around. print('gr2 started work: <>'.format(tic())) await asyncio.sleep(2) print('gr2 Ended work: <>'.format(tic())) async def gr3(): print("Let's do some stuff while the coroutines are blocked, <>".format(tic())) await asyncio.sleep(1) print("Done!") ioloop = asyncio.get_event_loop() tasks = [ ioloop.create_task(gr1()), ioloop.create_task(gr2()), ioloop.create_task(gr3()) ] ioloop.run_until_complete(asyncio.wait(tasks)) ioloop.close()
$ python3 1b-cooperatively-scheduled-asyncio-await.py gr1 started work: at 0.0 seconds gr2 started work: at 0.0 seconds Lets do some stuff while the coroutines are blocked, at 0.0 seconds Done! gr1 ended work: at 2.0 seconds gr2 Ended work: at 2.0 seconds
Обратите внимание как происходит работа с вводом-выводом и планированием выполнения, позволяя всё это уместить в один поток. Пока две задачи заблокированы ожиданием I/O, третья функция может занимать всё процессорное время.
Порядок выполнения
В синхронном мире мы мыслим последовательно. Если у нас есть список задач, выполнение которых занимает разное время, то они завершатся в том же порядке, в котором поступили в обработку. Однако, в случае конкурентности нельзя быть в этом уверенным.
import random from time import sleep import asyncio def task(pid): """Synchronous non-deterministic task. """ sleep(random.randint(0, 2) * 0.001) print('Task %s done' % pid) async def task_coro(pid): """Coroutine non-deterministic task """ await asyncio.sleep(random.randint(0, 2) * 0.001) print('Task %s done' % pid) def synchronous(): for i in range(1, 10): task(i) async def asynchronous(): tasks = [asyncio.ensure_future(task_coro(i)) for i in range(1, 10)] await asyncio.wait(tasks) print('Synchronous:') synchronous() ioloop = asyncio.get_event_loop() print('Asynchronous:') ioloop.run_until_complete(asynchronous()) ioloop.close()
$ python3 1c-determinism-sync-async-asyncio-await.py Synchronous: Task 1 done Task 2 done Task 3 done Task 4 done Task 5 done Task 6 done Task 7 done Task 8 done Task 9 done Asynchronous: Task 2 done Task 5 done Task 6 done Task 8 done Task 9 done Task 1 done Task 4 done Task 3 done Task 7 done
Разумеется, ваш результат будет иным, поскольку каждая задача будет засыпать на случайное время, но заметьте, что результат выполнения полностью отличается, хотя мы всегда ставим задачи в одном и том же порядке.
Также обратите внимание на корутину для нашей довольно простой задачи. Это важно для понимания, что в asyncio нет никакой магии при реализации неблокирующих задач. Во время реализации asyncio стоял отдельно в стандартной библиотеке, т.к. остальные модули предоставляли только блокирующую функциональность. Вы можете использовать модуль concurrent.futures для оборачивания блокирующих задач в потоки или процессы и получения футуры для использования в asyncio. Несколько таких примеров доступны на GitHub.
Это, наверно, главный недостаток сейчас при использовании asyncio, однако уже есть несколько библиотек, помогающих решить эту проблему.
Самая популярная блокирующая задача — получение данных по HTTP-запросу. Рассмотрим работу с великолепной библиотекой aiohttp на примере получения информации о публичных событиях на GitHub.
import time import urllib.request import asyncio import aiohttp URL = 'https://api.github.com/events' MAX_CLIENTS = 3 def fetch_sync(pid): print('Fetch sync process <> started'.format(pid)) start = time.time() response = urllib.request.urlopen(URL) datetime = response.getheader('Date') print('Process <>: <>, took: seconds'.format( pid, datetime, time.time() - start)) return datetime async def fetch_async(pid): print('Fetch async process <> started'.format(pid)) start = time.time() response = await aiohttp.request('GET', URL) datetime = response.headers.get('Date') print('Process <>: <>, took: seconds'.format( pid, datetime, time.time() - start)) response.close() return datetime def synchronous(): start = time.time() for i in range(1, MAX_CLIENTS + 1): fetch_sync(i) print("Process took: seconds".format(time.time() - start)) async def asynchronous(): start = time.time() tasks = [asyncio.ensure_future( fetch_async(i)) for i in range(1, MAX_CLIENTS + 1)] await asyncio.wait(tasks) print("Process took: seconds".format(time.time() - start)) print('Synchronous:') synchronous() print('Asynchronous:') ioloop = asyncio.get_event_loop() ioloop.run_until_complete(asynchronous()) ioloop.close()
$ python3 1d-async-fetch-from-server-asyncio-await.py Synchronous: Fetch sync process 1 started Process 1: Wed, 17 Feb 2016 13:10:11 GMT, took: 0.54 seconds Fetch sync process 2 started Process 2: Wed, 17 Feb 2016 13:10:11 GMT, took: 0.50 seconds Fetch sync process 3 started Process 3: Wed, 17 Feb 2016 13:10:12 GMT, took: 0.48 seconds Process took: 1.54 seconds Asynchronous: Fetch async process 1 started Fetch async process 2 started Fetch async process 3 started Process 3: Wed, 17 Feb 2016 13:10:12 GMT, took: 0.50 seconds Process 2: Wed, 17 Feb 2016 13:10:12 GMT, took: 0.52 seconds Process 1: Wed, 17 Feb 2016 13:10:12 GMT, took: 0.54 seconds Process took: 0.54 seconds
Тут стоит обратить внимание на пару моментов.
Во-первых, разница во времени — при использовании асинхронных вызовов мы запускаем запросы одновременно. Как говорилось ранее, каждый из них передавал управление следующему и возвращал результат по завершении. То есть скорость выполнения напрямую зависит от времени работы самого медленного запроса, который занял как раз 0.54 секунды. Круто, правда?
Во-вторых, насколько код похож на синхронный. Это же по сути одно и то же! Основные отличия связаны с реализацией библиотеки для выполнения запросов, созданием и ожиданием завершения задач.
Создание конкурентности
До сих пор мы использовали единственный метод создания и получения результатов из корутин, создания набора задач и ожидания их завершения. Однако, корутины могут быть запланированы для запуска и получения результатов несколькими способами. Представьте ситуацию, когда нам надо обрабатывать результаты GET-запросов по мере их получения; на самом деле реализация очень похожа на предыдущую:
import time import random import asyncio import aiohttp URL = 'https://api.github.com/events' MAX_CLIENTS = 3 async def fetch_async(pid): start = time.time() sleepy_time = random.randint(2, 5) print('Fetch async process <> started, sleeping for <> seconds'.format( pid, sleepy_time)) await asyncio.sleep(sleepy_time) response = await aiohttp.request('GET', URL) datetime = response.headers.get('Date') response.close() return 'Process <>: <>, took: seconds'.format( pid, datetime, time.time() - start) async def asynchronous(): start = time.time() futures = [fetch_async(i) for i in range(1, MAX_CLIENTS + 1)] for i, future in enumerate(asyncio.as_completed(futures)): result = await future print('<> <>'.format(">>" * (i + 1), result)) print("Process took: seconds".format(time.time() - start)) ioloop = asyncio.get_event_loop() ioloop.run_until_complete(asynchronous()) ioloop.close()
$ python3 2a-async-fetch-from-server-as-completed-asyncio-await.py Fetch async process 1 started, sleeping for 4 seconds Fetch async process 3 started, sleeping for 5 seconds Fetch async process 2 started, sleeping for 3 seconds >> Process 2: Wed, 17 Feb 2016 13:55:19 GMT, took: 3.53 seconds >>>> Process 1: Wed, 17 Feb 2016 13:55:20 GMT, took: 4.49 seconds >>>>>> Process 3: Wed, 17 Feb 2016 13:55:21 GMT, took: 5.48 seconds Process took: 5.48 seconds
Посмотрите на отступы и тайминги — мы запустили все задачи одновременно, однако они обработаны в порядке завершения выполнения. Код в данном случае немного отличается: мы пакуем корутины, каждая из которых уже подготовлена для выполнения, в список. Функция as_completed возвращает итератор, который выдаёт результаты корутин по мере их выполнения. Круто же, правда?! Кстати, и as_completed, и wait — функции из пакета concurrent.futures.
Ещё один пример — что если вы хотите узнать свой IP адрес. Есть куча сервисов для этого, но вы не знаете какой из них будет доступен в момент работы программы. Вместо того, чтобы последовательно опрашивать каждый из списка, можно запустить все запросы конкурентно и выбрать первый успешный.
Что ж, для этого в нашей любимой функции wait есть специальный параметр return_when. До сих пор мы игнорировали то, что возвращает wait, т.к. только распараллеливали задачи. Но теперь нам надо получить результат из корутины, так что будем использовать набор футур done и pending.
from collections import namedtuple import time import asyncio from concurrent.futures import FIRST_COMPLETED import aiohttp Service = namedtuple('Service', ('name', 'url', 'ip_attr')) SERVICES = ( Service('ipify', 'https://api.ipify.org?format=json', 'ip'), Service('ip-api', 'http://ip-api.com/json', 'query') ) async def fetch_ip(service): start = time.time() print('Fetching IP from <>'.format(service.name)) response = await aiohttp.request('GET', service.url) json_response = await response.json() ip = json_response[service.ip_attr] response.close() return '<> finished with result: <>, took: seconds'.format( service.name, ip, time.time() - start) async def asynchronous(): futures = [fetch_ip(service) for service in SERVICES] done, pending = await asyncio.wait( futures, return_when=FIRST_COMPLETED) print(done.pop().result()) ioloop = asyncio.get_event_loop() ioloop.run_until_complete(asynchronous()) ioloop.close()
$ python3 2c-fetch-first-ip-address-response-await.py Fetching IP from ip-api Fetching IP from ipify ip-api finished with result: 82.34.76.170, took: 0.09 seconds Unclosed client session client_session: Task was destroyed but it is pending! task: wait_for=>
Что же случилось? Первый сервис ответил успешно, но в логах какое-то предупреждение!
На самом деле мы запустили выполнение двух задач, но вышли из цикла уже после первого результата, в то время как вторая корутина ещё выполнялась. Asyncio подумал что это баг и предупредил нас. Наверно, стоит прибираться за собой и явно убивать ненужные задачи. Как? Рад, что вы спросили.
Состояния футур
- ожидание (pending)
- выполнение (running)
- выполнено (done)
- отменено (cancelled)
Вы можете узнать состояние футуры с помощью методов done, cancelled или running, но не забывайте, что в случае done вызов result может вернуть как ожидаемый результат, так и исключение, которое возникло в процессе работы. Для отмены выполнения футуры есть метод cancel. Это подходит для исправления нашего примера.
from collections import namedtuple import time import asyncio from concurrent.futures import FIRST_COMPLETED import aiohttp Service = namedtuple('Service', ('name', 'url', 'ip_attr')) SERVICES = ( Service('ipify', 'https://api.ipify.org?format=json', 'ip'), Service('ip-api', 'http://ip-api.com/json', 'query') ) async def fetch_ip(service): start = time.time() print('Fetching IP from <>'.format(service.name)) response = await aiohttp.request('GET', service.url) json_response = await response.json() ip = json_response[service.ip_attr] response.close() return '<> finished with result: <>, took: seconds'.format( service.name, ip, time.time() - start) async def asynchronous(): futures = [fetch_ip(service) for service in SERVICES] done, pending = await asyncio.wait( futures, return_when=FIRST_COMPLETED) print(done.pop().result()) for future in pending: future.cancel() ioloop = asyncio.get_event_loop() ioloop.run_until_complete(asynchronous()) ioloop.close()
$ python3 2c-fetch-first-ip-address-response-no-warning-await.py Fetching IP from ipify Fetching IP from ip-api ip-api finished with result: 82.34.76.170, took: 0.08 seconds
Простой и аккуратный вывод — как раз то, что я люблю!
Если вам нужна некоторая дополнительная логика по обработке футур, то вы можете подключать коллбэки, которые будут вызваны при переходе в состояние done. Это может быть полезно для тестов, когда некоторые результаты надо переопределить какими-то своими значениями.
Обработка исключений
asyncio — это целиком про написание управляемого и читаемого конкурентного кода, что хорошо заметно при обработке исключений. Вернёмся к примеру, чтобы продемонстрировать.
Допустим, мы хотим убедиться, что все запросы к сервисам по определению IP вернули одинаковый результат. Однако, один из них может быть оффлайн и не ответить нам. Просто применим try. except как обычно:
from collections import namedtuple import time import asyncio import aiohttp Service = namedtuple('Service', ('name', 'url', 'ip_attr')) SERVICES = ( Service('ipify', 'https://api.ipify.org?format=json', 'ip'), Service('ip-api', 'http://ip-api.com/json', 'query'), Service('borken', 'http://no-way-this-is-going-to-work.com/json', 'ip') ) async def fetch_ip(service): start = time.time() print('Fetching IP from <>'.format(service.name)) try: response = await aiohttp.request('GET', service.url) except: return '<> is unresponsive'.format(service.name) json_response = await response.json() ip = json_response[service.ip_attr] response.close() return '<> finished with result: <>, took: seconds'.format( service.name, ip, time.time() - start) async def asynchronous(): futures = [fetch_ip(service) for service in SERVICES] done, _ = await asyncio.wait(futures) for future in done: print(future.result()) ioloop = asyncio.get_event_loop() ioloop.run_until_complete(asynchronous()) ioloop.close()
$ python3 3a-fetch-ip-addresses-fail-await.py Fetching IP from ip-api Fetching IP from borken Fetching IP from ipify ip-api finished with result: 85.133.69.250, took: 0.75 seconds ipify finished with result: 85.133.69.250, took: 1.37 seconds borken is unresponsive
Мы также можем обработать исключение, которое возникло в процессе выполнения корутины:
from collections import namedtuple import time import asyncio import aiohttp import traceback Service = namedtuple('Service', ('name', 'url', 'ip_attr')) SERVICES = ( Service('ipify', 'https://api.ipify.org?format=json', 'ip'), Service('ip-api', 'http://ip-api.com/json', 'this-is-not-an-attr'), Service('borken', 'http://no-way-this-is-going-to-work.com/json', 'ip') ) async def fetch_ip(service): start = time.time() print('Fetching IP from <>'.format(service.name)) try: response = await aiohttp.request('GET', service.url) except: return '<> is unresponsive'.format(service.name) json_response = await response.json() ip = json_response[service.ip_attr] response.close() return '<> finished with result: <>, took: seconds'.format( service.name, ip, time.time() - start) async def asynchronous(): futures = [fetch_ip(service) for service in SERVICES] done, _ = await asyncio.wait(futures) for future in done: try: print(future.result()) except: print("Unexpected error: <>".format(traceback.format_exc())) ioloop = asyncio.get_event_loop() ioloop.run_until_complete(asynchronous()) ioloop.close()
$ python3 3b-fetch-ip-addresses-future-exceptions-await.py Fetching IP from ipify Fetching IP from borken Fetching IP from ip-api ipify finished with result: 85.133.69.250, took: 0.91 seconds borken is unresponsive Unexpected error: Traceback (most recent call last): File “3b-fetch-ip-addresses-future-exceptions.py”, line 39, in asynchronous print(future.result()) File “3b-fetch-ip-addresses-future-exceptions.py”, line 26, in fetch_ip ip = json_response[service.ip_attr] KeyError: ‘this-is-not-an-attr’
Точно также, как и запуск задачи без ожидания её завершения является ошибкой, так и получение неизвестных исключений оставляет свои следы в выводе:
from collections import namedtuple import time import asyncio import aiohttp Service = namedtuple('Service', ('name', 'url', 'ip_attr')) SERVICES = ( Service('ipify', 'https://api.ipify.org?format=json', 'ip'), Service('ip-api', 'http://ip-api.com/json', 'this-is-not-an-attr'), Service('borken', 'http://no-way-this-is-going-to-work.com/json', 'ip') ) async def fetch_ip(service): start = time.time() print('Fetching IP from <>'.format(service.name)) try: response = await aiohttp.request('GET', service.url) except: print('<> is unresponsive'.format(service.name)) else: json_response = await response.json() ip = json_response[service.ip_attr] response.close() print('<> finished with result: <>, took: seconds'.format( service.name, ip, time.time() - start)) async def asynchronous(): futures = [fetch_ip(service) for service in SERVICES] await asyncio.wait(futures) # intentionally ignore results ioloop = asyncio.get_event_loop() ioloop.run_until_complete(asynchronous()) ioloop.close()
$ python3 3c-fetch-ip-addresses-ignore-exceptions-await.py Fetching IP from ipify Fetching IP from borken Fetching IP from ip-api borken is unresponsive ipify finished with result: 85.133.69.250, took: 0.78 seconds Task exception was never retrieved future: exception=KeyError(‘this-is-not-an-attr’,)> Traceback (most recent call last): File “3c-fetch-ip-addresses-ignore-exceptions.py”, line 25, in fetch_ip ip = json_response[service.ip_attr] KeyError: ‘this-is-not-an-attr’
Вывод выглядит также, как и в предыдущем примере за исключением укоризненного сообщения от asyncio.
Таймауты
А что, если информация о нашем IP не так уж важна? Это может быть хорошим дополнением к какому-то составному ответу, в котором эта часть будет опциональна. В таком случае не будем заставлять пользователя ждать. В идеале мы бы ставили таймаут на вычисление IP, после которого в любом случае отдавали ответ пользователю, даже без этой информации.
И снова у wait есть подходящий аргумент:
import time import random import asyncio import aiohttp import argparse from collections import namedtuple from concurrent.futures import FIRST_COMPLETED Service = namedtuple('Service', ('name', 'url', 'ip_attr')) SERVICES = ( Service('ipify', 'https://api.ipify.org?format=json', 'ip'), Service('ip-api', 'http://ip-api.com/json', 'query'), ) DEFAULT_TIMEOUT = 0.01 async def fetch_ip(service): start = time.time() print('Fetching IP from <>'.format(service.name)) await asyncio.sleep(random.randint(1, 3) * 0.1) try: response = await aiohttp.request('GET', service.url) except: return '<> is unresponsive'.format(service.name) json_response = await response.json() ip = json_response[service.ip_attr] response.close() print('<> finished with result: <>, took: seconds'.format( service.name, ip, time.time() - start)) return ip async def asynchronous(timeout): response = < "message": "Result from asynchronous.", "ip": "not available" >futures = [fetch_ip(service) for service in SERVICES] done, pending = await asyncio.wait( futures, timeout=timeout, return_when=FIRST_COMPLETED) for future in pending: future.cancel() for future in done: response["ip"] = future.result() print(response) parser = argparse.ArgumentParser() parser.add_argument( '-t', '--timeout', help='Timeout to use, defaults to <>'.format(DEFAULT_TIMEOUT), default=DEFAULT_TIMEOUT, type=float) args = parser.parse_args() print("Using a <> timeout".format(args.timeout)) ioloop = asyncio.get_event_loop() ioloop.run_until_complete(asynchronous(args.timeout)) ioloop.close()
Я также добавил аргумент timeout к строке запуска скрипта, чтобы проверить что же произойдёт, если запросы успеют обработаться. Также я добавил случайные задержки, чтобы скрипт не завершался слишком быстро, и было время разобраться как именно он работает.
$ python 4a-timeout-with-wait-kwarg-await.py Using a 0.01 timeout Fetching IP from ipify Fetching IP from ip-api
$ python 4a-timeout-with-wait-kwarg-await.py -t 5 Using a 5.0 timeout Fetching IP from ip-api Fetching IP from ipify ipify finished with result: 82.34.76.170, took: 1.24 seconds
Заключение
Asyncio укрепил мою и так уже большую любовь к python. Если честно, я влюбился в сопрограммы, ещё когда познакомился с ними в Tornado, но asyncio сумел взять всё лучшее из него и других библиотек по реализации конкурентности. Причём настолько, что были предприняты особые усилия, чтобы они могли использовать основной цикл ввода-вывода. Так что если вы используете Tornado или Twisted, то можете подключать код, предназначенный для asyncio!
Как я уже упоминал, основная проблема заключается в том, что стандартные библиотеки пока ещё не поддерживают неблокирующее поведение. Также и многие популярные библиотеки работают пока лишь в синхронном стиле, а те, что используют конкурентность, пока ещё молоды и экспериментальны. Однако, их число растёт.
Надеюсь, в этом уроке я показал, насколько приятно работать с asyncio, и эта технология подтолкнёт вас к переходу на python 3, если вы по какой-то причине застряли на python 2.7. Одно точно — будущее Python полностью изменилось.
От переводчика:
Оригинальная статья была опубликована 20 февраля 2016, за это время многое произошло. Вышел Python 3.6, в котором помимо оптимизаций была улучшена работа asyncio, API переведено в стабильное состояние. Были выпущены библиотеки для работы с Postgres, Redis, Elasticsearch и пр. в неблокирующем режиме. Даже новый фреймворк — Sanic, который напоминает Flask, но работает в асинхронном режиме. В конце концов даже event loop был оптимизирован и переписан на Cython, что получилось раза в 2 быстрее. Так что я не вижу причин игнорировать эту технологию!
