Spec-Zone.ru › Python 3.8

Очереди

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

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

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

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

Очередь

class asyncio.Queue(maxsize=0, *, loop=None)

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

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

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

Устаревшее с версии 3.8, будет удалено в версии 3.10: Параметр loop.

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

maxsize

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

empty()

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

full()

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

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

coroutine get()

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

get_nowait()

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

coroutine join()

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

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

coroutine put(item)

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

put_nowait(item)

Добавляет элемент в очередь без блокировки.

Если свободное место немедленно недоступно, генерирует исключение QueueFull.

qsize()

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

task_done()

Указывает, что завершена ранее помещённая в очередь задача.

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

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

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

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

class asyncio.PriorityQueue

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

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

Очередь LIFO

class asyncio.LifoQueue

Вариант Queue, который извлекает записи, добавленные последними, в первую очередь (LIFO).

Исключения

exception asyncio.QueueEmpty

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

exception asyncio.QueueFull

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

Примеры

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

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–2022 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.8/library/asyncio-queue.html

Spec-Zone.ru

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