Spec-Zone.ru › Python 3.7

queue — Класс синхронизированной очереди

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

Модуль queue реализует очереди с множественными производителями и потребителями. Он особенно полезен в многопоточной программировании, когда информация должна безопасно обмениваться между несколькими потоками. Класс Queue в этом модуле реализует все необходимые семантики блокировки. Он зависит от наличия поддержки потоков в Python; см. модуль threading.

Модуль реализует три типа очереди, которые различаются только порядком извлечения элементов. В очереди FIFO первыми извлекаются задачи, добавленные первыми. В очереди LIFO первой извлекается последняя добавленная запись (работает как стек). В очереди с приоритетами записи хранятся в отсортированном виде (используя модуль heapq), и первой извлекается запись с наименьшим значением.

Внутренне эти три типа очередей используют блокировки для временного блокирования конкурирующих потоков; однако они не предназначены для обработки повторного входа в один и тот же поток.

Кроме того, модуль реализует тип «простой» очереди FIFO, SimpleQueue, чья конкретная реализация обеспечивает дополнительные гарантии в обмен на меньшие функциональные возможности.

Модуль queue определяет следующие классы и исключения:

class queue.Queue(maxsize=0)

Конструктор очереди FIFO. maxsize — целое число, которое устанавливает верхнюю границу количества элементов, которые могут быть помещены в очередь. Вставка будет заблокирована, когда этот размер будет достигнут, до тех пор, пока элементы очереди не будут извлечены. Если maxsize меньше или равно нулю, размер очереди бесконечен.

class queue.LifoQueue(maxsize=0)

Конструктор очереди LIFO. maxsize — целое число, которое устанавливает верхнюю границу количества элементов, которые могут быть помещены в очередь. Вставка будет заблокирована, когда этот размер будет достигнут, до тех пор, пока элементы очереди не будут извлечены. Если maxsize меньше или равно нулю, размер очереди бесконечен.

class queue.PriorityQueue(maxsize=0)

Конструктор очереди с приоритетами. maxsize — целое число, которое устанавливает верхнюю границу количества элементов, которые могут быть помещены в очередь. Вставка будет заблокирована, когда этот размер будет достигнут, до тех пор, пока элементы очереди не будут извлечены. Если maxsize меньше или равно нулю, размер очереди бесконечен.

Первыми извлекаются записи с наименьшим значением (запись с наименьшим значением — та, которая возвращается sorted(list(entries))[0]). Типичный шаблон для записей — кортеж в форме: (priority_number, data).

Если элементы data не сравнимы, данные можно обернуть в класс, который игнорирует элемент данных и сравнивает только число приоритета:

from dataclasses import dataclass, field
from typing import Any

@dataclass(order=True)
class PrioritizedItem:
    priority: int
    item: Any=field(compare=False)
class queue.SimpleQueue

Конструктор неограниченной очереди FIFO. Простые очереди лишены расширенной функциональности, такой как отслеживание задач.

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

exception queue.Empty

Исключение, генерируемое при вызове неблокирующего get() (или get_nowait()) для объекта Queue, который пуст.

exception queue.Full

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

Объекты очереди

Объекты очереди (Queue, LifoQueue или PriorityQueue) предоставляют описанные ниже публичные методы.

Queue.qsize()

Возвращает приблизительный размер очереди. Обратите внимание, что qsize() > 0 не гарантирует, что последующий вызов get() не заблокируется, и qsize() < maxsize не гарантирует, что put() не заблокируется.

Queue.empty()

Возвращает True, если очередь пуста, False в противном случае. Если empty() возвращает True, это не гарантирует, что последующий вызов put() не заблокируется. Аналогично, если empty() возвращает False, это не гарантирует, что последующий вызов get() не заблокируется.

Queue.full()

Возвращает True, если очередь заполнена, False в противном случае. Если full() возвращает True, это не гарантирует, что последующий вызов get() не заблокируется. Аналогично, если full() возвращает False, это не гарантирует, что последующий вызов put() не заблокируется.

Queue.put(item, block=True, timeout=None)

