Spec-Zone.ru › Python 3.14

Корутины и задачи

В этом разделе описаны высокоуровневые API asyncio для работы с корутинами и задачами.

  • Корутины
  • Объекты, допускающие ожидание
  • Создание задач
  • Отмена задач
  • Группы задач
  • Ожидание
  • Параллельный запуск задач
  • Фабрика для немедленного запуска задач
  • Защита от отмены
  • Тайм-ауты
  • Примитивы ожидания
  • Запуск в потоках
  • Планирование из других потоков
  • Интроспекция
  • Объект Task

Корутины

Исходный код: Lib/asyncio/coroutines.py

Корутины, объявленные с помощью синтаксиса async/await, — предпочтительный способ написания приложений asyncio. Например, следующий фрагмент кода выводит «hello», ждёт 1 секунду, а затем выводит «world»:

>>> import asyncio

>>> async def main():
...     print('hello')
...     await asyncio.sleep(1)
...     print('world')

>>> asyncio.run(main())
hello
world

Обратите внимание, что простой вызов корутины не планирует её выполнение:

>>> main()
<coroutine object main at 0x1053bb7c8>

Для фактического запуска корутины asyncio предоставляет следующие механизмы:

  • Функция asyncio.run() для запуска функции «main()», являющейся точкой входа верхнего уровня (см. пример выше).
  • Ожидание корутины. Следующий фрагмент кода после ожидания в течение 1 секунды выведет «hello», а затем, подождав ещё 2 секунды, выведет «world»:

    import asyncio
    import time
    
    async def say_after(delay, what):
        await asyncio.sleep(delay)
        print(what)
    
    async def main():
        print(f"started at {time.strftime('%X')}")
    
        await say_after(1, 'hello')
        await say_after(2, 'world')
    
        print(f"finished at {time.strftime('%X')}")
    
    asyncio.run(main())
    

    Ожидаемый результат:

    started at 17:13:52
    hello
    world
    finished at 17:13:55
    
  • Функция asyncio.create_task() для параллельного выполнения корутин в виде Tasks asyncio.

    Изменим предыдущий пример и запустим две корутины say_after параллельно:

    async def main():
        task1 = asyncio.create_task(
            say_after(1, 'hello'))
    
        task2 = asyncio.create_task(
            say_after(2, 'world'))
    
        print(f"started at {time.strftime('%X')}")
    
        # Wait until both tasks are completed (should take
        # around 2 seconds.)
        await task1
        await task2
    
        print(f"finished at {time.strftime('%X')}")
    

    Обратите внимание, что ожидаемый результат теперь показывает, что фрагмент выполняется на 1 секунду быстрее, чем раньше:

    started at 17:14:32
    hello
    world
    finished at 17:14:34
    
  • Класс asyncio.TaskGroup предоставляет более современную альтернативу create_task(). При использовании этого API последний пример будет выглядеть так:

    async def main():
        async with asyncio.TaskGroup() as tg:
            task1 = tg.create_task(
                say_after(1, 'hello'))
    
            task2 = tg.create_task(
                say_after(2, 'world'))
    
            print(f"started at {time.strftime('%X')}")
    
        # The await is implicit when the context manager exits.
    
        print(f"finished at {time.strftime('%X')}")
    

    Время выполнения и результат должны совпадать с предыдущей версией.

    Добавлено в версии 3.11: asyncio.TaskGroup.

Объекты, допускающие ожидание

Объект называется объектом, допускающим ожидание, если его можно использовать в выражении await. Многие API asyncio предназначены для работы с такими объектами.

Существует три основных типа объектов, допускающих ожидание: корутины, задачи и объекты Future.

Корутины

Корутины Python допускают ожидание, поэтому их можно ожидать из других корутин:

import asyncio

async def nested():
    return 42

async def main():
    # Nothing happens if we just call "nested()".
    # A coroutine object is created but not awaited,
    # so it *won't run at all*.
    nested()  # will raise a "RuntimeWarning".

    # Let's do it differently now and await it:
    print(await nested())  # will print "42".

asyncio.run(main())

Важно

В этой документации термин «корутина» может обозначать два тесно связанных понятия:

  • функция-корутина: функция async def;
  • объект-корутина: объект, возвращаемый при вызове функции-корутины.

Задачи

Задачи используются для параллельного планирования корутин.

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

import asyncio

async def nested():
    return 42

async def main():
    # Schedule nested() to run soon concurrently
    # with "main()".
    task = asyncio.create_task(nested())

    # "task" can now be used to cancel "nested()", or
    # can simply be awaited to wait until it is complete:
    await task

asyncio.run(main())

Объекты Future

Future — это специальный низкоуровневый объект, допускающий ожидание и представляющий результат, который будет получен в будущем в результате асинхронной операции.

Ожидание объекта Future означает, что корутина будет ждать, пока Future не будет обработан в другом месте.

Объекты Future в asyncio нужны для совместного использования кода, основанного на обратных вызовах, с синтаксисом async/await.

Обычно на уровне кода приложения нет необходимости создавать объекты Future.

Объекты Future, иногда предоставляемые библиотеками и некоторыми API asyncio, можно ожидать:

async def main():
    await function_that_returns_a_future_object()

    # this is also valid:
    await asyncio.gather(
        function_that_returns_a_future_object(),
        some_python_coroutine()
    )

