Очереди
Исходный код: 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, который извлекает недавно добавленные записи в первую очередь (последний вошел, первый вышел).
Исключения
-
exception asyncio.QueueEmpty -
Это исключение возникает, когда метод
get_nowait()вызывается на пустой очереди.
-
exception asyncio.QueueFull -
Исключение, возникающее при вызове метода
put_nowait()на очереди, достигшей своего maxsize.
Примеры
Очереди могут использоваться для распределения рабочей нагрузки между несколькими concurrent задачами:
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.9/library/asyncio-queue.html