Привет, Хабр! Асинхронный код на питоне пишут почти все, а про отмену задач знают на удивление мало. Оно и понятно: пока всё работает, отмена не всплывает. Всплывает она в момент, когда сервис не хочет останавливаться, таймаут не срабатывает, транзакция остаётся открытой, а в логах — CancelledError без всякого контекста.

Отмена в asyncio устроена не совсем так, как ожидаешь. Это не «убить задачу», а «попросить задачу закончиться, бросив в неё исключение в следующей точке ожидания». Из этой формулировки следует довольно много неочевидных последствий.

Разберём механику, а потом пройдёмся по местам, где она ломает код.

Как отмена работает физически

task.cancel() не останавливает задачу. Он ставит флаг и планирует бросить CancelledError внутрь корутины при следующем проходе цикла событий.

async def worker():
    print("начали")
    await asyncio.sleep(10)
    print("это не напечатается")

task = asyncio.create_task(worker())
await asyncio.sleep(0.1)
task.cancel()

Исключение прилетит в точку await asyncio.sleep(10) — именно туда, где корутина отдала управление. Если бы там был вызов к базе, исключение прилетело бы в него.

Из этого сразу следует главное ограничение: задача без единого await неотменяема.

async def busy():
    total = 0
    for i in range(10_000_000):
        total += i          # ни одной точки приостановки
    return total

cancel() на такой задаче не сделает ничего. Цикл событий не получит управление, пока цикл не досчитает, а к тому моменту задача уже завершится сама. То же с блокирующим вызовом — time.sleep(30), синхронный запрос через requests, чтение большого файла. Всё это не отменяется и попутно вешает весь цикл событий.

Проверить, что задача действительно отменилась, а не просто получила запрос, можно только дождавшись её:

task.cancel()
try:
    await task
except asyncio.CancelledError:
    pass

Без этого await вы не знаете, добралась отмена до задачи или нет.

CancelledError не наследуется от Exception

С версии 3.8 CancelledError наследуется напрямую от BaseException.

async def handler():
    try:
        await do_work()
    except Exception as e:          # отмену НЕ поймает
        log.error("упало: %s", e)

Это сделано специально. Отмена — не ошибка приложения, а служебный сигнал, и глотать его вместе с обычными исключениями нельзя.

А вот так делать нельзя почти никогда:

try:
    await do_work()
except BaseException:
    log.error("что-то случилось")     # проглотили отмену

Задача, съевшая отмену и продолжившая работать — источник самых мутных багов. Тот, кто её отменял, ждёт завершения и не дожидается. Сервис не останавливается, а из логов непонятно, почему.

Если отмену действительно надо обработать — что-то дописать, откатить, залогировать, её положено пробросить дальше:

try:
    await do_work()
except asyncio.CancelledError:
    await cleanup()
    raise                    # обязательно

Без raise вы врёте вызывающему коду о том, что задача отменилась.

await внутри finally

Отдельная проблема, которая проявляется только при отмене.

async def process(conn):
    try:
        await conn.execute("BEGIN")
        await do_something(conn)
    finally:
        await conn.execute("ROLLBACK")     # а если нас уже отменяют?

Задачу отменили, исключение прилетело в do_something, управление ушло в finally. Там снова await — и если отмену запросят повторно, ROLLBACK тоже прервётся. Соединение останется в подвешенном состоянии.

Защита от повторной отмены на время уборки — asyncio.shield, либо в более общем виде — выполнение критичной части вне зоны отмены:

finally:
    await asyncio.shield(conn.execute("ROLLBACK"))

shield защищает внутреннюю операцию: отмена придёт снаружи, но обёрнутая задача доработает до конца. Внешний await при этом получит CancelledError сразу, то есть вы узнаете об отмене, а операция всё равно завершится.

Пользоваться этим надо аккуратно. Заэкранированная задача продолжает жить, и если про неё забыть, вы получите работающий кусок кода, за которым никто не следит.

gather и брошенные задачи

Классика, которая живёт в половине кодовых баз.

results = await asyncio.gather(fetch_a(), fetch_b(), fetch_c())

Если fetch_b() упадёт, gather немедленно пробросит исключение наружу. А fetch_a() и fetch_c() при этом продолжат выполняться. Никто их не отменяет и никто их результат не заберёт.

Дальше сценарии зависят от того, что там внутри. В лучшем случае задача просто доработает вхолостую. В худшем — она держит соединение, пишет в базу или ждёт ответа, а исключение из неё вылезет позже и в другом месте, с сообщением Task exception was never retrieved.