Хороший пример низкоуровневой функции, возвращающей объект Future, — loop.run_in_executor().

Создание задач

Исходный код: Lib/asyncio/tasks.py

asyncio.create_task(coro, *, name=None, context=None, eager_start=None, **kwargs)

Оборачивает корутину coro coroutine в Task и планирует её выполнение. Возвращает объект Task.

Полная сигнатура функции в основном совпадает с сигнатурой конструктора (или фабрики) Task — все именованные аргументы этой функции передаются в этот интерфейс.

Необязательный именованный аргумент context позволяет указать пользовательский contextvars.Context, в котором будет выполняться coro. Если context не указан, создаётся копия текущего контекста.

Необязательный именованный аргумент eager_start позволяет указать, должна ли задача выполняться немедленно при вызове create_task или быть запланирована позже. Если eager_start не передан, используется режим, установленный с помощью loop.set_task_factory().

Задача выполняется в цикле, возвращённом get_running_loop(); если в текущем потоке нет работающего цикла, вызывается RuntimeError.

Примечание

asyncio.TaskGroup.create_task() — это новая альтернатива, использующая структурную конкурентность; она позволяет ожидать завершения группы связанных задач с надёжными гарантиями безопасности.

Важно

Сохраните ссылку на результат этой функции, чтобы задача не исчезла во время выполнения. Цикл событий хранит только слабые ссылки на задачи. Задача, на которую больше нигде нет ссылок, может быть собрана сборщиком мусора в любой момент, даже до завершения. Чтобы надёжно запускать фоновые задачи в режиме «запустить и забыть», соберите их в коллекцию:

background_tasks = set()

for i in range(10):
    task = asyncio.create_task(some_coro(param=i))

    # Add task to the set. This creates a strong reference.
    background_tasks.add(task)

    # To prevent keeping references to finished tasks forever,
    # make each task remove its own reference from the set after
    # completion:
    task.add_done_callback(background_tasks.discard)

Обратите внимание, что при таком подходе задачи никогда не ожидаются. Поэтому, если задача завершится с ошибкой, её исключение не будет получено, и при сборке задачи сборщик мусора asyncio выведет сообщение «Task exception was never retrieved». Чтобы избежать этого, используйте asyncio.TaskGroup, которая хранит сильную ссылку на каждую задачу, ожидает их завершения и передаёт их исключения дальше:

async with asyncio.TaskGroup() as tg:
    for i in range(10):
        tg.create_task(some_coro(param=i))

Добавлено в версии 3.7.

Изменено в версии 3.8: Добавлен параметр name.

Изменено в версии 3.11: Добавлен параметр context.

Изменено в версии 3.14: Параметр eager_start добавлен посредством передачи всех kwargs.

Отмена задач

Задачи можно легко и безопасно отменить. При отмене задачи при первой возможности в ней будет вызвано исключение asyncio.CancelledError.

Рекомендуется использовать в корутинах блоки try/finally для надёжного выполнения очистки. Если исключение asyncio.CancelledError перехватывается явно, обычно его следует повторно вызвать после завершения очистки. asyncio.CancelledError напрямую наследуется от BaseException, поэтому большинству кода не нужно учитывать это исключение.

Компоненты asyncio, обеспечивающие структурную конкурентность, например asyncio.TaskGroup и asyncio.timeout(), внутри используют отмену и могут работать некорректно, если корутина подавляет asyncio.CancelledError. Аналогично, пользовательскому коду обычно не следует вызывать uncancel. Однако, если подавление asyncio.CancelledError действительно необходимо, нужно также вызвать uncancel(), чтобы полностью сбросить состояние отмены.

Группы задач

Группы задач объединяют API создания задач с удобным и надёжным способом ожидания завершения всех задач в группе.

class asyncio.TaskGroup

Асинхронный менеджер контекста, содержащий группу задач. Задачи можно добавлять в группу с помощью create_task(). При выходе из менеджера контекста ожидается завершение всех задач.

Добавлено в версии 3.11.

create_task(coro, *, name=None, context=None, eager_start=None, **kwargs)

Создаёт задачу в этой группе задач. Сигнатура совпадает с сигнатурой asyncio.create_task(). Если группа задач неактивна (например, менеджер ещё не был активирован, группа уже завершила работу или находится в процессе завершения), указанная coro будет закрыта.

Изменено в версии 3.13: Закрывает указанную корутину, если группа задач неактивна.

Изменено в версии 3.14: Передаёт все kwargs в loop.create_task()

Пример:

async def main():
    async with asyncio.TaskGroup() as tg:
        task1 = tg.create_task(some_coro(...))
        task2 = tg.create_task(another_coro(...))
    print(f"Both tasks have completed now: {task1.result()}, {task2.result()}")

Инструкция async with будет ожидать завершения всех задач в группе. Во время ожидания в группу всё ещё можно добавлять новые задачи (например, передав tg одной из корутин и вызвав tg.create_task() в этой корутине). После завершения последней задачи и выхода из блока async with добавлять новые задачи в группу уже нельзя.

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

После завершения всех задач, если какие-либо из них завершились с исключением, отличным от asyncio.CancelledError, эти исключения объединяются в ExceptionGroup или BaseExceptionGroup (в зависимости от ситуации; см. документацию к ним), после чего это исключение вызывается.

