Asyncio, threading, and multiprocessing in python
Содержание:
Transport and Protocol with SSL¶
import asyncio
import ssl
def make_header():
head = b"HTTP/1.1 200 OK\r\n"
head += b"Content-Type: text/html\r\n"
head += b"\r\n"
return head
def make_body():
resp = b"<html>"
resp += b"<h1>Hello SSL</h1>"
resp += b"</html>"
return resp
sslctx = ssl.SSLContext(ssl.PROTOCOL_SSLv23)
sslctx.load_cert_chain(
certfile="./root-ca.crt", keyfile="./root-ca.key"
)
class Service(asyncio.Protocol):
def connection_made(self, tr):
self.tr = tr
self.total =
def data_received(self, data):
if data
resp = make_header()
resp += make_body()
self.tr.write(resp)
self.tr.close()
async def start():
server = await loop.create_server(
Service, "localhost", 4433, ssl=sslctx
)
await server.wait_closed()
try
loop = asyncio.get_event_loop()
loop.run_until_complete(start())
finally
loop.close()
output:
Coroutine objects
This section applies only to native coroutines with CO_COROUTINE
flag, i.e. defined with the new async def syntax.
The behavior of existing *generator-based coroutines* in asyncio
remains unchanged.
Great effort has been made to make sure that coroutines and
generators are treated as distinct concepts:
-
Native coroutine objects do not implement __iter__ and
__next__ methods. Therefore, they cannot be iterated over or
passed to iter(), list(), tuple() and other built-ins.
They also cannot be used in a for..in loop.An attempt to use __iter__ or __next__ on a native
coroutine object will result in a TypeError. -
Plain generators cannot yield from native coroutines:
doing so will result in a TypeError. -
generator-based coroutines (for asyncio code must be decorated
with @asyncio.coroutine) can yield from native coroutine
objects. -
inspect.isgenerator() and inspect.isgeneratorfunction()
return False for native coroutine objects and native
coroutine functions.
Debugging Features
A common beginner mistake is forgetting to use yield from on
coroutines:
@asyncio.coroutine
def useful():
asyncio.sleep(1) # this will do nothing without 'yield from'
For debugging this kind of mistakes there is a special debug mode in
asyncio, in which @coroutine decorator wraps all functions with a
special object with a destructor logging a warning. Whenever a wrapped
generator gets garbage collected, a detailed logging message is
generated with information about where exactly the decorator function
was defined, stack trace of where it was collected, etc. Wrapper
object also provides a convenient __repr__ function with detailed
information about the generator.
The only problem is how to enable these debug capabilities. Since
debug facilities should be a no-op in production mode, @coroutine
decorator makes the decision of whether to wrap or not to wrap based on
an OS environment variable PYTHONASYNCIODEBUG. This way it is
possible to run asyncio programs with asyncio’s own functions
instrumented. EventLoop.set_debug, a different debug facility, has
no impact on @coroutine decorator’s behavior.
Функция обратного вызова (callback)
В Python много библиотек для асинхронного программирования, наиболее популярными являются Tornado, Asyncio и Gevent. Давайте посмотрим, как работает Tornado. Он использует стиль обратного вызова (callbacks) для асинхронного сетевого ввода-вывода. Обратный вызов — это функция, которая означает: «Как только это будет сделано, выполните эту функцию». Другими словами, вы звоните в службу поддержки и оставляете свой номер, чтобы они, когда будут доступны, перезвонили, вместо того, чтобы ждать их ответа.
Давайте посмотрим, как сделать то же самое, что и выше, используя Tornado:
Предпоследняя строка кода вызывает метод , который получает данные по URL-адресу неблокирующим способом. Этот метод выполняется и возвращается немедленно. Поскольку каждая следующая строка будет выполнена до того, как будет получен ответ по URL-адресу, невозможно получить объект, как результат выполнения метода. Решение этой проблемы заключается в том, что метод вместо того, чтобы возвращать объект, вызывает функцию с результатом или обратный вызов. Обратный вызов в этом примере — .
В примере вы можете заметить, что первая строка функции проверяет наличие ошибки. Это необходимо, потому что невозможно обработать исключение. Если исключение было создано, то оно не будет отрабатываться в коде из-за цикла событий. Когда выполняется, он запускает HTTP-запрос, а затем обрабатывает ответ в цикле событий. К тому моменту, когда возникнет ошибка, стек вызовов будет содержать только цикл событий и текущую функцию, при этом нигде в коде не сработает исключение. Таким образом, любые исключения, созданные в функции обратного вызова, прерывают цикл событий и останавливают выполнение программы. Поэтому все ошибки должны быть переданы как объекты, а не обработаны в виде исключений. Это означает, что если вы не проверили наличие ошибок, то они не будут обрабатываться.
Другая проблема с обратными вызовами заключается в том, что в асинхронном программировании единственный способ избегать блокировок — это обратный вызов. Это может привести к очень длинной цепочке: обратный вызов после обратного вызова после обратного вызова. Поскольку теряется доступ к стеку и переменным, вы в конечном итоге переносите большие объекты во все ваши обратные вызовы, но если вы используете сторонние API-интерфейсы, то не можете передать что-либо в обратный вызов, если он этого не может принять. Это также становится проблемой, потому что каждый обратный вызов действует как поток. Например, вы хотели бы вызвать три API-интерфейса и дождаться, пока все три вернут результат, чтобы его обобщить. В Gevent вы можете это сделать, но не с обратными вызовами. Вам придется немного поколдовать, сохраняя результат в глобальной переменной и проверяя в обратном вызове, является ли результат окончательным.
Simple HTTPS Web server (low-level api)¶
import asyncio
import socket
import ssl
def make_header():
head = b'HTTP/1.1 200 OK\r\n'
head += b'Content-type: text/html\r\n'
head += b'\r\n'
return head
def make_body():
resp = b'<html>'
resp += b'<h1>Hello SSL</h1>'
resp += b'</html>'
return resp
sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM, )
sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
sock.setblocking(False)
sock.bind(('localhost' , 4433))
sock.listen(10)
sslctx = ssl.SSLContext(ssl.PROTOCOL_SSLv23)
sslctx.load_cert_chain(certfile='./root-ca.crt',
keyfile='./root-ca.key')
def do_handshake(loop, sock, waiter):
sock_fd = sock.fileno()
try
sock.do_handshake()
except ssl.SSLWantReadError
loop.remove_reader(sock_fd)
loop.add_reader(sock_fd, do_handshake,
loop, sock, waiter)
return
except ssl.SSLWantWriteError
loop.remove_writer(sock_fd)
loop.add_writer(sock_fd, do_handshake,
loop, sock, waiter)
return
loop.remove_reader(sock_fd)
loop.remove_writer(sock_fd)
waiter.set_result(None)
def handle_read(loop, conn, waiter):
try
req = conn.recv(1024)
except ssl.SSLWantReadError
loop.remove_reader(conn.fileno())
loop.add_reader(conn.fileno(), handle_read,
loop, conn, waiter)
return
loop.remove_reader(conn.fileno())
waiter.set_result(req)
def handle_write(loop, conn, msg, waiter):
try
resp = make_header()
resp += make_body()
ret = conn.send(resp)
except ssl.SSLWantReadError
loop.remove_writer(conn.fileno())
loop.add_writer(conn.fileno(), handle_write,
loop, conn, waiter)
return
loop.remove_writer(conn.fileno())
conn.close()
waiter.set_result(None)
async def server(loop):
while True
conn, addr = await loop.sock_accept(sock)
conn.setblocking(False)
sslconn = sslctx.wrap_socket(conn,
server_side=True,
do_handshake_on_connect=False)
# wait SSL handshake
waiter = loop.create_future()
do_handshake(loop, sslconn, waiter)
await waiter
# wait read request
waiter = loop.create_future()
handle_read(loop, sslconn, waiter)
msg = await waiter
# wait write response
waiter = loop.create_future()
handle_write(loop, sslconn, msg, waiter)
await waiter
loop = asyncio.get_event_loop()
try
loop.run_until_complete(server(loop))
finally
loop.close()
output:
绑定回调
绑定回调,在task执行完成的时候可以获取执行的结果,回调的最后一个参数是future对象,通过该对象可以获取协程返回值。
import time
import asyncio
now = lambda : time.time()
async def do_some_work(x):
print("waiting:",x)
return "Done after {}s".format(x)
def callback(future):
print("callback:",future.result())
start = now()
coroutine = do_some_work(2)
loop = asyncio.get_event_loop()
task = asyncio.ensure_future(coroutine)
print(task)
task.add_done_callback(callback)
print(task)
loop.run_until_complete(task)
print("Time:", now()-start)
结果为:
<Task pending coro=<do_some_work() running at /app/py_code/study_asyncio/simple_ex3.py:13>> <Task pending coro=<do_some_work() running at /app/py_code/study_asyncio/simple_ex3.py:13> cb=[callback() at /app/py_code/study_asyncio/simple_ex3.py:18]> waiting: 2 callback: Done after 2s Time: 0.00039196014404296875
通过add_done_callback方法给task任务添加回调函数,当task(也可以说是coroutine)执行完成的时候,就会调用回调函数。并通过参数future获取协程执行的结果。这里我们创建 的task和回调里的future对象实际上是同一个对象
Зеленые потоки
Зеленые потоки (green threads) являются примитивным уровнем асинхронного программирования. Зеленый поток — это обычный поток, за исключением того, что переключения между потоками производятся в коде приложения, а не в процессоре. Gevent — известная Python-библиотека для использования зеленых потоков. Gevent — это зеленые потоки и сетевая библиотека неблокирующего ввода-вывода Eventlet. изменяет поведение стандартных библиотек Python таким образом, что они позволяют выполнять неблокирующие операции ввода-вывода. Вот пример использования Gevent для одновременного обращения к нескольким URL-адресам:
Как видите, API-интерфейс Gevent выглядит так же, как и потоки. Однако за кадром он использует сопрограммы (coroutines), а не потоки, и запускает их в цикле событий (event loop) для постановки в очередь. Это значит, что вы получаете преимущества потоков, без понимания сопрограмм, но вы не избавляетесь от проблем, связанных с потоками. Gevent — хорошая библиотека, но только для тех, кто понимает, как работают потоки.
Студенческое соревнование по кибербезопасности «Кибервызов: новый уровень»
29–31 августа, онлайн, беcплатно
tproger.ru
События и курсы на tproger.ru
Давайте рассмотрим некоторые аспекты асинхронного программирования. Один из таких аспектов — это цикл событий. Цикл событий — это очередь событий/заданий и цикл, который вытягивает задания из очереди и запускает их. Эти задания называются сопрограммами. Они представляют собой небольшой набор команд, содержащих, помимо прочего, инструкции о том, какие события при необходимости нужно возвращать в очередь.
Почему GIL всё ещё используют?
Разработчики языка получили уйму жалоб касательно GIL. Но такой популярный язык как Python не может провести такое радикальное изменение, как удаление GIL, ведь это, естественно, повлечёт за собой кучу проблем несовместимости.
В прошлом разработчиками были предприняты попытки удаления GIL. Но все эти попытки разрушались существующими расширениями на C, которые плотно зависели от существующих GIL-решений. Естественно, есть и другие варианты, схожие с GIL. Однако они либо снижают производительность однопоточных и многопоточных I/O-приложений, либо попросту сложны в реализации. Вам бы не хотелось, чтобы в новых версиях ваша программа работала медленней, чем сейчас, ведь так?
Создатель Python, Guido van Rossum, в сентябре 2007 года высказался по поводу этого в статье «It isn’t Easy to remove the GIL»:
С тех пор ни одна из предпринятых попыток не удовлетворяла это условие.
Asynchronous Context Managers and «async with»
An asynchronous context manager is a context manager that is able to
suspend execution in its enter and exit methods.
To make this possible, a new protocol for asynchronous context managers
is proposed. Two new magic methods are added: __aenter__ and
__aexit__. Both must return an awaitable.
An example of an asynchronous context manager:
class AsyncContextManager:
async def __aenter__(self):
await log('entering context')
async def __aexit__(self, exc_type, exc, tb):
await log('exiting context')
A new statement for asynchronous context managers is proposed:
async with EXPR as VAR:
BLOCK
which is semantically equivalent to:
mgr = (EXPR)
aexit = type(mgr).__aexit__
aenter = type(mgr).__aenter__
VAR = await aenter(mgr)
try:
BLOCK
except:
if not await aexit(mgr, *sys.exc_info()):
raise
else:
await aexit(mgr, None, None, None)
As with regular with statements, it is possible to specify multiple
context managers in a single async with statement.
It is an error to pass a regular context manager without __aenter__
and __aexit__ methods to async with. It is a SyntaxError
to use async with outside of an async def function.
新线程协程
import asyncio
import time
from threading import Thread
now = lambda :time.time()
def start_loop(loop):
asyncio.set_event_loop(loop)
loop.run_forever()
async def do_some_work(x):
print('Waiting {}'.format(x))
await asyncio.sleep(x)
print('Done after {}s'.format(x))
def more_work(x):
print('More work {}'.format(x))
time.sleep(x)
print('Finished more work {}'.format(x))
start = now()
new_loop = asyncio.new_event_loop()
t = Thread(target=start_loop, args=(new_loop,))
t.start()
print('TIME: {}'.format(time.time() - start))
asyncio.run_coroutine_threadsafe(do_some_work(6), new_loop)
asyncio.run_coroutine_threadsafe(do_some_work(4), new_loop)
上述的例子,主线程中创建一个new_loop,然后在另外的子线程中开启一个无限事件循环。 主线程通过run_coroutine_threadsafe新注册协程对象。这样就能在子线程中进行事件循环的并发操作,同时主线程又不会被block。一共执行的时间大概在6s左右。
Пример 2. Простой кооперативный параллелизм
Следующая версия программы демонстрирует возможности двух задач работать совместно с использованием генераторов. Добавление оператора yield в функцию task означает, что после выполнения этого оператора, функция завершает свою работу, но сохраняя свой контекст до следующего запуска. Затем цикл выполнения задачи возобновляет выполнение программы используя вызвав метод t.next(). Этот оператор перезапускает задачу в том месте, где она ранее вызывалась.
Это разновидность кооперативного параллелизма. Программа предоставляет контроль над своим контекстом, позволяя работать и чему-то другому. Это позволяет нашему примитивному планировщику запускать задачи по два экземпляра функции task, каждая из которых обращается в работе к одной очереди. Этот пример более продвинутый, но требует больше кода, чтобы получить схожие результаты, что и в первом примере.
После рассмотрения вывода работы программы, становится видно, что поочерёдно выполняются задачи One и Two, расходуя содержимое work_queue. Как и было задумано, обе task выполняют свою работу, и каждая из них заканчивает обработку двух элементов из очереди. Но опять же, довольно много работы для достижения результата.
Изюминка заключается в использовании оператора yield, который превращает функцию task в функция генератор, «переключатель контекста». Программа использует переключение контекста для запуска двух экземпляров task.
Синхронное и асинхронное выполнение
В видео “Конкурентность — это не параллелизм, это лучше” Роб Пайк обращает ваше внимание на ключевую вещь. Разбиение задач на конкурентные подзадачи возможно только при таком параллелизме, когда он же и управляет этими подзадачами
Asyncio делает тоже самое — вы можете разбивать ваш код на процедуры, которые определять как корутины, что даёт возможность управлять ими как пожелаете, включая и одновременное выполнение. Корутины содержат операторы yield, с помощью которых мы определяем места, где можно переключиться на другие ожидающие выполнения задачи.
За переключение контекста в asyncio отвечает yield, который передаёт управление обратно в event loop, а тот в свою очередь — к другой корутине. Рассмотрим базовый пример:
- Сначала мы объявили пару простейших корутин, которые притворяются неблокирующими, используя sleep из asyncio
- Корутины могут быть запущены только из другой корутины, или обёрнуты в задачу с помощью create_task
- После того, как у нас оказались 2 задачи, объединим их, используя wait
- И, наконец, отправим на выполнение в цикл событий через run_until_complete
Используя await в какой-либо корутине, мы таким образом объявляем, что корутина может отдавать управление обратно в event loop, который, в свою очередь, запустит какую-либо следующую задачу: bar. В bar произойдёт тоже самое: на await asyncio.sleep управление будет передано обратно в цикл событий, который в нужное время вернётся к выполнению foo.
Представим 2 блокирующие задачи: gr1 и gr2, как будто они обращаются к неким сторонним сервисам, и, пока они ждут ответа, третья функция может работать асинхронно.
Обратите внимание как происходит работа с вводом-выводом и планированием выполнения, позволяя всё это уместить в один поток. Пока две задачи заблокированы ожиданием I/O, третья функция может занимать всё процессорное время
TLS Upgrade¶
New in Python 3.7
import asyncio
import ssl
class HttpClient(asyncio.Protocol):
def __init__(self, on_con_lost):
self.on_con_lost = on_con_lost
self.resp = b""
def data_received(self, data):
self.resp += data
def connection_lost(self, exc):
resp = self.resp.decode()
print(resp.split("\r\n")[])
self.on_con_lost.set_result(True)
async def main():
paths = ssl.get_default_verify_paths()
sslctx = ssl.SSLContext()
sslctx.verify_mode = ssl.CERT_REQUIRED
sslctx.check_hostname = True
sslctx.load_verify_locations(paths.cafile)
loop = asyncio.get_running_loop()
on_con_lost = loop.create_future()
tr, proto = await loop.create_connection(
lambda HttpClient(on_con_lost), "github.com", 443
)
new_tr = await loop.start_tls(tr, proto, sslctx)
req = f"GET / HTTP/1.1\r\n"
req += "Host: github.com\r\n"
req += "Connection: close\r\n"
req += "\r\n"
new_tr.write(req.encode())
await on_con_lost
new_tr.close()
asyncio.run(main())
output:
Потоки
Еще одна эффективная оптимизация скорости — это потоковая передача запросов. При отправке запроса по умолчанию все тело ответа загружается немедленно. Лучший способ не загружать весь контент в память сразу при запросе. Для этого есть параметра stream, в библиотеке requests или атрибут content в aiohttp.
Потоковая передача с requests
import requests
# Use `with` to make sure the response stream is closed and the connection can
# be returned back to the pool.
with requests.get('http://example.org', stream=True) as r:
print(list(r.iter_content()))
Потоковая передача с aiohttp
import aiohttp
import asyncio
async def get(url):
async with aiohttp.ClientSession() as session:
async with session.get(url) as response:
return await response.content.read()
loop = asyncio.get_event_loop()
tasks = [asyncio.ensure_future(get("http://example.com"))]
loop.run_until_complete(asyncio.wait(tasks))
print("Results: %s" % )
Не загружать полный контент крайне важно, чтобы избежать ненужного выделения сотен мегабайт памяти. Если вашей программе не требуется доступ ко всему содержимому в целом, но она может работать с частями
Например, если вы собираетесь сохранить и записать содержимое в файл, чтение только куска и одновременная запись будет гораздо более эффективным, чем чтение всего тела HTTP, выделяя огромную кучу памяти , и только после этого записать его на диск.
Упрощенный web сервер
Его основная работа такая же, как и приведённый выше пакетных обработчиков, т. е. получить некоторые входные данные, обработать их и вернуть выходные. Если реализовать его в виде синхронной программы, то это был бы абсолютно ужасный веб-сервер.
Почему? Потому что web сервер должен обрабатывать сотни, а иногда и тысячи подключений от пользователей одновременно, а не обслуживать только одного клиента.
Можно ли как-то улучшить синхронный web сервер? Конечно, можно оптимизировать шаги исполнения, сделав их как можно быстрее. К сожалению, нужного эффекта по улучшению работы web сервера это не даст, и он не сможет возвращать ответы достаточно быстро, и не сможет обслуживать достаточное количество пользователей.
Каковы реальные пределы оптимизации указанного подхода? Скорость сети, скорость файлового IO, скорость запроса к базе данных, скорость других подключенных служб и т. д. Общей особенностью этого списка являются то, что все они являются функциями ввода-вывода. Они все на много порядков медленнее, чем скорость работы CPU.
Например, если выполняется запрос к базе данных в синхронной программе, прежде чем будет возвращён ответ клиенту и переход к следующему шагу, CPU будет находиться в состоянии длительного ожидания.
Файловый IO, сеть, база работают достаточно быстро, но намного медленнее, чем CPU. Технологии асинхронного программирования позволяют программам воспользоваться относительно медленными IO процессами, при этом нагружая CPU выполнением другими вычислениями, освобождая его от необходимости ожидать.
Разработка асинхронных программ сложнее синхронных. И это странно, потому что мир, в котором мы живем и с которым взаимодействуем, почти полностью асинхронен.
Пример из жизни. Многие из нас являются родителями, поэтому чтобы больше успеть мы делаем несколько вещей одновременно – домашняя бухгалтерия, стирка и присмотр за детьми.
A quick concurrent.futures summary
The provides a high-level abstraction for the and modules, which is why we won’t discuss those modules in detail within this post. In fact the module is a very low-level API that the module is itself built on top of (again, this is why we won’t be covering that either).
Now we’ve already mentioned that asyncio helps us avoid using threads so why would we want to use if it’s just an abstraction on top of threads (and multiprocessing)? Well, because not all libraries/modules/APIs support the asyncio model.
For example, if you use and interact with AWS S3, then you’ll find those are synchronous operations. You can wrap those calls in multi-threaded code, but it would be better to use as it means you not only benefit from traditional threads but an asyncio friendly package.
The module is also designed to interop with the asyncio event loop, making it easier to work with a pool of threads/subprocesses within an otherwise asyncio driven application.
Additionally you’ll also want to utilize when you require a pool of threads or a pool of subprocesses, while also using a clean and modern Python API (as apposed to the more flexible but low-level or modules).
Awaitables
The driving force behind asyncio is the ability to schedule asynchronous ‘tasks’. There are a few different types of objects in Python that help support this, and they are generally grouped by the term ‘awaitable’.
Ultimately, something is awaitable if it can be used in an expression.
There are three main types of awaitables:
- Coroutines
- Tasks
- Futures
Coroutines
There are two closely related terms used here:
- a coroutine function: an function.
- a coroutine object: an object returned by calling a coroutine function.
Tasks
are used to schedule coroutines concurrently.
All asyncio applications will typically have (at least) a single ‘main’ entrypoint task that will be scheduled to run immediately on the event loop. This is done using the function (see ‘’).
A coroutine function is expected to be passed to , while internally asyncio will check this using the helper function (see: ). If not a coroutine, then an error is raised, otherwise the coroutine will be passed to (see: ).
The function expects a (see below section for what a Future is) and uses another helper function to check the type provided. If not a Future, then the low-level API is used to convert the coroutine into a Future (see ).
In older versions of Python, if you were going to manually create your own Future and schedule it onto the event loop, then you would have used (now considered to be a low-level API), but with Python 3.7+ this has been superseded by .
Additionally with Python 3.7, the idea of interacting with the event loop directly (e.g. getting the event loop, creating a task with and then passing it to the event loop) has been replaced with , which abstracts it all away for you (see ‘’ to understand what that means).
The following APIs let you see the state of the tasks running on the event loop:
Futures
A Future is a low-level awaitable object that represents an eventual result of an asynchronous operation.
To use an analogy: it’s like an empty postbox. At some point in the future the postman will arrive and stick a letter into the postbox.
This API exists to enable callback-based code to be used with /, while is an example of an asyncio low-level API function that returns a Future (see also some of the APIs listed in ).
Эффективный секретарь
Теперь давайте рассмотрим эти понятия на примерах из жизни. Представьте секретаря, который настолько эффективен, что не тратит время впустую. У него есть пять заданий, которые он выполняет одновременно: отвечает на телефонные звонки, принимает посетителей, пытается забронировать билеты на самолет, контролирует графики встреч и заполняет документы. Теперь представьте, что такие задачи, как контроль графиков встреч, прием телефонных звонков и посетителей, повторяются не часто и распределены во времени. Таким образом, большую часть времени секретарь разговаривает по телефону с авиакомпанией, заполняя при этом документы. Это легко представить. Когда поступит телефонный звонок, он поставит разговор с авиакомпанией на паузу, ответит на звонок, а затем вернется к разговору с авиакомпанией. В любое время, когда новая задача потребует внимания секретаря, заполнение документов будет отложено, поскольку оно не критично. Секретарь, выполняющий несколько задач одновременно, переключает контекст в нужное ему время. Он асинхронный.
Потоки — это пять секретарей, у каждого из которых по одной задаче, но только одному из них разрешено работать в определенный момент времени. Для того, чтобы секретари работали в потоковом режиме, необходимо устройство, которое контролирует их работу, но ничего не понимает в самих задачах. Поскольку устройство не понимает характер задач, оно постоянно переключалось бы между пятью секретарями, даже если трое из них сидят, ничего не делая. Около 57% (чуть меньше, чем 3/5) переключения контекста были бы напрасны. Несмотря на то, что переключение контекста процессора является невероятно быстрым, оно все равно отнимает время и ресурсы процессора.
Создание конкурентности
До сих пор мы использовали единственный метод создания и получения результатов из корутин, создания набора задач и ожидания их завершения. Однако, корутины могут быть запланированы для запуска и получения результатов несколькими способами. Представьте ситуацию, когда нам надо обрабатывать результаты GET-запросов по мере их получения; на самом деле реализация очень похожа на предыдущую:
Посмотрите на отступы и тайминги — мы запустили все задачи одновременно, однако они обработаны в порядке завершения выполнения. Код в данном случае немного отличается: мы пакуем корутины, каждая из которых уже подготовлена для выполнения, в список. Функция возвращает итератор, который выдаёт результаты корутин по мере их выполнения. Круто же, правда?! Кстати, и as_completed, и wait — функции из пакета .
Ещё один пример — что если вы хотите узнать свой IP адрес. Есть куча сервисов для этого, но вы не знаете какой из них будет доступен в момент работы программы. Вместо того, чтобы последовательно опрашивать каждый из списка, можно запустить все запросы конкурентно и выбрать первый успешный.
Что ж, для этого в нашей любимой функции есть специальный параметр return_when. До сих пор мы игнорировали то, что возвращает wait, т.к. только распараллеливали задачи. Но теперь нам надо получить результат из корутины, так что будем использовать набор футур done и pending.
Что же случилось? Первый сервис ответил успешно, но в логах какое-то предупреждение!
На самом деле мы запустили выполнение двух задач, но вышли из цикла уже после первого результата, в то время как вторая корутина ещё выполнялась. Asyncio подумал что это баг и предупредил нас. Наверно, стоит прибираться за собой и явно убивать ненужные задачи. Как? Рад, что вы спросили.
Why the asend() and athrow() methods are necessary
They make it possible to implement concepts similar to
contextlib.contextmanager using asynchronous generators.
For instance, with the proposed design, it is possible to implement
the following pattern:
@async_context_manager
async def ctx():
await open()
try:
yield
finally:
await close()
async with ctx():
await ...
Another reason is that it is possible to push data and throw exceptions
into asynchronous generators using the object returned from
__anext__ object, but it is hard to do that correctly. Adding
explicit asend() and athrow() will pave a safe way to
accomplish that.
In terms of implementation, asend() is a slightly more generic
version of __anext__, and athrow() is very similar to
aclose(). Therefore having these methods defined for asynchronous
generators does not add any extra complexity.