Флаг return_exceptions=True меняет поведение: gather дождётся всех и вернёт исключения в списке вместо того, чтобы их бросать.

results = await asyncio.gather(*coros, return_exceptions=True)
for r in results:
    if isinstance(r, Exception):
        log.error("одна из задач упала: %s", r)

Это фиксит потерянные задачи, но не фиксит другое: gather не отменяет остальных, когда одна упала. Если вам нужно именно «упало одно — сворачиваемся все», нужна другая конструкция.

TaskGroup

С версии 3.11 в стандартной библиотеке есть структурная конкурентность.

async with asyncio.TaskGroup() as tg:
    a = tg.create_task(fetch_a())
    b = tg.create_task(fetch_b())
    c = tg.create_task(fetch_c())

# сюда попадаем, только когда все три закончились
print(a.result(), b.result(), c.result())

Правила такие. Выйти из блока async with, пока внутри есть незавершённые задачи, невозможно. Если любая из задач падает, остальные отменяются, группа дожидается их завершения (давая отработать уборке в finally) и только потом пробрасывает исключение наружу.

Брошенных задач тут не бывает по построению — в этом весь смысл.

Исключение приезжает не одно, а завёрнутым в ExceptionGroup, потому что упасть могло несколько задач сразу. Разбирается это синтаксисом except*:

try:
    async with asyncio.TaskGroup() as tg:
        tg.create_task(fetch_a())
        tg.create_task(fetch_b())
except* ConnectionError as eg:
    for exc in eg.exceptions:
        log.error("сеть: %s", exc)
except* ValueError as eg:
    log.error("данные: %s", eg.exceptions)

Каждый except* ловит все исключения своего типа из группы. Если в группе оказались и ConnectionError, и ValueError, сработают оба блока — в отличие от обычного try/except, где отработал бы только первый подходящий.

Одно неудобство при переезде с gather: если у вас снаружи стоит except ConnectionError, после перехода на TaskGroup он перестанет ловить, потому что теперь прилетает группа.

Ограничение параллельности внутри группы

Частый вопрос при переезде: TaskGroup запускает всё сразу, а надо не больше N одновременно. Встроенного ограничителя нет, делается семафором:

sem = asyncio.Semaphore(10)

async def limited(coro):
    async with sem:
        return await coro

async with asyncio.TaskGroup() as tg:
    tasks = [tg.create_task(limited(fetch(url))) for url in urls]

Задачи создаются все, но реально работают по десять — остальные висят на семафоре. Отмена при этом работает нормально: ожидание семафора — обычная точка await.

Если задач тысячи, создавать их все сразу расточительно даже с семафором — каждая занимает память. Тогда лучше очередь и фиксированное число воркеров:

async def worker(queue: asyncio.Queue):
    while True:
        url = await queue.get()
        try:
            await fetch(url)
        finally:
            queue.task_done()

async with asyncio.TaskGroup() as tg:
    queue = asyncio.Queue(maxsize=100)
    workers = [tg.create_task(worker(queue)) for _ in range(10)]

    for url in urls:
        await queue.put(url)
    await queue.join()

    for w in workers:
        w.cancel()

Воркеры крутятся в бесконечном цикле и сами не завершатся, поэтому их надо отменить явно, иначе TaskGroup будет ждать их вечно.

timeout

До 3.11 таймауты делали через wait_for, и это работало, но выглядело не очень, когда нужно было ограничить целый блок кода. Сейчас есть контекстный менеджер:

async with asyncio.timeout(5):
    data = await fetch(url)
    parsed = await parse(data)
    await save(parsed)

По истечении времени в текущую задачу летит отмена, а на выходе из блока она превращается в TimeoutError. То есть снаружи вы ловите нормальную ошибку таймаута, а не CancelledError.

Есть вариант с абсолютным временем:

deadline = loop.time() + 30

async with asyncio.timeout_at(deadline):
    await step_one()

async with asyncio.timeout_at(deadline):
    await step_two()

Оба этапа укладываются в общие тридцать секунд, а не по тридцать каждый.

Таймаут можно двигать изнутри блока:

async with asyncio.timeout(5) as cm:
    await first_part()
    cm.reschedule(loop.time() + 20)     # дальше нужен запас побольше
    await second_part()

Работает всё это на той же механике отмены, а значит, наследует её ограничения. Блокирующий вызов внутри asyncio.timeout не прервётся — таймаут сработает только когда управление вернётся в цикл событий.

Задачи, которые исчезают

Неочевидный источник багов: цикл событий держит на задачи только слабые ссылки.