Два базовых исключения обрабатываются особым образом: если какая-либо задача завершается с исключением KeyboardInterrupt или SystemExit, группа задач всё равно отменяет оставшиеся задачи и ожидает их завершения, но затем повторно вызывает исходное исключение KeyboardInterrupt или SystemExit вместо ExceptionGroup или BaseExceptionGroup.

Если тело инструкции async with завершается исключением (то есть __aexit__() вызывается при установленном исключении), это обрабатывается так же, как сбой одной из задач: оставшиеся задачи отменяются и затем ожидаются, а исключения, не связанные с отменой, объединяются в группу исключений и вызываются. Исключение, переданное в __aexit__(), также включается в группу исключений, если только это не asyncio.CancelledError. Как и в предыдущем абзаце, для KeyboardInterrupt и SystemExit применяется особая обработка.

Группы задач тщательно отделяют внутреннюю отмену, используемую для «пробуждения» их __aexit__(), от запросов на отмену задачи, в которой они выполняются, поступающих от других сторон. В частности, если одна группа задач синтаксически вложена в другую и в обеих одновременно возникает исключение в одной из дочерних задач, внутренняя группа задач обработает свои исключения, после чего внешняя группа получит ещё один запрос на отмену и обработает собственные исключения.

Если группа задач отменена извне и одновременно должна вызвать ExceptionGroup, она вызовет метод cancel() родительской задачи. Это гарантирует, что при следующем await будет вызвано исключение asyncio.CancelledError, и запрос на отмену не будет потерян.

Группы задач сохраняют счётчик отмен, возвращаемый asyncio.Task.cancelling().

Изменено в версии 3.13: Улучшена обработка одновременных внутренних и внешних запросов на отмену, а также обеспечено корректное сохранение счётчиков отмен.

Завершение группы задач

Хотя стандартная библиотека не поддерживает непосредственное завершение группы задач, этого можно добиться, добавив в группу задачу, вызывающую исключение, и игнорируя вызванное исключение:

import asyncio
from asyncio import TaskGroup

class TerminateTaskGroup(Exception):
    """Exception raised to terminate a task group."""

async def force_terminate_task_group():
    """Used to force termination of a task group."""
    raise TerminateTaskGroup()

async def job(task_id, sleep_time):
    print(f'Task {task_id}: start')
    await asyncio.sleep(sleep_time)
    print(f'Task {task_id}: done')

async def main():
    try:
        async with TaskGroup() as group:
            # spawn some tasks
            group.create_task(job(1, 0.5))
            group.create_task(job(2, 1.5))
            # sleep for 1 second
            await asyncio.sleep(1)
            # add an exception-raising task to force the group to terminate
            group.create_task(force_terminate_task_group())
    except* TerminateTaskGroup:
        pass

asyncio.run(main())

Ожидаемый результат:

Task 1: start
Task 2: start
Task 1: done

Ожидание

async asyncio.sleep(delay, result=None)

Приостанавливает выполнение на delay секунд.

Если задан параметр result, он возвращается вызывающему коду после завершения корутины.

Функция sleep() всегда приостанавливает текущую задачу, позволяя выполняться другим задачам.

Если установить задержку 0, будет использован оптимизированный способ передачи управления другим задачам. Это можно использовать в длительно выполняющихся функциях, чтобы не блокировать цикл событий на всё время выполнения вызова функции.

Пример корутины, которая в течение 5 секунд каждую секунду выводит текущую дату:

import asyncio
import datetime as dt

async def display_date():
    loop = asyncio.get_running_loop()
    end_time = loop.time() + 5.0
    while True:
        print(dt.datetime.now())
        if (loop.time() + 1.0) >= end_time:
            break
        await asyncio.sleep(1)

asyncio.run(display_date())

Изменено в версии 3.10: Удалён параметр loop.

Изменено в версии 3.13: Вызывает ValueError, если delay равен nan.

Параллельный запуск задач

awaitable asyncio.gather(*aws, return_exceptions=False)

Параллельно выполняет объекты, допускающие ожидание, входящие в последовательность aws.

Если какой-либо объект в aws является корутиной, она автоматически планируется как задача.

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

Если return_exceptions равно False (значение по умолчанию), первое возникшее исключение немедленно передаётся задаче, ожидающей gather(). Остальные объекты в последовательности aws не будут отменены и продолжат выполняться.

Если return_exceptions равно True, исключения обрабатываются так же, как успешные результаты, и добавляются в общий список результатов.

Если gather() отменяется, все переданные объекты, допускающие ожидание (которые ещё не завершились), также отменяются.

Если какая-либо задача или Future из последовательности aws отменяется, это обрабатывается так, как если бы она вызвала CancelledError; вызов gather() в этом случае не отменяется. Это нужно, чтобы отмена одной переданной задачи/Future не приводила к отмене других задач/Future.

Примечание

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

Пример:

import asyncio

async def factorial(name, number):
    f = 1
    for i in range(2, number + 1):
        print(f"Task {name}: Compute factorial({number}), currently i={i}...")
        await asyncio.sleep(1)
        f *= i
    print(f"Task {name}: factorial({number}) = {f}")
    return f

async def main():
    # Schedule three calls *concurrently*:
    L = await asyncio.gather(
        factorial("A", 2),
        factorial("B", 3),
        factorial("C", 4),
    )
    print(L)

asyncio.run(main())

