Spec-Zone.ru › Python 3.13

Очереди

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

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

Обратите внимание, что методы очередей asyncio не имеют параметра timeout; используйте функцию asyncio.wait_for() для выполнения операций с очередью с указанием таймаута.

См. также раздел Примеры ниже.

Очередь

class asyncio.Queue(maxsize=0)

Очередь «первый пришёл, первый ушёл» (FIFO).

Если maxsize меньше или равно нулю, размер очереди бесконечен. Если это целое число, большее 0, то await put() блокируется, когда очередь достигает maxsize, до тех пор, пока элемент не будет удалён методом get().

В отличие от потоковой очереди стандартной библиотеки queue, размер очереди всегда известен и может быть возвращён вызовом метода qsize().

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

Этот класс не потокобезопасен.

maxsize

Количество элементов, разрешённых в очереди.

empty()

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

full()

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

Если очередь была инициализирована значением maxsize=0 (по умолчанию), то full() никогда не возвращает True.

coroutine get()

Удалить и вернуть элемент из очереди. Если очередь пуста, дождаться, пока элемент не станет доступным.

Вызывает QueueShutDown, если очередь закрыта и пуста или если очередь закрыта немедленно.

get_nowait()

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

coroutine join()

Заблокировать, пока все элементы в очереди не будут получены и обработаны.

Счёт незавершенных задач увеличивается при добавлении элемента в очередь. Счёт уменьшается, когда сопрограмма-потребитель вызывает task_done() для указания того, что элемент был извлечен и все работы с ним завершены. Когда счёт незавершенных задач снизится до нуля, join() разблокируется.

coroutine put(item)

Поместить элемент в очередь. Если очередь полна, подождать, пока не освободится место, прежде чем добавить элемент.

Вызывает QueueShutDown, если очередь закрыта.

put_nowait(item)

Поместить элемент в очередь без блокировки.

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

qsize()

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

shutdown(immediate=False)

Закрыть очередь, заставляя get() и put() вызывать QueueShutDown.

По умолчанию, get() для закрытой очереди будет вызывать исключение только после того, как очередь станет пустой. Установите immediate в true, чтобы get() вызвало исключение сразу же.

Все заблокированные вызывающие стороны put() и get() будут разблокированы. Если immediate равно true, для каждого оставшегося элемента в очереди задача будет помечена как завершенная, что может разблокировать вызывающих сторон join().

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

task_done()

Указать, что задача, ранее добавленная в очередь, завершена.

Используется потребителями очереди. Для каждого вызова get(), используемого для извлечения задачи, последующий вызов task_done() сообщает очереди, что обработка задачи завершена.

Если join() в данный момент заблокирован, он возобновит работу, когда все элементы будут обработаны (то есть, что вызов task_done() был получен для каждого элемента, который был добавлен в очередь с помощью put()).

shutdown(immediate=True) вызывает task_done() для каждого оставшегося элемента в очереди.

Вызывает ValueError, если вызов производится большее число раз, чем элементов было помещено в очередь.

Очередь с приоритетами

class asyncio.PriorityQueue

Вариант Queue; извлекает элементы в порядке приоритета (наименьший в первую очередь).

Элементы обычно представляют собой кортежи вида (priority_number, data).

Очередь LIFO

class asyncio.LifoQueue

Вариант Queue, который извлекает элементы, добавленные последними (последний пришёл, первый ушёл).

Исключения

exception asyncio.QueueEmpty

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

exception asyncio.QueueFull

Исключение, возникающее при вызове метода put_nowait() для очереди, достигшей maxsize.

exception asyncio.QueueShutDown

Исключение, возникающее при вызове put() или get() для очереди, которая была закрыта.

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

Примеры

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

import asyncio
import random
import time


async def worker(name, queue):
    while True:
        # Get a "work item" out of the queue.
        sleep_for = await queue.get()

        # Sleep for the "sleep_for" seconds.
        await asyncio.sleep(sleep_for)

        # Notify the queue that the "work item" has been processed.
        queue.task_done()

        print(f'{name} has slept for {sleep_for:.2f} seconds')


async def main():
    # Create a queue that we will use to store our "workload".
    queue = asyncio.Queue()

    # Generate random timings and put them into the queue.
    total_sleep_time = 0
    for _ in range(20):
        sleep_for = random.uniform(0.05, 1.0)
        total_sleep_time += sleep_for
        queue.put_nowait(sleep_for)

    # Create three worker tasks to process the queue concurrently.
    tasks = []
    for i in range(3):
        task = asyncio.create_task(worker(f'worker-{i}', queue))
        tasks.append(task)

    # Wait until the queue is fully processed.
    started_at = time.monotonic()
    await queue.join()
    total_slept_for = time.monotonic() - started_at

    # Cancel our worker tasks.
    for task in tasks:
        task.cancel()
    # Wait until all worker tasks are cancelled.
    await asyncio.gather(*tasks, return_exceptions=True)

    print('====')
    print(f'3 workers slept in parallel for {total_slept_for:.2f} seconds')
    print(f'total expected sleep time: {total_sleep_time:.2f} seconds')


asyncio.run(main())

© 2001–2024 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.13/library/asyncio-queue.html

Spec-Zone.ru

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