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 запускает корутину сразу. Можно сохранить ссылку для отмены или получения результата.

asyncio только слабо ссылается на Task. Если потерять сильную ссылку, Task может быть собран garbage collector до завершения, что приведёт к молчаливому исчезновению фоновой работы.
# Плохо - ссылка не сохранена
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

  1. Используй async-библиотеки: aiohttp вместо requests, asyncpg вместо psycopg2, aiofiles для файлов.

  2. Бенчмарк прежде чем оптимизировать: asyncio даёт прирост только для IO-bound. Для CPU-bound лучше multiprocessing.

  3. Ограничивай concurrency через Semaphore: тысячи одновременных запросов могут перегрузить downstream сервисы.

  4. Не забывай await: легко забыть и получить unhandled coroutine warning.

  5. Используй 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.

Мини-задание

  1. Параллельная обработка с 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())
  1. 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())
  1. 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 не подходит, и как использовать потоки правильно.

Зарегистрируйтесь бесплатно, чтобы пройти квиз, решить задание с автопроверкой и вести прогресс.