# Expected output:
#
#     Task A: Compute factorial(2), currently i=2...
#     Task B: Compute factorial(3), currently i=2...
#     Task C: Compute factorial(4), currently i=2...
#     Task A: factorial(2) = 2
#     Task B: Compute factorial(3), currently i=3...
#     Task C: Compute factorial(4), currently i=3...
#     Task B: factorial(3) = 6
#     Task C: Compute factorial(4), currently i=4...
#     Task C: factorial(4) = 24
#     [2, 6, 24]

Примечание

Если return_exceptions равно false, отмена gather() после того, как вызов уже помечен как завершённый, не отменит ни один из переданных объектов, допускающих ожидание. Например, gather может быть помечен как завершённый после передачи исключения вызывающему коду, поэтому вызов gather.cancel() после перехвата исключения (вызванного одним из объектов, допускающих ожидание) из gather не отменит остальные такие объекты.

Изменено в версии 3.7: Если отменяется сам gather, запрос на отмену передаётся дальше независимо от значения return_exceptions.

Изменено в версии 3.10: Удалён параметр loop.

Устарело с версии 3.10: Предупреждение об устаревании выводится, если не указаны позиционные аргументы либо не все позиционные аргументы являются объектами, подобными Future, и при этом нет работающего цикла событий.

Фабрика для немедленного запуска задач

asyncio.eager_task_factory(loop, coro, *, name=None, context=None)

Фабрика задач для их немедленного выполнения.

При использовании этой фабрики (через loop.set_task_factory(asyncio.eager_task_factory)) корутины начинают выполняться синхронно во время создания Task. Задачи планируются в цикле событий только в том случае, если они блокируются. Это может повысить производительность, поскольку для корутин, завершающихся синхронно, не требуется планирование циклом.

Типичный пример, когда это полезно, — корутины, использующие кэширование или мемоизацию, чтобы по возможности избежать фактического ввода-вывода.

Примечание

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

Добавлено в версии 3.12.

asyncio.create_eager_task_factory(custom_task_constructor)

Создаёт фабрику задач для немедленного выполнения, аналогичную eager_task_factory(), которая при создании новой задачи использует предоставленный custom_task_constructor вместо стандартного Task.

custom_task_constructor должен быть вызываемым объектом с сигнатурой, соответствующей сигнатуре Task.__init__. Вызываемый объект должен возвращать объект, совместимый с asyncio.Task.

Эта функция возвращает вызываемый объект, предназначенный для использования в качестве фабрики задач цикла событий через loop.set_task_factory(factory)).

Добавлено в версии 3.12.

Защита от отмены

awaitable asyncio.shield(aw)

Защищает объект, допускающий ожидание, от cancelled.

Если aw является корутиной, она автоматически планируется как задача.

Инструкция:

task = asyncio.create_task(something())
res = await shield(task)

эквивалентна:

res = await something()

за исключением того, что если содержащая её корутина отменяется, задача, выполняющаяся в something(), не отменяется. С точки зрения something() отмены не произошло. При этом её вызывающая сторона всё равно отменяется, поэтому выражение «await» по-прежнему вызывает CancelledError.

Если something() отменяется другим способом (то есть изнутри себя), это также отменит shield().

Если требуется полностью игнорировать отмену (не рекомендуется), функцию shield() следует использовать вместе с конструкцией try/except, как показано ниже:

task = asyncio.create_task(something())
try:
    res = await shield(task)
except CancelledError:
    res = None

Важно

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

Изменено в версии 3.10: Удалён параметр loop.

Устарело с версии 3.10: Предупреждение об устаревании выводится, если aw не является объектом, подобным Future, и при этом нет работающего цикла событий.

Тайм-ауты

asyncio.timeout(delay)

Возвращает асинхронный менеджер контекста, который можно использовать, чтобы ограничить время ожидания.

delay может быть равен None или представлять собой число секунд ожидания с плавающей точкой или целое число. Если delay равен None, ограничение по времени не применяется; это может быть полезно, если при создании менеджера контекста время ожидания неизвестно.

В обоих случаях после создания менеджера контекста можно изменить время ожидания с помощью Timeout.reschedule().

Пример:

async def main():
    async with asyncio.timeout(10):
        await long_running_task()

Если выполнение long_running_task занимает более 10 секунд, менеджер контекста отменит текущую задачу и обработает возникшее исключение asyncio.CancelledError внутри себя, преобразовав его в исключение TimeoutError, которое можно перехватить и обработать.

Примечание

Менеджер контекста asyncio.timeout() преобразует исключение asyncio.CancelledError в TimeoutError, поэтому исключение TimeoutError можно перехватить только за пределами менеджера контекста.

Пример перехвата TimeoutError:

async def main():
    try:
        async with asyncio.timeout(10):
            await long_running_task()
    except TimeoutError:
        print("The long operation timed out, but we've handled it.")

    print("This statement will run regardless.")

Менеджер контекста, созданный с помощью asyncio.timeout(), можно перенастроить на другой срок и проверить его состояние.

class asyncio.Timeout(when)

Асинхронный менеджер контекста для отмены корутин, превысивших отведённое время.

Предпочтительнее использовать asyncio.timeout() или asyncio.timeout_at(), а не создавать Timeout напрямую.

