asyncio: gather, create_task, Queue, timeout, отмена
В прошлом уроке мы запускали несколько корутин через gather. Теперь разберём более сложные паттерны: гибкое управление через Tasks, очереди для producer/consumer, таймауты, корректная отмена, обработка ошибок. Это инструменты, которые потребуются в реальном backend-коде.
gather - запуск с ожиданием всех
import asyncio
async def fetch(url):
await asyncio.sleep(1)
return f"data from {url}"
async def main():
results = await asyncio.gather(
fetch("a"),
fetch("b"),
fetch("c"),
)
print(results) # ['data from a', 'data from b', 'data from c']
asyncio.run(main())
Порядок результатов соответствует порядку аргументов. Все ждут параллельно, общее время ~max времени отдельной задачи.
return_exceptions для гибкой обработки
По умолчанию при ошибке в одной задаче gather отменяет остальные:
async def main():
try:
results = await asyncio.gather(
fetch("a"),
failing_fetch(),
fetch("c"),
)
except Exception as e:
print(f"Error: {e}")
# fetch_a и fetch_c могут быть в любом состоянии
С return_exceptions=True ошибки возвращаются как обычные значения:
results = await asyncio.gather(
fetch("a"),
failing_fetch(),
fetch("c"),
return_exceptions=True,
)
# results = ["data from a", SomeException(...), "data from c"]
for r in results:
if isinstance(r, Exception):
log.error(r)
else:
process(r)
TaskGroup (Python 3.11+)
С Python 3.11 появился TaskGroup - более структурированная альтернатива gather:
async def main():
async with asyncio.TaskGroup() as tg:
task1 = tg.create_task(fetch("a"))
task2 = tg.create_task(fetch("b"))
task3 = tg.create_task(fetch("c"))
# После выхода из async with все task завершены
print(task1.result())
print(task2.result())
При ошибке в одной задаче группа отменяет остальные и поднимает ExceptionGroup (тоже Py 3.11+) со всеми ошибками. Это надстройка над обычными исключениями, ловится отдельным синтаксисом except*. Это более явное и безопасное API.
create_task для фоновых задач
async def background_worker():
while True:
await asyncio.sleep(1)
print("tick")
async def main():
task = asyncio.create_task(background_worker())
# Делаем что-то ещё, worker работает в фоне
await asyncio.sleep(5)
task.cancel()
try:
await task
except asyncio.CancelledError:
print("worker cancelled")
asyncio.run(main())
create_task запускает корутину сразу. Можно сохранить ссылку для отмены или получения результата.
# Плохо - ссылка не сохранена
async def main():
asyncio.create_task(background_worker()) # может исчезнуть
# Хорошо
async def main():
task = asyncio.create_task(background_worker())
await asyncio.sleep(5)
task.cancel()
# Или сохранить в set
background_tasks = set()
async def main():
task = asyncio.create_task(background_worker())
background_tasks.add(task)
task.add_done_callback(background_tasks.discard)
wait_for - таймаут
async def slow_fetch():
await asyncio.sleep(10)
return "data"
async def main():
try:
result = await asyncio.wait_for(slow_fetch(), timeout=3)
except asyncio.TimeoutError:
print("Timeout!")
asyncio.run(main())
wait_for отменяет задачу при превышении таймаута и бросает TimeoutError. Это безопасный способ ограничить время IO-операций. В Go ту же роль играет context с дедлайном.
С Python 3.11 есть более универсальный asyncio.timeout():
async with asyncio.timeout(3):
result = await slow_fetch()
Контекстный менеджер выглядит чище и работает с любым кодом внутри.
wait - частичный сбор
asyncio.wait возвращает множества выполненных и невыполненных task:
async def main():
tasks = [asyncio.create_task(fetch(url)) for url in urls]
# Ждём пока хотя бы одна завершится
done, pending = await asyncio.wait(tasks, return_when=asyncio.FIRST_COMPLETED)
for task in done:
print(task.result())
# Отменяем остальные
for task in pending:
task.cancel()
Опции return_when:
ALL_COMPLETED(default) - ждём всехFIRST_COMPLETED- первойFIRST_EXCEPTION- первой ошибки
Полезно для гонок (use the fastest result), скан портов, fan-out paradigms.
as_completed - результаты по мере готовности
async def main():
tasks = [fetch(url) for url in urls]
for coro in asyncio.as_completed(tasks):
result = await coro
process(result)
# Обрабатываем результаты в порядке завершения, не запуска
Удобно когда хочется реагировать на быстрые ответы сразу, не ждать всех.
Queue для producer/consumer
import asyncio
async def producer(queue, n):
for i in range(n):
await asyncio.sleep(0.5)
await queue.put(f"item-{i}")
await queue.put(None) # сигнал конца
async def consumer(queue):
while True:
item = await queue.get()
if item is None:
break
print(f"Consumed: {item}")
queue.task_done()
async def main():
queue = asyncio.Queue(maxsize=10)
await asyncio.gather(
producer(queue, 5),
consumer(queue),
)
asyncio.run(main())
Queue - thread-safe-аналог для async. put блокирует если очередь полна (maxsize > 0), get блокирует если пуста. Идеально для разделения работы между producer и consumer корутинами. Именно на такой связке строятся консумеры брокеров - подробно в уроке про consumer patterns.
Другие очереди:
LifoQueue- LIFO (стек)PriorityQueue- с приоритетом
Отмена и CancelledError
async def task():
try:
while True:
print("working")
await asyncio.sleep(1)
except asyncio.CancelledError:
print("cleanup")
raise # ОБЫЧНО нужно re-raise
async def main():
t = asyncio.create_task(task())
await asyncio.sleep(3)
t.cancel()
try:
await t
except asyncio.CancelledError:
print("task was cancelled")
asyncio.run(main())
task.cancel() бросает CancelledError внутрь корутины при следующем await. Корутина может поймать для cleanup, но обычно должна re-raise чтобы asyncio мог корректно отслеживать отмену.
С Python 3.11+ asyncio.CancelledError стал подклассом BaseException, не Exception - чтобы случайно не поймать в except Exception.
shield - защита от отмены
async def critical_operation():
await save_to_db()
async def main():
try:
await asyncio.wait_for(
asyncio.shield(critical_operation()),
timeout=5
)
except asyncio.TimeoutError:
print("Timed out, but critical_operation продолжает")
shield защищает внутреннюю корутину от отмены извне. Полезно для операций которые должны завершиться даже если внешний таймаут истёк.
Lock и Semaphore
asyncio.Lock для критических секций:
shared_counter = 0
lock = asyncio.Lock()
async def increment():
global shared_counter
async with lock:
# критическая секция
current = shared_counter
await asyncio.sleep(0.1) # симулируем какое-то IO
shared_counter = current + 1
Хотя asyncio однопоточен, между await'ами другие корутины могут модифицировать shared state. Lock гарантирует exclusive доступ.
asyncio.Semaphore ограничивает количество одновременных операций:
sem = asyncio.Semaphore(10) # макс 10 одновременно
async def fetch_limited(url):
async with sem:
return await fetch(url)
async def main():
tasks = [fetch_limited(url) for url in URLs] # 1000 URLs
results = await asyncio.gather(*tasks)
# Не больше 10 запросов одновременно
Полезно для rate limiting, защиты внешних сервисов.
Event для координации
ready_event = asyncio.Event()
async def waiter():
print("waiting...")
await ready_event.wait()
print("ready!")
async def setter():
await asyncio.sleep(2)
ready_event.set()
await asyncio.gather(waiter(), setter())
Event это простой сигнал «случилось/не случилось». Несколько waiter могут ждать одного события.
Обработка ошибок
async def task():
try:
await do_work()
except SpecificError as e:
# Логируем и продолжаем
log.error(e)
return default_value
except Exception as e:
# Логируем и пробрасываем
log.exception("unexpected")
raise
Внутри асинхронной функции обработка ошибок работает как в синхронной. Главное помнить про cleanup в case CancelledError.
Производительность - tips
-
Используй async-библиотеки:
aiohttpвместоrequests,asyncpgвместоpsycopg2,aiofilesдля файлов. -
Бенчмарк прежде чем оптимизировать: asyncio даёт прирост только для IO-bound. Для CPU-bound лучше multiprocessing.
-
Ограничивай concurrency через Semaphore: тысячи одновременных запросов могут перегрузить downstream сервисы.
-
Не забывай await: легко забыть и получить unhandled coroutine warning.
-
Используй gather с return_exceptions для устойчивости: одна ошибка не должна валить весь batch.
Распространённые ошибки
1. Создать Task и потерять ссылку
Уже обсуждали - сохранять в set с auto-cleanup через done_callback.
2. Блокирующий код в async
async def bad():
requests.get("https://api.example.com") # БЛОКИРУЕТ loop
Используй aiohttp или loop.run_in_executor для wrapping sync вызовов.
3. await в цикле вместо gather
# Плохо - последовательно
async def slow():
results = []
for url in urls:
result = await fetch(url) # ждём каждый отдельно
results.append(result)
# Хорошо - параллельно
async def fast():
results = await asyncio.gather(*[fetch(url) for url in urls])
4. Незаявленный return в async
async def task():
await asyncio.sleep(1)
# забыли return - вернёт None
Это не специфично для async, но в async коде чаще встречается.
5. Отмена без cleanup
async def task():
file = open("data.txt")
await something()
file.close() # НЕ выполнится если task отменена
# Правильно - context manager или try/finally
async def task():
file = open("data.txt")
try:
await something()
finally:
file.close()
Сравнение с Go
В Go параллельные запросы через горутины:
results := make([]string, 3)
var wg sync.WaitGroup
for i, url := range urls {
wg.Add(1)
go func(i int, url string) {
defer wg.Done()
results[i] = fetch(url)
}(i, url)
}
wg.Wait()
asyncio проще для IO-bound, но Go имеет настоящую параллельность через горутины - не ограничен одним CPU как Python asyncio.
Мини-задание
- Параллельная обработка с timeout:
import asyncio
async def slow_task(name, duration):
await asyncio.sleep(duration)
return f"{name} done"
async def main():
try:
result = await asyncio.wait_for(
slow_task("A", 5),
timeout=2
)
print(result)
except asyncio.TimeoutError:
print("Too slow")
asyncio.run(main())
- Producer/Consumer через Queue:
import asyncio
async def producer(queue):
for i in range(5):
await asyncio.sleep(0.5)
await queue.put(i)
print(f"Produced {i}")
await queue.put(None)
async def consumer(queue):
while True:
item = await queue.get()
if item is None:
break
print(f"Consumed {item}")
await asyncio.sleep(0.3)
async def main():
queue = asyncio.Queue(maxsize=3)
await asyncio.gather(producer(queue), consumer(queue))
asyncio.run(main())
- Semaphore для rate-limiting:
import asyncio
async def fetch(url, semaphore):
async with semaphore:
print(f"Fetching {url}")
await asyncio.sleep(1)
return f"data from {url}"
async def main():
semaphore = asyncio.Semaphore(3) # максимум 3 одновременно
urls = [f"url-{i}" for i in range(10)]
tasks = [fetch(url, semaphore) for url in urls]
results = await asyncio.gather(*tasks)
print(f"Got {len(results)} results")
asyncio.run(main())
Что дальше
Освоили продвинутые паттерны asyncio. В следующем уроке - GIL и threading: когда asyncio не подходит, и как использовать потоки правильно.