Очереди
Исходный код: Lib/asyncio/queues.py
Очереди asyncio предназначены для использования так же, как классы модуля queue. Хотя очереди asyncio не являются потокобезопасными, они предназначены специально для использования в коде с async/await.
Обратите внимание, что у методов очередей asyncio нет параметра timeout; чтобы выполнять операции с очередью с тайм-аутом, используйте функцию asyncio.wait_for().
См. также раздел Примеры ниже.
Queue
-
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.
-
async get() -
Удаляет и возвращает элемент из очереди. Если очередь пуста, ожидает появления элемента.
Вызывает исключение
QueueShutDown, если очередь была остановлена и пуста либо если очередь была остановлена немедленно.
-
get_nowait() -
Возвращает элемент, если он доступен немедленно, иначе вызывает исключение
QueueEmpty.Вызывает исключение
QueueShutDown, если очередь была остановлена и пуста.
-
async join() -
Блокируется, пока все элементы очереди не будут получены и обработаны.
Счётчик незавершённых задач увеличивается каждый раз, когда в очередь добавляется элемент. Счётчик уменьшается каждый раз, когда сопрограмма-потребитель вызывает
task_done(), указывая, что элемент был получен и вся работа с ним завершена. Когда счётчик незавершённых задач становится равен нулю, выполнениеjoin()разблокируется.
-
async put(item) -
Помещает элемент в очередь. Если очередь заполнена, ожидает появления свободного места, прежде чем добавить элемент.
Вызывает исключение
QueueShutDown, если очередь была остановлена.
-
put_nowait(item) -
Помещает элемент в очередь, не блокируясь.
Если свободного места нет, вызывает исключение
QueueFull.Вызывает исключение
QueueShutDown, если очередь была остановлена.
-
qsize() -
Возвращает количество элементов в очереди.
-
shutdown(immediate=False) -
Переводит экземпляр
Queueв режим остановки.Очередь больше не может увеличиваться. Последующие вызовы
put()вызывают исключениеQueueShutDown. Вызовыput(), которые в данный момент заблокированы, будут разблокированы и вызовутQueueShutDownв ожидающей задаче.Если immediate имеет значение false (по умолчанию), очередь можно штатно опустошить, вызвав
get(), чтобы извлечь уже помещённые в неё задачи.Если для каждой оставшейся задачи вызвать
task_done(), ожидающий вызовjoin()будет штатно разблокирован.Когда очередь опустеет, последующие вызовы
get()будут вызывать исключениеQueueShutDown.Если immediate имеет значение true, очередь немедленно прекращает работу. Очередь опустошается полностью, а счётчик незавершённых задач уменьшается на количество извлечённых задач. Если число незавершённых задач равно нулю, ожидающие вызовы
join()разблокируются. Кроме того, заблокированные вызовыget()разблокируются и вызывают исключениеQueueShutDown, поскольку очередь пуста.Будьте осторожны при использовании
join()со значением immediate, равным true. В этом случае вызов join разблокируется, даже если задачи не были обработаны, что нарушает обычный инвариант при ожидании завершения работы с очередью.Добавлено в версии 3.13.
-
task_done() -
Указывает, что ранее помещённый в очередь элемент работы обработан.
Используется потребителями очереди. Для каждого вызова
get(), использованного для получения элемента работы, последующий вызовtask_done()сообщает очереди, что обработка этого элемента завершена.Если вызов
join()в данный момент заблокирован, он возобновится после обработки всех элементов (то есть когда для каждого элемента, который былput()в очередь, будет получен вызов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(),put_nowait(),get()илиget_nowait().Добавлено в версии 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 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.14/library/asyncio-queue.html