Spec-Zone.ru › Python 3.13

Потоки и задачи

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

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

Потоки

Исходный код: 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()» — точки входа верхнего уровня (см. пример выше).
  • Ожидание потока. Следующий фрагмент кода выведет «hello» после ожидания 1 секунды, а затем выведет «world» после ожидания ещё 2 секунд:

    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() для одновременного выполнения потоков как asyncio Tasks.

    Давайте изменим пример выше и запустим два 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.

Объекты-awaitables

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

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

Потоки

Python-потоки являются объектами-awaitable и поэтому могут быть ожидаемыми другими потоками:

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 — специальный низкоуровневый объект-awaitable, представляющий конечный результат асинхронной операции.

Когда объект 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)

Обернуть coro поток в Task и запланировать его выполнение. Вернуть объект Task.

Если name не None, он устанавливается как имя задачи с помощью Task.set_name().

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

Задача выполняется в цикле, возвращаемом 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)

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

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

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

Отмена задач

Задачи можно легко и безопасно отменять. При отмене задачи 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)

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

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

Пример:

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() родительской задачи. Это гарантирует, что asyncio.CancelledError будет поднят в следующем await, чтобы отмена не была потеряна.

Группы задач сохраняют счётчик отмен, отчёт о котором даёт 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

Ожидание

coroutine asyncio.sleep(delay, result=None)

Ожидание в течение delay секунд.

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

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

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

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

import asyncio
import datetime

async def display_date():
    loop = asyncio.get_running_loop()
    end_time = loop.time() + 5.0
    while True:
        print(datetime.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() отменён, все поданные айтейблы (которые ещё не завершены) также отменяются.

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

Примечание

Альтернативным способом создания и выполнения задач одновременно и ожидания их завершения является 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)

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

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.

coroutine asyncio.wait_for(aw, timeout)

Ожидать завершения aw awaitable с таймаутом.

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

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

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

Чтобы избежать отмены задачи cancellation, оберните её в shield().

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

Если ожидание отменено, будущее 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.

Ожидающие примитивы

coroutine 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 (значение типа float или int), если указано, может использоваться для контроля максимального времени ожидания в секундах перед возвратом.

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

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

Константа

Описание

asyncio.FIRST_COMPLETED

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

asyncio.FIRST_EXCEPTION

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

asyncio.ALL_COMPLETED

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

В отличие от wait_for(), wait() не отменяет объекты Future при истечении времени ожидания.

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

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

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

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

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

Возвращаемый объектом as_completed() можно итерировать как асинхронный итератор или обычный итератор. При использовании асинхронной итерации, исходные объекты awaitables возвращаются, если это задачи или объекты 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.")

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

При использовании в качестве обычного итератора, каждая итерация возвращает новую корутину, которая возвращает результат или вызывает исключение следующего завершённого объекта awaitable. Эта схема совместима с версиями 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 генерируется, если истекает время ожидания до завершения всех awaitables. Оно генерируется в цикле async for во время асинхронной итерации или корутинами, возвращаемыми при обычной итерации.

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

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

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

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

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

coroutine 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() в любой корутине заблокирует цикл событий на всё время его выполнения, что приведёт к дополнительному 1 секунде времени выполнения. Вместо этого, используя asyncio.to_thread(), мы можем запустить его в отдельном потоке, не блокируя цикл событий.

Примечание

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

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

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

asyncio.run_coroutine_threadsafe(coro, loop)

Передать корутину в заданный цикл событий. Потокобезопасная функция.

Возвращает concurrent.futures.Future для ожидания результата из другого потока операционной системы.

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

# 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) == 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 не указан, используется get_running_loop() для получения текущего цикла.

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

asyncio.all_tasks(loop=None)

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

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

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

asyncio.iscoroutine(obj)

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

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

END_OF_DOCUMENT_MARKER

Объект задачи

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)

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

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

Аргумент 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.

END_OF_DOCUMENT_MARKER
cancel(msg=None)

Запросить отмену задачи.

Это организует исключение 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–2024 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.13/library/asyncio-task.html

Spec-Zone.ru

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