async/await, async генераторы, asyncio интеграция
- Маршрут: Работа в эксплуатации · тема 3 из 6
- Навыки:
M5,M6,S4,S5- До урока:
functions,exceptions,context_managers,python_testing,test_doubles,type_system.- Результат: ограниченная конкурентность с однозначными правилами срока ожидания, повторов, отмены и очистки.
- Основа: корутины, задачи и базовая координация. Работа в эксплуатации: ограничение конкурентности, общий срок, повторы и частичные сбои. Углубление: структурированная конкурентность и выбор среды выполнения; вернитесь после
python_performance.- Подтверждение: устойчивый пакетный клиент и тесты отмены/очистки.
Асинхронность организует конкурентное ожидание I/O. Она может работать в одном потоке и не означает параллельного выполнения Python-кода на нескольких ядрах.
| Подход | Модель | Годится для |
|---|---|---|
threading | Потоки ОС | I/O-bound (GIL освобождается) |
multiprocessing | Процессы ОС | CPU-bound |
asyncio | Event loop, обычно один поток | Много конкурентных I/O-ожиданий |
Синхронно: A──────────────────── B──────────────────── C────────────
Асинхронно: A────await──────A B────await──────B C────await──────C
↑ переключение ↑ переключение
Каждый await — это возможность event loop запустить другую корутину пока первая ждёт ответа.
Event loop держит очередь готовых callback-функций и таймеров. Для сетевого
I/O реализация обычно ждёт готовности файловых дескрипторов через системный
selector: epoll в Linux или kqueue в BSD/macOS. Когда дескриптор готов к
чтению или записи, соответствующий callback попадает в ready queue. Это
конкурентность в одном потоке, а не параллельное исполнение Python-кода.
import asyncio
async def greet(name: str) -> str:
await asyncio.sleep(1) # I/O-ожидание; event loop переключается
return f"Hello, {name}!"
# Вызов async def создаёт объект coroutine — код ещё не выполнялся
coro = greet("Alice") # <coroutine object greet at 0x...>
result = asyncio.run(coro) # запускаем и ждёмasync def main():
# create_task() планирует корутину для выполнения в event loop
task = asyncio.create_task(greet("Bob"))
# Код задачи начнётся, когда текущая корутина отдаст управление event loop.
await asyncio.sleep(0) # даём event loop шанс переключиться
result = await task # дожидаемся результата
print(result)Для новой корутины предпочитайте asyncio.create_task(). Более старый
asyncio.ensure_future() также принимает уже созданный Future и возвращает
его без оборачивания, поэтому он всё ещё встречается в коде интеграций.
await task возвращает результат или повторно поднимает исключение задачи.
После завершения тот же объект ошибки можно получить через task.exception();
у незавершённой задачи этот вызов поднимет InvalidStateError, а у отменённой
— CancelledError.
import asyncio
# Future — низкоуровневый объект; обычно не создаётся напрямую
# Task наследует от Future
async def future_example():
loop = asyncio.get_running_loop()
future = loop.create_future()
future.set_result(42)
return await future # → 42| Тип | Что это | Когда нужен |
|---|---|---|
coroutine | Функция async def, не запущенная | Определение логики |
Task | Запланированная корутина | Конкурентное выполнение |
Future | Placeholder для результата | Интеграция с callback-кодом |
asyncio.run() — точка входаimport asyncio
async def main():
print("Start")
await asyncio.sleep(1)
print("End")
# Python 3.7+: запускает новый event loop
asyncio.run(main())
# Проблема: Нельзя вызывать asyncio.run() внутри уже запущенного event loop
# (например, в Jupyter — там уже есть event loop, используйте await)asyncio.gather() — запуск нескольких корутинimport asyncio
import aiohttp
async def fetch(session, url):
async with session.get(url) as resp:
return await resp.text()
async def fetch_all(urls):
async with aiohttp.ClientSession() as session:
# Все запросы выполняются конкурентно во время I/O-ожиданий
results = await asyncio.gather(
*[fetch(session, url) for url in urls]
)
return results
# Для большого списка добавьте Semaphore из раздела 6, чтобы не перегрузить
# собственный клиент и удалённый сервер.
urls = [f"https://example.com/page/{i}" for i in range(20)]
asyncio.run(fetch_all(urls))Обработка исключений в gather():
async def fetch_all_with_errors(urls):
async with aiohttp.ClientSession() as session:
# По умолчанию первое исключение сразу передаётся ожидающему коду.
# Остальные awaitable при этом не отменяются автоматически.
results = await asyncio.gather(
*(fetch(session, url) for url in urls),
return_exceptions=True, # ошибки как результаты, не исключения
)
for result in results:
if isinstance(result, Exception):
print(f"Ошибка: {result}")
return resultsasyncio.wait() — гибкое ожиданиеasync def fetch_until_first(urls):
async with aiohttp.ClientSession() as session:
tasks = [
asyncio.create_task(fetch(session, url))
for url in urls
]
done, pending = await asyncio.wait(
tasks,
return_when=asyncio.FIRST_COMPLETED,
)
# cancel() только запрашивает отмену. Задачи нужно дождаться, чтобы они
# выполнили finally и исключения не потерялись.
for task in pending:
task.cancel()
await asyncio.gather(*pending, return_exceptions=True)
return [task.result() for task in done]asyncio.wait_for() — таймаутasync def with_timeout():
try:
return await asyncio.wait_for(slow_operation(), timeout=5.0)
except TimeoutError: # встроенный тип начиная с Python 3.11
print("Превышено время ожидания!")
return NoneПри таймауте wait_for() отменяет ожидаемую задачу или корутину и ждёт
завершения отмены, затем поднимает TimeoutError. Если операция не должна
получать отмену от этого ожидания, её нужно передать через shield() и отдельно
организовать дальнейший жизненный цикл.
asyncio.TaskGroup (Python 3.11+) — structured concurrencyasync def main():
async with aiohttp.ClientSession() as session:
async with asyncio.TaskGroup() as tg:
tasks = [
tg.create_task(fetch(session, url))
for url in urls
]
# Здесь все задачи завершены
# Если любая задача завершилась с ошибкой — все остальные отменяются
print([task.result() for task in tasks])async def long_task():
try:
await asyncio.sleep(100)
except asyncio.CancelledError:
print("Задача отменена — выполняю cleanup")
await cleanup()
raise # важно: перебросить CancelledError
async def cancellation_example():
task = asyncio.create_task(long_task())
await asyncio.sleep(1)
task.cancel() # запрашивает отмену
try:
await task
except asyncio.CancelledError:
print("Задача была отменена")
# asyncio.shield() — не передавать внешний cancel внутренней Task
critical_task = asyncio.create_task(critical_op()) # сильная ссылка обязательна
try:
return await asyncio.wait_for(
asyncio.shield(critical_task),
timeout=5,
)
except TimeoutError:
# critical_task не отменена этим timeout; её нужно явно дождаться,
# отменить при shutdown и обязательно прочитать возможную ошибку.
critical_task.cancel()
await asyncio.gather(critical_task, return_exceptions=True)
return Nonelock = asyncio.Lock()
async def safe_operation():
async with lock: # ждёт пока блокировка свободна
await modify_shared_state()# Максимум 10 одновременных запросов
semaphore = asyncio.Semaphore(10)
async def fetch_limited(session, url):
async with semaphore:
return await fetch(session, url)
async def fetch_many_limited(urls):
# Клиентская сессия переиспользуется; максимум 10 операций одновременно.
async with aiohttp.ClientSession() as session:
return await asyncio.gather(
*(fetch_limited(session, url) for url in urls)
)Semaphore ограничивает concurrency, а не частоту «N запросов за секунду». Для rate limiting нужен алгоритм, учитывающий время, например token bucket.
ready = asyncio.Event()
async def producer():
await prepare_data()
ready.set() # сигнал
async def consumer():
await ready.wait() # ждёт сигнала
process_data()
async def event_example():
await asyncio.gather(producer(), consumer())async def producer(queue: asyncio.Queue):
for i in range(10):
await queue.put(i)
await queue.put(None) # sentinel
async def consumer(queue: asyncio.Queue):
while True:
item = await queue.get()
try:
if item is None:
return
await process(item)
finally:
# Каждый успешный get(), включая sentinel, имеет ровно один
# task_done(). Иначе queue.join() зависнет.
queue.task_done()
async def main():
queue = asyncio.Queue(maxsize=5)
consumer_task = asyncio.create_task(consumer(queue))
await producer(queue)
await queue.join()
await consumer_taskasyncio.to_thread() (Python 3.9+)import asyncio
def blocking_io(filename):
with open(filename) as f:
return f.read()
async def main():
# Запускает синхронную функцию в thread pool
content = await asyncio.to_thread(blocking_io, "large_file.txt")
return contentОсновное применение to_thread() — блокирующий ввод-вывод или синхронная
библиотека без асинхронного API. В обычной сборке с GIL чистый Python-код с
процессорной нагрузкой не получает от пула потоков привычного многоядерного
ускорения. Исключения: нативный код, освобождающий GIL, и сборка CPython без GIL.
loop.run_in_executor()import concurrent.futures
async def main():
loop = asyncio.get_running_loop()
# ThreadPoolExecutor — для I/O-bound синхронного кода
with concurrent.futures.ThreadPoolExecutor() as pool:
result = await loop.run_in_executor(pool, blocking_function, arg)
# ProcessPoolExecutor — для CPU-bound кода
with concurrent.futures.ProcessPoolExecutor() as pool:
result = await loop.run_in_executor(pool, cpu_intensive, data)# Async generator
async def fetch_pages(base_url, pages):
async with aiohttp.ClientSession() as session:
for page in range(pages):
url = f"{base_url}?page={page}"
async with session.get(url) as resp:
yield await resp.json()
# Использование через async for
async def process_pages():
async for page_data in fetch_pages("https://api.example.com/items", 100):
await store(page_data)
# Async iterator через класс
class AsyncCounter:
def __init__(self, stop):
self.current = 0
self.stop = stop
def __aiter__(self):
return self
async def __anext__(self):
if self.current >= self.stop:
raise StopAsyncIteration
await asyncio.sleep(0.1)
self.current += 1
return self.currentasyncio.timeout() (Python 3.11+)async def main():
try:
async with asyncio.timeout(10):
await long_running_task()
except TimeoutError:
print("Задача не завершилась за 10 секунд")
# Можно обновить дедлайн:
async with asyncio.timeout(10) as deadline:
await first_part()
# Увеличиваем таймаут если первая часть быстрая
loop = asyncio.get_running_loop()
deadline.reschedule(loop.time() + 5) # абсолютный deadline
await second_part()# Включение debug mode — более детальные ошибки
asyncio.run(main(), debug=True)
# или
# Для ручного управления жизненным циклом используйте Runner:
with asyncio.Runner(debug=True) as runner:
runner.run(main())
# asyncio.sleep(0) — передать управление event loop
async def cooperative_task():
for i in range(1_000_000):
if i % 1000 == 0:
await asyncio.sleep(0) # дать шанс другим задачам
process(i)
# Получить текущий запущенный event loop внутри coroutine/callback
loop = asyncio.get_running_loop()В Python 3.14 get_running_loop() возвращает только активный loop.
get_event_loop() сначала возвращает активный loop, а вне callback/coroutine —
явно установленный current loop; если нет ни одного, он поднимает
RuntimeError и больше не создаёт loop автоматически. Event-loop policies
deprecated; для верхнего уровня используйте asyncio.run() или
asyncio.Runner.
# Проблема: Забытый await — coroutine создана но не выполнена
async def main():
result = fetch_data() # coroutine object! await забыт
# Python выдаст RuntimeWarning: coroutine was never awaited
# Вариант: Всегда await coroutine
async def fixed():
return await fetch_data()
# Проблема: Блокирующий код в корутине — замораживает весь event loop
async def bad():
time.sleep(5) # блокирует ВСЕ корутины!
requests.get(url) # тоже блокирующий!
# Вариант: Используйте async-версии
async def good():
await asyncio.sleep(5)
async with aiohttp.ClientSession() as session:
await session.get(url)
# Проблема: asyncio.run() внутри корутины
async def nested():
asyncio.run(other_coro()) # RuntimeError!
# Вариант: Просто await
async def nested():
await other_coro()
# Проблема: Смешение asyncio.run() и nest_asyncio без необходимости
# В Jupyter используйте: await coro() (Jupyter имеет свой event loop)import asyncio
import time
from concurrent.futures import ThreadPoolExecutor
COUNT = 100
DELAY = 0.01
async def async_job():
await asyncio.sleep(DELAY)
async def async_test():
await asyncio.gather(*(async_job() for _ in range(COUNT)))
def thread_job(_):
time.sleep(DELAY)
def thread_test():
with ThreadPoolExecutor(max_workers=20) as pool:
list(pool.map(thread_job, range(COUNT)))
started = time.perf_counter()
asyncio.run(async_test())
print("async:", time.perf_counter() - started)
started = time.perf_counter()
thread_test()
print("threads:", time.perf_counter() - started)Этот опыт измеряет накладные расходы двух моделей на искусственном ожидании,
но не предсказывает скорость реального HTTP-клиента. Для решения измеряйте свой
DNS, TLS, сервер, лимиты соединений и одинаковую конкурентность. asyncio
обычно позволяет обслуживать много I/O-ожиданий меньшим числом потоков, но не
делает удалённую систему быстрее.
import asyncio
import random
from functools import wraps
def async_retry(max_attempts=3, delay=1.0, backoff=2.0,
exceptions=(Exception,)):
def decorator(func):
@wraps(func)
async def wrapper(*args, **kwargs):
current_delay = delay
for attempt in range(max_attempts):
try:
return await func(*args, **kwargs)
except exceptions as e:
if attempt + 1 >= max_attempts:
raise
jitter = random.uniform(0, current_delay * 0.1)
await asyncio.sleep(current_delay + jitter)
current_delay *= backoff
return wrapper
return decorator
@async_retry(max_attempts=3, delay=0.5, exceptions=(aiohttp.ClientError,))
async def fetch_with_retry(session, url):
async with session.get(url, timeout=aiohttp.ClientTimeout(total=10)) as resp:
resp.raise_for_status()
return await resp.json()import asyncio
import time
class AsyncRateLimiter:
"""Ограничивает до N запросов в секунду с burst."""
def __init__(self, rate: float, burst: int = 1):
self._rate = rate # запросов в секунду
self._burst = burst # максимальный burst
self._tokens = burst
self._last_update = time.monotonic()
self._lock = asyncio.Lock()
async def acquire(self):
async with self._lock:
now = time.monotonic()
elapsed = now - self._last_update
self._tokens = min(
self._burst,
self._tokens + elapsed * self._rate
)
self._last_update = now
if self._tokens < 1:
sleep_time = (1 - self._tokens) / self._rate
await asyncio.sleep(sleep_time)
self._tokens = 0
else:
self._tokens -= 1
limiter = AsyncRateLimiter(rate=10, burst=5) # 10 req/s, burst до 5
async def limited_request(session, url):
await limiter.acquire()
return await fetch(session, url)import asyncio
import threading
# Вариант 1: asyncio.run() — блокирует текущий поток
result = asyncio.run(async_function())
# Вариант 2: передать корутину в loop, работающий в другом потоке
def sync_caller(loop, coro):
future = asyncio.run_coroutine_threadsafe(coro, loop)
return future.result(timeout=10) # блокирует до результата
# Для обычного callback без coroutine:
loop.call_soon_threadsafe(callback, argument)Оба thread-safe метода нужны только при обращении к loop из другого потока.
Внутри корутины вызывайте await или create_task() напрямую. В Jupyter уже
работает event loop, поэтому на верхнем уровне ячейки используйте await, а не
вкладывайте новый asyncio.run().
Для тестирования async-функций нужен pytest-asyncio — плагин, который
управляет event loop'ом для каждого теста.
pip install pytest-asyncioНастройка в pyproject.toml или pytest.ini:
# pyproject.toml
[tool.pytest.ini_options]
asyncio_mode = "auto" # все async-тесты автоматически получают @pytest.mark.asyncioБез asyncio_mode = "auto" каждый async-тест нужно помечать явно:
import pytest
import asyncio
@pytest.mark.asyncio
async def test_fetch():
result = await fetch_user(user_id=1)
assert result["id"] == 1
# С asyncio_mode = "auto" декоратор не нужен:
async def test_fetch_auto():
result = await fetch_user(user_id=1)
assert result["id"] == 1Мокирование async-функций:
from unittest.mock import AsyncMock, patch
@pytest.mark.asyncio
async def test_with_mock():
with patch("mymodule.fetch_user", new_callable=AsyncMock) as mock:
mock.return_value = {"id": 1, "name": "Alice"}
result = await process_user(1)
mock.assert_awaited_once_with(user_id=1)Тестирование таймаутов:
@pytest.mark.asyncio
async def test_timeout():
with pytest.raises(TimeoutError):
await asyncio.wait_for(asyncio.sleep(10), timeout=0.1)Подробнее про двойники для корутин — AsyncMock, assert_awaited_once_with и
асинхронные фикстуры через pytest_asyncio.fixture — в теме test_doubles.
На заметку: pytest-asyncio поддерживает режимы обнаружения async-тестов
strict и auto. В strict режиме каждый тест требует явной
метки; в auto — все async-тесты работают автоматически. Для нового проекта
auto — удобнее; для миграции старого кодобаза — strict безопаснее.
Асинхронный код обязан иметь владельца задач и ограничение конкурентности. Результат основы — пакетная операция, у которой измеряется верхняя граница одновременной работы и нет потерянных фоновых задач.
import asyncio
async def probe() -> None:
semaphore = asyncio.Semaphore(2)
active = 0
maximum = 0
async def worker() -> None:
nonlocal active, maximum
async with semaphore:
active += 1
maximum = max(maximum, active)
await asyncio.sleep(0)
active -= 1
async with asyncio.TaskGroup() as group:
for _ in range(10):
group.create_task(worker())
assert maximum == 2
assert active == 0
asyncio.run(probe())Практика с подсказками: внедрите транспорт, разделите срок одной операции и срок всего пакета, повторяйте только классифицированные временные ошибки. Проверьте ограничение конкурентности, отмену родителя, очистку и отсутствие секретов в метриках.
Сравните gather и TaskGroup на одновременных сбоях нескольких задач.
Зафиксируйте версии Python и наблюдаемую структуру исключений; затем добавьте
ограниченную очередь и докажите наличие обратного давления.
async def download_all(urls) загружающий все URL параллельно с ограничением 10 одновременных запросов.asyncio.Queue с несколькими потребителями.async def retry(coro_factory, max_attempts=3) — повторяет корутину при ошибке.asyncio.sleep(0) полезно в длинных синхронных циклах.asyncio.Semaphore, затем отдельно реализуйте rate limiter N запросов в секунду через token bucket. Объясните различие.AsyncCircuitBreaker — после N ошибок подряд блокирует запросы на T секунд.asyncio.CancelledError.try/finally: максимум
не должен превышать 10 даже при ошибках.get() получает task_done(), join() завершается, а
sentinel доходит до каждого потребителя.cancel() тест дожидается задачи и доказывает выполнение cleanup;
CancelledError не должен превращаться в обычный успешный результат.Лаборатория
«Устойчивый конкурентный пакетный клиент»
доступна в публичном GitLab. Она проверяет TaskGroup, ограничение
конкуренции, время ожидания операции, общий срок выполнения, выборочные
повторы, частичные отказы и корректное распространение отмены — без реальной
сети.
Склонируйте репозиторий и запустите тесты из его корня:
git clone --branch v0.6.2 --depth 1 https://gitlab.potapov.me/courses/python-labs.git
cd python-labs
uv sync --group test
uv run --group test pytest async_python/testsДалее: Безопасность в Python