when должно задавать абсолютное время, в которое должен истечь срок действия контекста, согласно часам цикла событий:

  • Если when равно None, тайм-аут никогда не сработает.
  • Если when < loop.time(), тайм-аут сработает на следующей итерации цикла событий.
when() → float | None

Возвращает текущий срок или None, если срок не задан.

reschedule(when: float | None)

Изменяет срок тайм-аута.

expired() → bool

Возвращает значение, указывающее, превысил ли менеджер контекста установленный срок (истёк ли тайм-аут).

Пример:

async def main():
    try:
        # We do not know the timeout when starting, so we pass ``None``.
        async with asyncio.timeout(None) as cm:
            # We know the timeout now, so we reschedule it.
            new_deadline = get_running_loop().time() + 10
            cm.reschedule(new_deadline)

            await long_running_task()
    except TimeoutError:
        pass

    if cm.expired():
        print("Looks like we haven't finished on time.")

Менеджеры контекста тайм-аутов можно безопасно вкладывать друг в друга.

Добавлено в версии 3.11.

asyncio.timeout_at(when)

Аналогично asyncio.timeout(), но when задаёт абсолютное время прекращения ожидания или None.

Пример:

async def main():
    loop = get_running_loop()
    deadline = loop.time() + 20
    try:
        async with asyncio.timeout_at(deadline):
            await long_running_task()
    except TimeoutError:
        print("The long operation timed out, but we've handled it.")

    print("This statement will run regardless.")

Добавлено в версии 3.11.

async asyncio.wait_for(aw, timeout)

Ожидает завершения объекта, допускающего ожидание aw в течение заданного времени.

timeout может быть равен None или представлять собой число секунд ожидания с плавающей точкой или целое число. Если timeout равен None, ожидание продолжается до завершения future.

При истечении тайм-аута aw отменяется и вызывается исключение TimeoutError.

Чтобы предотвратить отмену aw, оберните его в shield().

Функция ожидает фактической отмены future, поэтому общее время ожидания может превысить timeout. Если во время отмены возникает исключение, оно передаётся вызывающему коду.

Если ожидание отменяется, future aw также отменяется.

Пример:

async def eternity():
    # Sleep for one hour
    await asyncio.sleep(3600)
    print('yay!')

async def main():
    # Wait for at most 1 second
    try:
        await asyncio.wait_for(eternity(), timeout=1.0)
    except TimeoutError:
        print('timeout!')

asyncio.run(main())

# Expected output:
#
#     timeout!

Изменено в версии 3.7: Если aw отменяется из-за тайм-аута, wait_for дожидается отмены aw. Ранее исключение TimeoutError вызывалось немедленно.

Изменено в версии 3.10: Удалён параметр loop.

Изменено в версии 3.11: Вызывается исключение TimeoutError вместо asyncio.TimeoutError.

Изменено в версии 3.12: Реализовано с помощью asyncio.timeout(); корутина, переданная в качестве aw, больше не оборачивается в Task, если значение timeout положительное.

Примитивы ожидания

async asyncio.wait(aws, *, timeout=None, return_when=ALL_COMPLETED)

Параллельно запускает экземпляры Future и Task из итерируемого объекта aws и блокирует выполнение до наступления условия, заданного параметром return_when.

Итерируемый объект aws не должен быть пустым.

Возвращает два множества задач/future: (done, pending).

Использование:

done, pending = await asyncio.wait(aws)

Параметр timeout (число с плавающей точкой или целое число), если он указан, задаёт максимальное количество секунд ожидания до возврата.

Обратите внимание: эта функция не вызывает исключение TimeoutError. Future или задачи, не завершившиеся к моменту истечения тайм-аута, просто возвращаются во втором множестве.

Параметр return_when задаёт условие возврата этой функции. Он должен принимать одно из следующих значений:

Константа

Описание

asyncio.FIRST_COMPLETED

Функция вернёт управление, когда любой future завершится или будет отменён.

asyncio.FIRST_EXCEPTION

Функция вернёт управление, когда выполнение любого future завершится с исключением. Если ни один future не вызывает исключение, это значение эквивалентно ALL_COMPLETED.

asyncio.ALL_COMPLETED

Функция вернёт управление, когда все future завершатся или будут отменены.

В отличие от wait_for(), wait() не отменяет future при истечении тайм-аута.

Если wait() отменяется, future в aws не отменяются и продолжают выполняться.

Изменено в версии 3.10: Удалён параметр loop.

Изменено в версии 3.11: Передавать объекты-корутины непосредственно в wait() запрещено.

Изменено в версии 3.12: Добавлена поддержка генераторов, выдающих задачи.

asyncio.as_completed(aws, *, timeout=None)

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

Объект, возвращённый as_completed(), можно перебирать как асинхронный итератор или обычный итератор. При асинхронном переборе изначально переданные объекты, допускающие ожидание, выдаются без изменений, если это задачи или future. Это позволяет легко сопоставлять ранее запланированные задачи с их результатами. Пример:

ipv4_connect = create_task(open_connection("127.0.0.1", 80))
ipv6_connect = create_task(open_connection("::1", 80))
tasks = [ipv4_connect, ipv6_connect]

async for earliest_connect in as_completed(tasks):
    # earliest_connect is done. The result can be obtained by
    # awaiting it or calling earliest_connect.result()
    reader, writer = await earliest_connect

    if earliest_connect is ipv6_connect:
        print("IPv6 connection established.")
    else:
        print("IPv4 connection established.")

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

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