async def main():
    for item in items:
        asyncio.create_task(process(item))     # ссылку никто не сохранил
    await asyncio.sleep(60)

Сборщик мусора имеет полное право собрать такую задачу в любой момент, в том числе посреди выполнения. Обработается не всё, ошибок при этом не будет, а воспроизводится оно нестабильно — зависит от того, когда сборщик решит поработать.

Лечится сохранением ссылок:

background = set()

def spawn(coro):
    task = asyncio.create_task(coro)
    background.add(task)
    task.add_done_callback(background.discard)
    return task

add_done_callback тут нужен, чтобы множество не росло бесконечно.

Если задачи связаны логически и должны жить вместе — лучше TaskGroup, он держит ссылки сам.

uncancel и вложенные таймауты

С 3.11 у задачи появился счётчик отмен. Каждый cancel() его увеличивает, uncancel() — уменьшает.

Нужно это для вложенных зон отмены. Представьте: внешний timeout(30), внутри него внутренний timeout(5). Внутренний сработал, отменил задачу, поймал отмену и превратил её в TimeoutError. Но как ему понять, что отмена была его собственная, а не пришла снаружи от внешнего таймаута?

По счётчику. Внутренний менеджер запоминает значение до отмены и сверяет после.

async def worker():
    try:
        async with asyncio.timeout(5):
            await slow_operation()
    except TimeoutError:
        log.warning("не уложились, продолжаем без этих данных")
        await fallback()

Здесь после срабатывания внутреннего таймаута задача продолжает работать — потому что asyncio.timeout вызвал uncancel() и вернул задачу в нормальное состояние.

Если вы пишете свой контекстный менеджер с отменой, эту механику придётся воспроизводить руками:

task = asyncio.current_task()
task.cancel()
# ...
if task.uncancel() == 0 and self._my_cancellation:
    raise TimeoutError
raise      # отмена была не наша, пробрасываем

Метод cancelling() возвращает текущее значение счётчика, если нужно только посмотреть.

Асинхронные контекстные менеджеры и генераторы

Пара мест, где отмена ведёт себя не так, как ожидаешь.

Асинхронный контекстный менеджер при отмене получает исключение в aexit, и там действуют те же правила, что для finally:

class Connection:
    async def __aexit__(self, exc_type, exc, tb):
        await self._close()          # может быть прервано следующей отменой

Если закрытие критично, его тоже надо экранировать.

С асинхронными генераторами сложнее. Они не закрываются в момент, когда вы перестали их читать, закрытие происходит при сборке мусора, асинхронно, через хук цикла событий.

async def read_rows(conn):
    async with conn.transaction():
        async for row in conn.cursor("SELECT ..."):
            yield row

async for row in read_rows(conn):
    if row.id == target:
        break            # генератор остался приостановленным внутри транзакции

После break генератор висит на yield внутри открытой транзакции. Закроется он неизвестно когда — на усмотрение сборщика мусора. При отмене всей задачи картина та же.

Надёжный способ — закрывать явно:

from contextlib import aclosing

async with aclosing(read_rows(conn)) as rows:
    async for row in rows:
        if row.id == target:
            break

aclosing вызовет aclose() на выходе из блока, генератор получит GeneratorExit в точке yield, отработает свой finally и корректно закроет транзакцию.

Блокирующий вызов, который вешает всё

Возвращаясь к тому, с чего начали. Отмена работает только через точки await, а блокирующий код таких точек не создаёт.

async def handler():
    data = requests.get(url).json()      # блокирует весь цикл событий
    return data

Пока requests ждёт ответа, ни одна другая корутина не выполняется. Не отменяется, не обрабатывается и не таймаутится. При медленном внешнем сервисе весь сервис встаёт.

Способ два. Первый — асинхронный клиент (httpx, aiohttp). Второй, если библиотека синхронная и альтернативы нет, — вынести в поток:

data = await asyncio.to_thread(blocking_call, arg)

Цикл событий при этом освобождается. Но отменить такую операцию всё равно нельзя: to_thread вернёт управление при отмене, а сам вызов в потоке продолжит выполняться до конца. Потоки в питоне не прерываются извне.

Найти такие места помогает режим отладки:

asyncio.run(main(), debug=True)

В нём цикл событий логирует все колбэки, выполнявшиеся дольше 100 миллисекунд. Порог настраивается через loop.slow_callback_duration.

Отмена и внешние ресурсы

Запрос к базе через асинхронный драйвер отменяется корректно с точки зрения питона — исключение прилетает в await. Но сам запрос на сервере продолжает выполняться. Postgres о вашей отмене ничего не знает: он получил SELECT, он его считает.