Поместить item в очередь. Если необязательный аргумент block равен true и timeout равен None (по умолчанию), блокировать, если необходимо, до тех пор, пока не появится свободное место. Если timeout — положительное число, блокировать не более timeout секунд и вызывать исключение Full, если свободное место не появилось в течение этого времени. В противном случае (block ложь), поместить элемент в очередь, если свободное место сразу доступно, иначе вызывать исключение Full (timeout игнорируется в этом случае).

Queue.put_nowait(item)

Эквивалентно put(item, False).

Queue.get(block=True, timeout=None)

Удалить и вернуть элемент из очереди. Если необязательные аргументы block истинны и timeout равен None (по умолчанию), блокировать, если необходимо, до тех пор, пока элемент не станет доступным. Если timeout — положительное число, блокировать не более timeout секунд и вызывать исключение Empty, если элемент не был доступен в течение этого времени. В противном случае (block ложь), вернуть элемент, если он доступен сразу, иначе вызывать исключение Empty (timeout игнорируется в этом случае).

До версии 3.0 на системах POSIX и для всех версий на Windows, если block истинно, а timeout — None, эта операция переходит в непрерываемый ожидание на базовой блокировке. Это означает, что исключения не могут возникнуть, и, в частности, SIGINT не вызовет KeyboardInterrupt.

Queue.get_nowait()

Эквивалентно get(False).

Два метода предложены для отслеживания того, были ли задачи, запущенные в очередь, полностью обработаны демонами потребительских потоков.

Queue.task_done()

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

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

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

Queue.join()

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

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

Пример того, как дождаться завершения задач в очереди:

def worker():
    while True:
        item = q.get()
        if item is None:
            break
        do_work(item)
        q.task_done()

q = queue.Queue()
threads = []
for i in range(num_worker_threads):
    t = threading.Thread(target=worker)
    t.start()
    threads.append(t)

for item in source():
    q.put(item)

# block until all tasks are done
q.join()

# stop workers
for i in range(num_worker_threads):
    q.put(None)
for t in threads:
    t.join()

Объекты SimpleQueue

SimpleQueue объекты предоставляют описанные ниже публичные методы.

SimpleQueue.qsize()

Возвращает приблизительный размер очереди. Обратите внимание, что qsize() > 0 не гарантирует, что последующий вызов get() не заблокируется.

SimpleQueue.empty()

Возвращает True, если очередь пуста, False в противном случае. Если empty() возвращает False, это не гарантирует, что последующий вызов get() не заблокируется.

SimpleQueue.put(item, block=True, timeout=None)

Поместить item в очередь. Метод никогда не блокируется и всегда выполняется успешно (кроме потенциальных низкоуровневых ошибок, таких как неудача при выделении памяти). Необязательные аргументы block и timeout игнорируются и предоставляются только для совместимости с Queue.put().

Деталь реализации CPython: Этот метод имеет реализацию на C, которая является реентерабельной. То есть вызов put() или get() может быть прерван другим вызовом put() в той же потоке без возникновения тупика или повреждения внутреннего состояния очереди. Это делает его подходящим для использования в деструкторах, таких как методы __del__ или weakref обратные вызовы.

SimpleQueue.put_nowait(item)

Эквивалентно put(item), предоставлено для совместимости с Queue.put_nowait().

SimpleQueue.get(block=True, timeout=None)

Удалить и вернуть элемент из очереди. Если необязательный аргумент block имеет значение true и timeout равен None (по умолчанию), заблокировать, если необходимо, до тех пор, пока элемент не станет доступен. Если timeout — положительное число, то блокировка выполняется не более timeout секунд, а при отсутствии элемента в течение этого времени генерируется исключение Empty. В противном случае (block равен false), вернуть элемент, если он сразу доступен, иначе сгенерировать исключение Empty (timeout игнорируется в этом случае).

SimpleQueue.get_nowait()

Эквивалентно get(False).

См. также

Class multiprocessing.Queue

Класс очереди для использования в многопроцессорной (а не многопоточной) среде.

collections.deque — альтернативная реализация неограниченных очередей с быстрыми атомарными операциями append() и popleft(), которые не требуют блокировки.

© 2001–2020 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.7/library/queue.html

Spec-Zone.ru

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