ipv4_connect = create_task(open_connection("127.0.0.1", 80))
ipv6_connect = create_task(open_connection("::1", 80))
tasks = [ipv4_connect, ipv6_connect]

for next_connect in as_completed(tasks):
    # next_connect is not one of the original task objects. It must be
    # awaited to obtain the result value or raise the exception of the
    # awaitable that finishes next.
    reader, writer = await next_connect

Если тайм-аут истекает до завершения всех объектов, допускающих ожидание, вызывается исключение TimeoutError. При асинхронном переборе оно вызывается в цикле async for, а при обычном переборе — корутинами, которые выдаются итератором.

as_completed() не отменяет задачи, выполняющие переданные объекты, допускающие ожидание: если истекает тайм-аут или перебор отменяется, оставшиеся задачи продолжают выполняться.

Изменено в версии 3.10: Удалён параметр loop.

Устарело с версии 3.10: Выдаётся предупреждение об устаревании, если не все объекты, допускающие ожидание, в итерируемом объекте aws являются объектами, подобными Future, и при этом отсутствует работающий цикл событий.

Изменено в версии 3.12: Добавлена поддержка генераторов, выдающих задачи.

Изменено в версии 3.13: Теперь результат можно использовать как асинхронный итератор или обычный итератор (ранее он был только обычным итератором).

Выполнение в потоках

async asyncio.to_thread(func, /, *args, **kwargs)

Асинхронно выполняет функцию func в отдельном потоке.

Все переданные этой функции *args и **kwargs передаются непосредственно в func. Кроме того, передаётся текущий contextvars.Context, благодаря чему переменные контекста из потока цикла событий доступны в отдельном потоке.

Возвращает корутину, которую можно ожидать, чтобы получить результат выполнения func.

Эта функция-корутина предназначена главным образом для выполнения функций/методов с интенсивным вводом-выводом, которые в противном случае блокировали бы цикл событий при выполнении в основном потоке. Например:

def blocking_io():
    print(f"start blocking_io at {time.strftime('%X')}")
    # Note that time.sleep() can be replaced with any blocking
    # IO-bound operation, such as file operations.
    time.sleep(1)
    print(f"blocking_io complete at {time.strftime('%X')}")

async def main():
    print(f"started main at {time.strftime('%X')}")

    await asyncio.gather(
        asyncio.to_thread(blocking_io),
        asyncio.sleep(1))

    print(f"finished main at {time.strftime('%X')}")


asyncio.run(main())

# Expected output:
#
# started main at 19:50:53
# start blocking_io at 19:50:53
# blocking_io complete at 19:50:54
# finished main at 19:50:54

Прямой вызов blocking_io() в любой корутине заблокирует цикл событий на время выполнения и добавит к времени работы ещё одну секунду. Вместо этого с помощью asyncio.to_thread() мы можем запустить функцию в отдельном потоке, не блокируя цикл событий.

Примечание

Из-за GIL asyncio.to_thread() обычно можно использовать только для того, чтобы функции с интенсивным вводом-выводом не блокировали выполнение. Однако в модулях расширения, освобождающих GIL, или альтернативных реализациях Python, в которых его нет, asyncio.to_thread() можно использовать и для функций с интенсивными вычислениями.

Добавлено в версии 3.9.

Планирование из других потоков

asyncio.run_coroutine_threadsafe(coro, loop)

Передаёт корутину указанному циклу событий. Потокобезопасна.

Возвращает concurrent.futures.Future, позволяющий дождаться результата из другого потока ОС.

Эта функция предназначена для вызова из потока ОС, отличного от того, в котором выполняется цикл событий. Пример:

def in_thread(loop: asyncio.AbstractEventLoop) -> None:
    # Run some blocking IO
    pathlib.Path("example.txt").write_text("hello world", encoding="utf8")

    # Create a coroutine
    coro = asyncio.sleep(1, result=3)

    # Submit the coroutine to a given loop
    future = asyncio.run_coroutine_threadsafe(coro, loop)

    # Wait for the result with an optional timeout argument
    assert future.result(timeout=2) == 3

async def amain() -> None:
    # Get the running loop
    loop = asyncio.get_running_loop()

    # Run something in a thread
    await asyncio.to_thread(in_thread, loop)

Можно также выполнить обратную операцию. Пример:

@contextlib.contextmanager
def loop_in_thread() -> Generator[asyncio.AbstractEventLoop]:
    loop_fut = concurrent.futures.Future[asyncio.AbstractEventLoop]()
    stop_event = asyncio.Event()

    async def main() -> None:
        loop_fut.set_result(asyncio.get_running_loop())
        await stop_event.wait()

    with concurrent.futures.ThreadPoolExecutor(1) as tpe:
        complete_fut = tpe.submit(asyncio.run, main())
        for fut in concurrent.futures.as_completed((loop_fut, complete_fut)):
            if fut is loop_fut:
                loop = loop_fut.result()
                try:
                    yield loop
                finally:
                    loop.call_soon_threadsafe(stop_event.set)
            else:
                fut.result()

# Create a loop in another thread
with loop_in_thread() as loop:
    # Create a coroutine
    coro = asyncio.sleep(1, result=3)

    # Submit the coroutine to a given loop
    future = asyncio.run_coroutine_threadsafe(coro, loop)

    # Wait for the result with an optional timeout argument
    assert future.result(timeout=2) == 3

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