async with asyncio.timeout(2):
    rows = await conn.fetch("SELECT * FROM huge_table")

Через две секунды вы получите TimeoutError и пойдёте дальше, а база будет ещё минуту гонять запрос вхолостую. При потоке таймаутящихся запросов вы просто набьёте базу мёртвой работой.

Драйверы обычно умеют это решать — asyncpg, например, при отмене отправляет запрос на отмену на стороне сервера. Но проверить это для своего драйвера стоит, а лучше — выставить таймаут на стороне базы, чтобы не полагаться на клиента:

SET statement_timeout = '2s';

С HTTP та же история: отмена закрывает соединение, а обработка запроса на чужой стороне продолжается. Если запрос был не идемпотентным — например, создавал заказ, вы не знаете, создался он или нет. Отмена не означает, что действие не произошло.

Для таких мест нужны либо ключи идемпотентности, либо явная проверка состояния после таймаута.

Что смотреть, когда сервис не останавливается

Нажали Ctrl+C или отправили сигнал контейнеру, а процесс продолжает висеть.

Первое, что стоит сделать — посмотреть, какие задачи вообще живы:

for task in asyncio.all_tasks():
    print(task.get_name(), task.get_coro())
    task.print_stack()

print_stack покажет, где именно задача сейчас стоит. Обычно этого достаточно, чтобы найти виноватого.

Второе — именовать задачи при создании, иначе в выводе будет Task-17 и никакой полезной информации:

asyncio.create_task(process(order), name=f"process-order-{order.id}")

Третье — корректно завершать при остановке. Шаблон примерно такой:

async def shutdown():
    tasks = [t for t in asyncio.all_tasks() if t is not asyncio.current_task()]
    for t in tasks:
        t.cancel()
    await asyncio.gather(*tasks, return_exceptions=True)

return_exceptions=True тут обязателен: половина задач завершится с CancelledError, и без флага gather бросит первое же исключение, не дождавшись остальных.

Как это тестировать

Тесты на отмену пишут редко, а зря — большая часть описанных проблем ловится довольно просто.

Проверка, что задача корректно прибирает за собой:

async def test_cleanup_on_cancel():
    released = False

    async def worker():
        nonlocal released
        try:
            await asyncio.sleep(100)
        finally:
            released = True

    task = asyncio.create_task(worker())
    await asyncio.sleep(0)          # даём задаче стартовать
    task.cancel()

    with pytest.raises(asyncio.CancelledError):
        await task

    assert released

Без await asyncio.sleep(0) задача может не успеть начаться, и отмена придёт до входа в try. Это же и есть, кстати, штатный способ отдать управление циклу событий на один проход.

Проверка, что отмену не глотают:

async def test_does_not_swallow_cancellation():
    task = asyncio.create_task(handler())
    await asyncio.sleep(0.01)
    task.cancel()

    with pytest.raises(asyncio.CancelledError):
        await task

Если внутри handler стоит except BaseException без проброса, тест упадёт — задача завершится нормально вместо отмены.

Проверять таймауты через реальное ожидание не стоит: тесты станут медленными и начнут мигать на нагруженном CI. Вместо этого либо подменяют время цикла событий, либо тестируют не сам таймаут, а реакцию на него — просто отменяя задачу вручную.


Если коротко, о чём тут всё

Отмена в asyncio кооперативная. Никто ничего не убивает: задаче бросают исключение в ближайшей точке ожидания, а дальше всё зависит от того, как она с ним обошлась.

Задача без await не отменяется вовсе. CancelledError не ловится через except Exception, и это правильно, а вот except BaseException его глотает и создаёт задачи-зомби. Уборка в finally сама по себе может быть прервана следующей отменой.

Из инструментов на сегодня почти всё закрывается двумя: TaskGroup вместо gather и asyncio.timeout вместо wait_for. Первый не оставляет брошенных задач, второй нормально работает с вложенностью.

Чтобы проверить, насколько уверенно вы ориентируетесь в Python за пределами asyncio, можно пройти вступительный тест продвинутого курса. Он поможет оценить текущий уровень и найти темы, которые стоит подтянуть.

Когда отмена задач, таймауты и блокирующие вызовы начинают влиять на стабильность сервиса, полезно расширить картину и разобраться, как Python распределяет работу между корутинами, потоками и процессами. На бесплатных уроках можно проверить своё понимание и задать вопросы преподавателю:

  • 4 августа в 20:00. «Многозадачность в Python: асинхронность, процессы, потоки». Записаться

  • 20 августа в 20:00. «Docker для Python-разработчика». Записаться