async/await на практике в Python: паттерны отмены, таймаутов и дедлайнов
Соберём несколько типовых сценариев (стриминг, параллельные запросы, внешние вызовы) и разберём, как сделать отмену и таймауты корректными.
Содержание
async/await на практике в Python: паттерны отмены, таймаутов и дедлайнов
Асинхронность в Python — это не «ускоритель ради ускорителя», а набор инструментов для управления ожиданием: запросы к сети, ожидание файлового ввода-вывода, стриминг данных, конкурентные вычисления с ограничениями. И почти в любой реальной системе рано или поздно встают три вопроса:
- Как корректно отменять асинхронные операции, когда результат уже не нужен?
- Как ставить таймауты, не оставляя «висящие» задачи и не теряя контроль над исключениями?
- Как работать с дедлайнами, то есть временем «до», а не «через N секунд».
В этой статье соберём типовые сценарии (стриминг, параллельные запросы, внешние вызовы) и разберём паттерны отмены, таймаутов и дедлайнов на уровне кода: как именно делать, какие ошибки типичны и как их избежать.
Принципы: что async/await «умеет», а что нужно делать вам
Отмена — это не «остановить поток», а «сообщить задаче»
В asyncio отмена реализуется через исключение asyncio.CancelledError, которое пробрасывается в корутину в момент, когда она делает «проверяемые» ожидания (например, await asyncio.sleep(...), await reader.read(...), await fetch(...) и т.п.).
Отсюда несколько практических правил:
- Не глотайте
CancelledErrorбездумно. Если вы поймали отмену, часто нужно либо пробросить её дальше, либо корректно завершить ресурсы, после чего всё равно вернуть отмену. - Если вы создаёте фоновые задачи (
asyncio.create_task), вы отвечаете за их жизнь: отменять, собирать результаты/исключения, предотвращать утечки. - Отмена может произойти в любой момент между
await. Значит, код должен быть устойчивым: закрывать соединения, сбрасывать буферы, корректно освобождать семафоры/локи.
Таймаут — это контроль ожидания, но не гарантия «остановки всего»
asyncio.wait_for() или таймауты в ClientSession прерывают ожидание корутины по таймеру. Но это не «магическое прекращение внешнего запроса» на уровне протокола: вы прерываете ваше ожидание, а дальше зависит от того, как реализован клиент.
Условно: таймаут — это про ваш контроль, а не про гарантированную остановку удалённого сервера. Чтобы система была надёжной, нужно уметь освобождать локальные ресурсы и отменять задачи.
Дедлайн — это «абсолютное время», а не просто N секунд
Для микросервисов и цепочек вызовов дедлайны критичны: иначе вы легко получите «таймауты, которые суммируются» (каждый сервис тратит свои N секунд, и итоговый запрос живёт дольше SLA).
Практический подход: дедлайн вычисляется как deadline = loop.time() + ttl, а дальше вы используете оставшееся время max(0, deadline - loop.time()) для каждого шага. Это переносит «время до» между слоями.
Базовые инструменты: CancelledError, Task, wait_for и shield
Отмена задачи и ожидание её завершения
Минимальный корректный паттерн для отмены фоновой задачи выглядит так:
- инициировать
task.cancel() - await task чтобы получить
CancelledError(и понять, что всё завершилось) - не подавлять отмену без нужды
Пример:
import asyncio
async def worker():
try:
while True:
await asyncio.sleep(0.5)
# имитация работы
except asyncio.CancelledError:
# здесь можно освободить ресурсы (закрыть поток, соединение и т.п.)
print("worker: cancelled, cleaning up")
raise # важно: пробросить отмену дальше
async def main():
task = asyncio.create_task(worker())
await asyncio.sleep(2)
task.cancel()
try:
await task
except asyncio.CancelledError:
print("main: worker task cancelled")
asyncio.run(main())
Почему asyncio.shield() иногда нужен
Случается, что вы хотите отменить «обёртку» (например, таймаут операции), но не отменять внутреннюю задачу, потому что она обслуживает общий ресурс. asyncio.shield() защищает от отмены именно на уровне корутины, но важно понимать последствия: вы можете создать ситуацию, где внутренние задачи продолжают жить дольше, чем вы ожидаете.
Применяйте shield осознанно: например, когда нужно дождаться безопасного освобождения ресурса или продолжить очистку.
Сценарий 1: стриминг — отмена потребителя и корректное завершение производителя
Стриминг (например, отправка/приём чанками) — один из самых частых сценариев, где неправильная отмена превращает систему в утечки задач.
Задача-производитель и задача-потребитель
Рассмотрим упрощённый пример: производитель генерирует куски данных и кладёт их в очередь, потребитель читает и отправляет дальше. При этом потребитель может остановиться раньше (например, клиент отвалился).
Критическое условие: при отмене потребителя нужно отменить производителя, иначе он продолжит заполнять очередь и расходовать ресурсы.
Корректный вариант
import asyncio
async def producer(queue: asyncio.Queue[int]):
i = 0
try:
while True:
await asyncio.sleep(0.2)
await queue.put(i)
i += 1
except asyncio.CancelledError:
# очистка: например, можно добавить логирование
print("producer: cancelled")
raise
async def consumer(queue: asyncio.Queue[int], stop_after: int):
received = 0
try:
while received < stop_after:
item = await queue.get()
print("consumer got:", item)
received += 1
except asyncio.CancelledError:
# здесь можно также сделать очистку
print("consumer: cancelled")
raise
async def main():
queue: asyncio.Queue[int] = asyncio.Queue(maxsize=10)
prod_task = asyncio.create_task(producer(queue))
try:
# потребитель сам завершится через stop_after
await consumer(queue, stop_after=5)
finally:
# гарантированно остановим производителя
prod_task.cancel()
try:
await prod_task
except asyncio.CancelledError:
pass
asyncio.run(main())
Что бывает, если сделать иначе
- Если вы просто перестали читать, не отменив производителя, он будет продолжать
putв очередь. Приmaxsizeэто может привести к зависанию производителя (он будет ждать место), а при отсутствииmaxsize— к росту памяти. - Если отмену ловить и «забыть», то
CancelledErrorперестаёт сигнализировать задаче корректно закончить работу.
Таймаут стриминга: таймаут на «тишину», а не на «общую длительность»
В стриминге часто важнее контролировать отсутствие данных. Например, сервер должен прислать хотя бы один чанку в течение T, иначе соединение считается зависшим.
Паттерн: таймаут вокруг queue.get() или ожидания следующего чанка.
import asyncio
import random
async def producer(queue: asyncio.Queue[str]):
# модель: иногда делает паузы подлиннее
for chunk in ["a", "b", "c", "d"]:
await asyncio.sleep(random.choice([0.1, 0.3, 1.5]))
await queue.put(chunk)
async def consumer_with_idle_timeout(queue: asyncio.Queue[str], idle_timeout: float):
while True:
try:
chunk = await asyncio.wait_for(queue.get(), timeout=idle_timeout)
except asyncio.TimeoutError:
print("idle timeout: no chunk for", idle_timeout, "seconds")
return
else:
print("got:", chunk)
async def main():
queue = asyncio.Queue(maxsize=10)
prod_task = asyncio.create_task(producer(queue))
try:
await consumer_with_idle_timeout(queue, idle_timeout=0.8)
finally:
prod_task.cancel()
try:
await prod_task
except asyncio.CancelledError:
pass
asyncio.run(main())
Сценарий 2: параллельные запросы — ограничение конкуренции и управляемые исключения
Параллельные запросы (gather, TaskGroup, семафоры) помогают ускорять I/O, но добавляют сложность: нужно решать, что считать успехом, что делать при исключениях и как корректно отменять «лишнее».
gather и семантика отмены
asyncio.gather(*coros) по умолчанию:
- запускает все,
- возвращает результаты,
- но если одна корутина падает исключением, остальные продолжат выполняться (пока не отменятся сами или не завершатся).
Это может быть не тем, что вы хотите в системах с дедлайнами. Обычно логика такая:
- если первый ответ пришёл — прекращаем остальные (race),
- если дедлайн истёк — отменяем всё,
- если упало ключевое — отменяем остальные, чтобы не тратить ресурсы.
Чтобы сделать это правильно, используйте TaskGroup (Python 3.11+) или явную отмену задач.
Паттерн: дедлайн на весь набор запросов + отмена остальных
Рассмотрим: нужно сходить к нескольким внешним источникам и взять первый успешный результат. Если дедлайн истёк — возвращаем ошибку.
Вариант через TaskGroup и FIRST_COMPLETED (race)
import asyncio
import time
async def query_source(name: str, delay: float, ok: bool) -> str:
await asyncio.sleep(delay)
if ok:
return f"{name}: OK"
raise RuntimeError(f"{name}: failed")
async def first_success(sources, deadline: float) -> str:
"""
deadline: время в монотонном отсчёте loop.time()
"""
loop = asyncio.get_running_loop()
tasks = []
async def remaining_timeout():
return max(0.0, deadline - loop.time())
async with asyncio.TaskGroup() as tg:
for name, delay, ok in sources:
async def run_one(n=name, d=delay, success=ok):
# ограничим каждую задачу оставшимся временем через sleep/await,
# но корректнее — отменить TaskGroup по дедлайну на уровне обёртки.
return await query_source(n, d, success)
# В TaskGroup нельзя напрямую сделать "race" с автоматической отменой,
# но мы можем дождаться результата через очередь/событие.
task = tg.create_task(run_one())
tasks.append(task)
# Ждём любой задаче успешного завершения.
done, pending = await asyncio.wait(tasks, timeout=remaining_timeout(), return_when=asyncio.FIRST_COMPLETED)
# Если в done нет успехов — ждём исключения/ошибки корректно.
# Здесь нужно аккуратно: FIRST_COMPLETED может завершиться исключением.
for t in done:
exc = t.exception()
if exc is None:
result = t.result()
# отменяем остальные "pending"
for p in pending:
p.cancel()
# выход из TaskGroup дождётся завершения отменённых задач
return result
# если все завершились с ошибкой/неуспехом — пусть исключение поднимется
# но если дедлайна не хватило — обработаем отдельно
if remaining_timeout() <= 0:
raise TimeoutError("deadline exceeded before any successful response")
# иначе: берём исключение первой завершившейся задачи
# (или можно собрать все)
raise list(done)[0].exception()
async def main():
loop = asyncio.get_running_loop()
deadline = loop.time() + 1.0
sources = [
("s1", 0.7, False),
("s2", 0.5, True),
("s3", 0.2, False),
]
try:
res = await first_success(sources, deadline)
print("result:", res)
except Exception as e:
print("error:", repr(e))
asyncio.run(main())
Нюанс: TaskGroup требует аккуратной структуры
TaskGroup отменяет «остальные» при исключении внутри группы (и наоборот — исключение может отменить всё). Поэтому для сложных схем race иногда удобнее:
- не рассчитывать на «гарантированное поведение»
TaskGroupпри первом исключении, - использовать явные
create_taskиcancel, - либо делать баланс через очередь результатов (успех/ошибка), не допуская «раннего вылета» группы.
На практике часто проще и контролируемее строить логику вне TaskGroup, используя asyncio.wait/очереди и ручную отмену.
Ограничение конкуренции: семафор вместо «создадим тысячу задач»
Если вы делаете 1000 одновременных запросов, даже если они асинхронные, это создаёт давление на:
- пул соединений,
- лимиты API,
- память (буферизация),
- планировщик событий.
Решение: семафор.
import asyncio
async def bounded_fetch(sem: asyncio.Semaphore, client, url: str):
async with sem:
return await client(url) # имитация I/O
async def run_all(urls):
sem = asyncio.Semaphore(20) # максимум 20 параллельных запросов
async def fake_client(url):
await asyncio.sleep(0.1)
return f"data:{url}"
tasks = [asyncio.create_task(bounded_fetch(sem, fake_client, u)) for u in urls]
return await asyncio.gather(*tasks, return_exceptions=True)
Сценарий 3: внешние вызовы с таймаутами и дедлайнами — «не накапливайте ожидание»
Самая дорогая ошибка в системах с несколькими слоями — это когда дедлайн теряется. Например:
- API-гейт отдаёт запрос сервису А с дедлайном 3 секунды.
- Сервис А делает вызов сервису B с таймаутом «по умолчанию 5 секунд».
- B вызывает сервис C с таймаутом «тоже по умолчанию 5 секунд».
И итог может превышать SLA на порядки, потому что каждый слой считает, что у него «ещё есть время».
Правильный подход: дедлайн «сверху вниз»
Подход:
- Сверху получить
deadline = loop.time() + ttl. - На каждом шаге вычислять
timeout = deadline - loop.time(). - Использовать
timeoutтолько на оставшееся время.
Общая утилита: дедлайн → таймаут для wait_for
import asyncio
from typing import Awaitable, TypeVar
T = TypeVar("T")
async def with_deadline(coro: Awaitable[T], deadline: float) -> T:
loop = asyncio.get_running_loop()
timeout = max(0.0, deadline - loop.time())
if timeout <= 0:
raise TimeoutError("deadline exceeded")
return await asyncio.wait_for(coro, timeout=timeout)
Теперь её можно использовать в любых местах.
Внешний вызов: отмена при таймауте + корректная очистка
Рассмотрим абстрактный клиент, который делает запрос и возвращает данные. Реальный код будет зависеть от используемой библиотеки (aiohttp, httpx, асинхронные клиенты SDK и т.д.), но логика одинакова: ограничить ожидание оставшимся временем и корректно закрыть ресурс.
import asyncio
async def external_call(name: str, duration: float) -> str:
# имитация I/O
await asyncio.sleep(duration)
return f"{name} result"
async def orchestrate(deadline: float):
try:
result = await asyncio.wait_for(
external_call("serviceB", duration=0.9),
timeout=max(0.0, deadline - asyncio.get_running_loop().time()),
)
return result
except asyncio.TimeoutError:
# Важно: здесь не обязательно подавлять отмену.
# TimeoutError ≠ CancelledError. Это сигнал таймера.
raise TimeoutError("orchestrate: timed out before deadline")
async def main():
loop = asyncio.get_running_loop()
deadline = loop.time() + 0.5
try:
print(await orchestrate(deadline))
except Exception as e:
print("error:", e)
asyncio.run(main())
Типичная ловушка: «поглощение» отмены в except
Частая анти-практика:
try:
await asyncio.wait_for(..., timeout=...)
except Exception:
# проглотили всё, включая CancelledError
...
Почему это плохо:
- если задачу отменили (
task.cancel()), вы потеряете корректный сигнал отмены; - приложение будет продолжать жить в неконсистентном состоянии.
Правильнее:
try:
await asyncio.wait_for(..., timeout=...)
except asyncio.CancelledError:
raise # отмену не подавляем
except asyncio.TimeoutError:
...
Сборка: «правильный» шаблон оркестратора с дедлайном, параллельностью и отменой
Теперь соберём все идеи в один практический каркас. Предположим, есть сценарий:
- Снаружи приходит запрос с дедлайном (например, 2.5 секунды до истечения SLA).
- Нужно параллельно:
- подтянуть метаданные,
- загрузить данные,
- посчитать что-то локально (или вызвать второй сервис).
- Если один из шагов не успел — мы прекращаем всё, возвращаем ошибку и освобождаем ресурсы.
Каркас
Комментарии
Пока нет комментариев