try:
    result = future.result(timeout)
except TimeoutError:
    print('The coroutine took too long, cancelling the task...')
    future.cancel()
except Exception as exc:
    print(f'The coroutine raised an exception: {exc!r}')
else:
    print(f'The coroutine returned: {result!r}')

См. раздел документации конкурентность и многопоточность.

В отличие от других функций asyncio, для этой функции требуется явно передать аргумент loop.

Добавлено в версии 3.5.1.

Интроспекция

asyncio.current_task(loop=None)

Возвращает экземпляр текущей выполняющейся Task или None, если ни одна задача не выполняется.

Если loop равен None, для получения текущего цикла используется get_running_loop().

Добавлено в версии 3.7.

asyncio.all_tasks(loop=None)

Возвращает множество ещё не завершённых объектов Task, выполняемых циклом.

Если loop равен None, для получения текущего цикла используется get_running_loop().

Добавлено в версии 3.7.

asyncio.iscoroutine(obj)

Возвращает True, если obj является объектом-корутиной.

Добавлено в версии 3.4.

Объект Task

class asyncio.Task(coro, *, loop=None, name=None, context=None, eager_start=False)

Объект Future-like, который выполняет Python корутину. Не является потокобезопасным.

Задачи используются для выполнения корутин в циклах событий. Если корутина ожидает Future, задача приостанавливает выполнение корутины и ждёт завершения Future. Когда Future завершается, выполнение обёрнутой корутины возобновляется.

Циклы событий используют кооперативное планирование: цикл событий выполняет одну задачу за раз. Пока задача ожидает завершения Future, цикл событий выполняет другие задачи, обратные вызовы или операции ввода-вывода.

Для создания задач используйте высокоуровневую функцию asyncio.create_task() или низкоуровневые функции loop.create_task() и ensure_future(). Ручное создание экземпляров задач не рекомендуется.

Чтобы отменить выполняющуюся задачу, используйте метод cancel(). При его вызове в обёрнутую корутину будет выброшено исключение CancelledError. Если во время отмены корутина ожидает объект Future, этот объект Future будет отменён.

Метод cancelled() можно использовать, чтобы проверить, была ли задача отменена. Метод возвращает True, если обёрнутая корутина не подавила исключение CancelledError и действительно была отменена.

asyncio.Task наследует все API класса Future, кроме Future.set_result() и Future.set_exception().

Необязательный аргумент только по ключевому слову context позволяет указать пользовательский contextvars.Context, в котором будет выполняться coro. Если context не указан, задача копирует текущий контекст и затем выполняет корутину в скопированном контексте.

Необязательный аргумент только по ключевому слову eager_start позволяет начать выполнение asyncio.Task сразу при создании задачи. Если задано значение True и цикл событий запущен, задача немедленно начнёт выполнение корутины и продолжит его до первой блокировки корутины. Если корутина завершается или вызывает исключение, не блокируясь, задача завершится немедленно, не передаваясь на планирование циклу событий.

Задачи являются обобщёнными по типу возвращаемого значения обёрнутых корутин.

Изменено в версии 3.7: Добавлена поддержка модуля contextvars.

Изменено в версии 3.8: Добавлен параметр name.

Устарело с версии 3.10: Выдаётся предупреждение об устаревании, если loop не указан и цикл событий не запущен.

Изменено в версии 3.11: Добавлен параметр context.

Изменено в версии 3.12: Добавлен параметр eager_start.

done()

Возвращает True, если задача завершена.

Задача считается завершённой, если обёрнутая корутина вернула значение, вызвала исключение или задача была отменена.

result()

Возвращает результат задачи.

Если задача завершена, возвращается результат обёрнутой корутины (или повторно вызывается исключение, если корутина вызвала его).

Если задача была отменена, этот метод вызывает исключение CancelledError.

Если результат задачи ещё недоступен, этот метод вызывает исключение InvalidStateError.

exception()

Возвращает исключение задачи.

Если обёрнутая корутина вызвала исключение, оно возвращается. Если обёрнутая корутина завершилась обычным образом, этот метод возвращает None.

Если задача была отменена, этот метод вызывает исключение CancelledError.

Если задача ещё не завершена, этот метод вызывает исключение InvalidStateError.

add_done_callback(callback, *, context=None)

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

Этот метод следует использовать только в низкоуровневом коде на основе обратных вызовов.

Подробнее см. документацию по Future.add_done_callback().

remove_done_callback(callback)

Удаляет callback из списка обратных вызовов.

Этот метод следует использовать только в низкоуровневом коде на основе обратных вызовов.

Подробнее см. документацию по Future.remove_done_callback().

get_stack(*, limit=None)

Возвращает список кадров стека для этой задачи.

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

Кадры всегда упорядочены от самых старых к самым новым.

Для приостановленной корутины возвращается только один кадр стека.

Необязательный аргумент limit задаёт максимальное количество возвращаемых кадров; по умолчанию возвращаются все доступные кадры. Порядок элементов в возвращаемом списке зависит от того, возвращается стек или трассировка: для стека возвращаются самые новые кадры, а для трассировки — самые старые. (Это соответствует поведению модуля traceback.)

print_stack(*, limit=None, file=None)

Выводит стек или трассировку для этой задачи.

Для кадров, полученных с помощью get_stack(), формируется вывод, аналогичный выводу модуля traceback.

Аргумент limit напрямую передаётся в get_stack().

Аргумент file — это поток ввода-вывода, в который записывается вывод; по умолчанию вывод записывается в sys.stdout.

get_coro()

Возвращает объект корутины, обёрнутый в Task.

Примечание

Для задач, которые уже завершились немедленно, будет возвращено None. См. Фабрика задач с немедленным запуском.

Добавлено в версии 3.8.

Изменено в версии 3.12: Добавленное выполнение задач с немедленным запуском означает, что результат может быть None.

get_context()

Возвращает объект contextvars.Context, связанный с задачей.

Добавлено в версии 3.12.

get_name()

Возвращает имя задачи.

Если задаче явно не присвоено имя, реализация asyncio Task по умолчанию создаёт имя при создании экземпляра.

Добавлено в версии 3.8.

set_name(value)

Задаёт имя задачи.

Аргумент value может быть любым объектом, который затем преобразуется в строку.

В реализации Task по умолчанию имя будет отображаться в выводе repr() для объекта задачи.

Добавлено в версии 3.8.

cancel(msg=None)

Запрашивает отмену задачи.

Если задача уже завершена или отменена, возвращает False, в противном случае возвращает True.

Метод обеспечивает выброс исключения CancelledError в обёрнутую корутину на следующем цикле цикла событий.

После этого корутина может выполнить очистку или даже отклонить запрос, подавив исключение с помощью блока try … … except CancelledError … finally. Поэтому, в отличие от Future.cancel(), Task.cancel() не гарантирует отмену задачи, хотя полное подавление отмены встречается редко и настоятельно не рекомендуется. Если корутина всё же решит подавить отмену, ей необходимо вызвать Task.uncancel() в дополнение к перехвату исключения.

Изменено в версии 3.9: Добавлен параметр msg.

Изменено в версии 3.11: Параметр msg передаётся от отменённой задачи ожидающей её корутине.

В следующем примере показано, как корутины могут перехватывать запрос на отмену:

async def cancel_me():
    print('cancel_me(): before sleep')

    try:
        # Wait for 1 hour
        await asyncio.sleep(3600)
    except asyncio.CancelledError:
        print('cancel_me(): cancel sleep')
        raise
    finally:
        print('cancel_me(): after sleep')

async def main():
    # Create a "cancel_me" Task
    task = asyncio.create_task(cancel_me())

    # Wait for 1 second
    await asyncio.sleep(1)

    task.cancel()
    try:
        await task
    except asyncio.CancelledError:
        print("main(): cancel_me is cancelled now")

asyncio.run(main())

# Expected output:
#
#     cancel_me(): before sleep
#     cancel_me(): cancel sleep
#     cancel_me(): after sleep
#     main(): cancel_me is cancelled now
cancelled()

Возвращает True, если задача отменена.

Задача считается отменённой, если отмена была запрошена с помощью cancel() и обёрнутая корутина передала дальше исключение CancelledError, выброшенное в неё.

uncancel()

Уменьшает счётчик запросов на отмену этой задачи.

Возвращает оставшееся количество запросов на отмену.

Обратите внимание: после завершения выполнения отменённой задачи дальнейшие вызовы uncancel() не имеют эффекта.

Добавлено в версии 3.11.

Этот метод используется внутренними механизмами asyncio и не предназначен для кода конечных пользователей. В частности, если отмену задачи удалось успешно снять, это позволяет таким элементам структурированного параллелизма, как группы задач и asyncio.timeout(), продолжить выполнение, ограничивая действие отмены соответствующим структурированным блоком. Например:

async def make_request_with_timeout():
    try:
        async with asyncio.timeout(1):
            # Structured block affected by the timeout:
            await make_request()
            await make_another_request()
    except TimeoutError:
        log("There was a timeout")
    # Outer code not affected by the timeout:
    await unrelated_code()

Хотя блок с make_request() и make_another_request() может быть отменён из-за тайм-аута, unrelated_code() должен продолжить выполнение даже при наступлении тайм-аута. Это реализовано с помощью uncancel(). Менеджеры контекста TaskGroup используют uncancel() аналогичным образом.

Если код конечного пользователя по какой-либо причине подавляет отмену, перехватывая CancelledError, ему необходимо вызвать этот метод, чтобы сбросить состояние отмены.

Когда этот метод уменьшает счётчик отмены до нуля, он проверяет, не подготовил ли предыдущий вызов cancel() выброс исключения CancelledError в задачу. Если исключение ещё не было выброшено, этот план отменяется (сбросом внутреннего флага _must_cancel).

Изменено в версии 3.13: При достижении нуля теперь отменяются ожидающие запросы на отмену.

cancelling()

Возвращает количество ожидающих запросов на отмену этой задачи, то есть количество вызовов cancel() за вычетом количества вызовов uncancel().

Обратите внимание: если это число больше нуля, но задача всё ещё выполняется, cancelled() всё равно вернёт False. Это связано с тем, что число можно уменьшить вызовом uncancel(); если количество запросов на отмену снизится до нуля, задача в итоге может не быть отменена.

Этот метод используется внутренними механизмами asyncio и не предназначен для кода конечных пользователей. Подробнее см. uncancel().

Добавлено в версии 3.11.

© 2001 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.14/library/asyncio-task.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API