Spec-Zone.ru › Python 3.8

threading — Параллелизм на основе потоков

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

Этот модуль строит интерфейсы потоков высокого уровня поверх модуля низкого уровня _thread. Также см. модуль queue.

Изменено в версии 3.7: Этот модуль раньше был необязательным, теперь он всегда доступен.

Примечание

Хотя они не указаны ниже, имена camelCase некоторых методов и функций в этом модуле в серии Python 2.x по-прежнему поддерживаются этим модулем.

Деталь реализации CPython: В CPython из-за глобальной блокировки интерпретатора только один поток может выполнять код Python одновременно (хотя некоторые ориентированные на производительность библиотеки могут обойти это ограничение). Если вы хотите, чтобы ваше приложение лучше использовало вычислительные ресурсы многоядерных машин, рекомендуется использовать multiprocessing или concurrent.futures.ProcessPoolExecutor. Однако потоки всё ещё являются подходящей моделью, если вы хотите одновременно выполнять несколько задач, связанных с вводом-выводом.

Этот модуль определяет следующие функции:

threading.active_count()

Возвращает количество объектов Thread, которые в данный момент активны. Возвращаемое значение равно длине списка, возвращаемого функцией enumerate().

threading.current_thread()

Возвращает текущий объект Thread, соответствующий потоку управления вызывающей стороны. Если поток управления вызывающей стороны не был создан через модуль threading, возвращается объект-заглушка потока с ограниченными возможностями.

threading.excepthook(args, /)

Обрабатывает необработанное исключение, поднятое Thread.run().

Аргумент args имеет следующие атрибуты:

  • exc_type: Тип исключения.
  • exc_value: Значение исключения, может быть None.
  • exc_traceback: Трассировка стека исключения, может быть None.
  • thread: Поток, который поднял исключение, может быть None.

Если exc_type равен SystemExit, исключение игнорируется. В противном случае исключение выводится на sys.stderr.

Если эта функция поднимает исключение, вызывается sys.excepthook() для его обработки.

threading.excepthook() может быть переопределён для управления обработкой необработанных исключений, поднятых Thread.run().

Хранение exc_value с помощью пользовательской обработки может создать цикл ссылок. Он должен быть очищен явно, чтобы разорвать цикл ссылок, когда исключение больше не нужно.

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

См. также

sys.excepthook() обрабатывает необработанные исключения.

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

threading.get_ident()

Возвращает идентификатор потока текущего потока. Это ненулевое целое число. Его значение не имеет прямого смысла; оно предназначено для использования, например, для индексирования словаря данных, специфичных для потока. Идентификаторы потоков могут быть повторно использованы, когда поток завершается, а другой поток создаётся.

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

threading.get_native_id()

Возвращает целое число идентификатора потока ядра текущего потока, назначенное ядром. Это целое число без знака. Его значение может использоваться для уникальной идентификации этого конкретного потока в системе (до завершения потока, после чего значение может быть повторно использовано ОС).

Доступность: Windows, FreeBSD, Linux, macOS, OpenBSD, NetBSD, AIX.

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

threading.enumerate()

Возвращает список всех объектов Thread, которые в данный момент активны. Список включает демонические потоки, объекты-заглушки потоков, созданные функцией current_thread(), и основной поток. Он исключает завершенные потоки и потоки, которые ещё не были запущены.

threading.main_thread()

Возвращает основной объект Thread. В нормальных условиях основной поток — это поток, из которого был запущен интерпретатор Python.

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

threading.settrace(func)

Устанавливает функцию трассировки для всех потоков, запущенных из модуля threading. Функция func будет передана в sys.settrace() для каждого потока перед вызовом его метода run().

threading.setprofile(func)

Устанавливает функцию профилирования для всех потоков, запущенных из модуля threading. Функция func будет передана в sys.setprofile() для каждого потока перед вызовом его метода run().

threading.stack_size([size])

Возвращает размер стека потока, используемый при создании новых потоков. Необязательный аргумент size задаёт размер стека для последующих создаваемых потоков и должен быть равен 0 (используется платформа или настроенное значение по умолчанию) или положительным целочисленным значением не меньше 32 768 (32 КБ). Если size не указан, используется 0. Если изменение размера стека потока не поддерживается, генерируется исключение RuntimeError. Если указанный размер стека недопустим, генерируется исключение ValueError, и размер стека не изменяется. 32 КБ в настоящее время является минимальным поддерживаемым размером стека, чтобы гарантировать достаточное пространство стека для самого интерпретатора. Обратите внимание, что на некоторых платформах могут быть определённые ограничения на значения размера стека, например, требование минимального размера стека > 32 КБ или необходимость выделения кратными размеру страницы системной памяти — для получения дополнительной информации следует обратиться к документации платформы (размер страницы 4 КБ является распространённым; использование кратных 4096 для размера стека является рекомендуемым подходом в отсутствие более специфической информации).

Доступность: Windows, системы с POSIX-потоками.

Этот модуль также определяет следующую константу:

threading.TIMEOUT_MAX

Максимальное значение, разрешённое для параметра timeout блокирующих функций (Lock.acquire(), RLock.acquire(), Condition.wait() и т.д.). Указание таймаута, превышающего это значение, приведёт к исключению OverflowError.

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

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

Архитектура этого модуля в общих чертах основана на модели потоков Java. Однако, где Java делает блокировки и переменные условия базовым поведением каждого объекта, в Python они являются отдельными объектами. Класс Python Thread поддерживает подмножество поведения класса Thread Java; в настоящее время нет приоритетов, нет групп потоков, и потоки нельзя уничтожить, остановить, приостановить, возобновить или прервать. Статические методы класса Thread Java, если они реализованы, сопоставляются с функциями на уровне модуля.

Все описанные ниже методы выполняются атомарно.

Данные, локальные для потока

Данные, локальные для потока, — это данные, значения которых специфичны для каждого потока. Для управления данными, локальными для потока, просто создайте экземпляр local (или подкласс) и сохраните в нем атрибуты:

mydata = threading.local()
mydata.x = 1

Значения экземпляра будут разными для разных потоков.

class threading.local

Класс, представляющий данные, локальные для потока.

Дополнительные сведения и подробные примеры см. в строке документации модуля _threading_local.

Объекты потоков

Класс Thread представляет активность, выполняющуюся в отдельном потоке управления. Есть два способа указать активность: передать вызываемый объект в конструктор или переопределить метод run() в подклассе. Никакие другие методы (кроме конструктора) в подклассе переопределять не следует. Другими словами, только переопределяйте методы __init__() и run() этого класса.

После создания объекта потока его активность необходимо запустить, вызвав метод start() потока. Это вызывает метод run() в отдельном потоке управления.

После запуска активности потока поток считается «активным». Он перестает быть активным, когда его метод run() завершается — либо нормально, либо сгенерировав необработанное исключение. Метод is_alive() проверяет, активен ли поток.

Другие потоки могут вызвать метод join() потока. Это блокирует вызывающий поток до тех пор, пока поток, для которого вызывается метод join(), не завершится.

Поток имеет имя. Имя можно передать в конструктор и прочитать или изменить с помощью атрибута name.

Если метод run() генерирует исключение, вызывается threading.excepthook() для его обработки. По умолчанию threading.excepthook() игнорирует SystemExit молчаливо.

Поток может быть помечен как «поток-демон». Значение этого флага состоит в том, что вся программа Python завершается, когда остаются только потоки-демоны. Начальное значение наследуется от создающего потока. Флаг можно установить с помощью свойства daemon или аргумента конструктора daemon.

Примечание

Потоки-демоны внезапно останавливаются при завершении работы. Их ресурсы (такие как открытые файлы, транзакции баз данных и т. д.) могут не быть корректно освобождены. Если вы хотите, чтобы ваши потоки завершались плавно, сделайте их не демоническими и используйте подходящий механизм сигнализации, например, Event.

Существует объект «главный поток»; он соответствует начальному потоку управления в программе Python. Он не является потоком-демоном.

Существует возможность создания «объектов фиктивных потоков». Это объекты потоков, соответствующие «чужеродным потокам», которые являются потоками управления, запущенными вне модуля потоков, например, непосредственно из кода C. Объекты фиктивных потоков имеют ограниченную функциональность; они всегда считаются активными и демоническими и не могут быть join()ed. Они никогда не удаляются, так как невозможно обнаружить завершение чужеродных потоков.

class threading.Thread(group=None, target=None, name=None, args=(), kwargs={}, *, daemon=None)

Этот конструктор всегда следует вызывать со ключевыми аргументами. Аргументы:

group должен быть None; зарезервирован для будущего расширения при реализации класса ThreadGroup.

target — вызываемый объект, который будет вызван методом run(). По умолчанию None, что означает ничего не вызывать.

name — имя потока. По умолчанию генерируется уникальное имя вида «Поток-N», где N — небольшое десятичное число.

args — кортеж аргументов для вызова target. По умолчанию ().

kwargs — словарь ключевых аргументов для вызова target. По умолчанию {}.

Если не None, daemon явно устанавливает, является ли поток демоническим. Если None (значение по умолчанию), свойство daemon наследуется от текущего потока.

Если подкласс переопределяет конструктор, он должен убедиться, что вызовет конструктор базового класса (Thread.__init__()) перед выполнением каких-либо других действий с потоком.

Изменено в версии 3.3: Добавлен аргумент daemon.

start()

Запустить активность потока.

Он должен вызываться не более одного раза на объект потока. Он организует вызов метода run() объекта в отдельном потоке управления.

Этот метод генерирует RuntimeError, если вызывается более одного раза для одного и того же объекта потока.

run()

Метод, представляющий активность потока.

Вы можете переопределить этот метод в подклассе. Стандартный метод run() вызывает вызываемый объект, переданный в конструктор объекта в качестве аргумента target, если таковой имеется, с позиционными и ключевыми аргументами, взятыми из аргументов args и kwargs соответственно.

join(timeout=None)

Ожидание завершения потока. Это блокирует вызывающий поток до тех пор, пока поток, чей метод join() вызывается, не завершится — либо нормально, либо из-за необработанного исключения — или до наступления опционального таймаута.

Если аргумент timeout присутствует и не None, он должен быть числом с плавающей запятой, определяющим таймаут операции в секундах (или дробных частях секунды). Так как join() всегда возвращает None, вы должны вызвать is_alive() после join(), чтобы определить, произошел ли таймаут — если поток все еще активен, вызов join() истек по времени.

Когда аргумент timeout отсутствует или None, операция будет блокироваться до завершения потока.

Поток можно join() многократно.

Вызов join() вызывает RuntimeError, если попытка объединить текущий поток, так как это приведет к тупику. Также является ошибкой join() потока до его запуска, и попытки сделать это вызывают то же исключение.

name

Строка, используемая только для идентификации. Она не имеет семантики. Несколько потоков могут иметь одно и то же имя. Начальное имя устанавливается конструктором.

getName()
setName()

Старая API для получения/установки name; используйте ее напрямую как свойство.

ident

Идентификатор потока или None , если поток еще не запущен. Это целое число, отличное от нуля. См. функцию get_ident(). Идентификаторы потоков могут быть переиспользованы при завершении потока и создании другого.

native_id

Целочисленный идентификатор потока для данного потока. Это неотрицательное целое число или None , если поток еще не запущен. См. функцию get_native_id(). Это идентификатор потока (TID) , присвоенный потоку ОС (ядро). Его значение может использоваться для уникальной идентификации данного потока во всей системе (до тех пор, пока поток не завершится, после чего значение может быть переиспользовано ОС).

Примечание

Аналогично ID процессов, ID потоков действительны (гарантируют уникальность во всей системе) только с момента создания потока до его завершения.

Доступность: Требуется функция get_native_id().

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

is_alive()

Возвращает, жив ли поток.

Этот метод возвращает True непосредственно перед началом метода run() и до окончания метода run(). Функция модуля enumerate() возвращает список всех активных потоков.

daemon

Булевое значение, указывающее, является ли этот поток демоническим потоком (True) или нет (False). Это должно быть установлено перед вызовом start(), иначе возникает RuntimeError. Его начальное значение наследуется от создающего потока; главный поток не является демоническим, поэтому все потоки, созданные в главном потоке, по умолчанию имеют daemon = False.

Вся программа Python завершается, когда не остается активных не демонических потоков.

isDaemon()
setDaemon()

Старая API для получения/установки daemon; используйте ее напрямую как свойство.

Объекты блокировок

Примитивная блокировка — это примитив синхронизации, который не принадлежит какому-либо конкретному потоку при блокировке. В Python в настоящее время это примитив синхронизации самого низкого уровня, реализованный непосредственно модулем расширения _thread.

Примитивная блокировка находится в одном из двух состояний: «заблокированная» или «разблокированная». Она создается в состоянии «разблокированная». У неё есть два основных метода, acquire() и release(). Когда состояние «разблокированная», acquire() изменяет состояние на «заблокированная» и возвращает значение немедленно. Когда состояние «заблокированная», acquire() блокируется до тех пор, пока вызов release() в другом потоке не изменит его на «разблокированная», после чего вызов acquire() устанавливает его обратно в «заблокированная» и возвращает значение. Метод release() должен вызываться только в состоянии «заблокированная»; он изменяет состояние на «разблокированная» и возвращает значение немедленно. Если попытка разблокировать разблокированную блокировку, будет поднято исключение RuntimeError.

Блокировки также поддерживают протокол управления контекстом протокол управления контекстом.

Когда более одного потока заблокированы в acquire() в ожидании изменения состояния на «разблокированная», только один поток продолжает работу, когда вызов release() устанавливает состояние обратно на «разблокированная»; какой из ожидающих потоков будет продолжен, не определено и может меняться в разных реализациях.

Все методы выполняются атомарно.

class threading.Lock

Класс, реализующий объекты примитивных блокировок. После того, как поток получил блокировку, последующие попытки получить её блокируют поток, до тех пор, пока она не будет освобождена; любой поток может её освободить.

Обратите внимание, что Lock на самом деле является функцией-фабрикой, которая возвращает экземпляр самой эффективной версии конкретного класса Lock, поддерживаемого платформой.

acquire(blocking=True, timeout=-1)

Получение блокировки, блокирующая или неблокирующая.

При вызове со значением аргумента blocking установленным в True (по умолчанию), заблокировать до тех пор, пока блокировка не будет разблокирована, а затем установить её в «заблокированную» и вернуть True.

При вызове со значением аргумента blocking установленным в False, не блокировать. Если вызов с blocking установленным в True заблокирует, вернуть False немедленно; в противном случае установить блокировку в «заблокированную» и вернуть True.

При вызове с аргументом timeout с плавающей точкой, равным положительному значению, заблокировать не более чем на заданное количество секунд, указанное в timeout, и до тех пор, пока блокировка не будет получена. Аргумент timeout со значением -1 указывает на неограниченное ожидание. Запрещается указывать timeout, когда blocking имеет значение false.

Возвращаемое значение — True если блокировка успешно получена, False если нет (например, если истекло время ожидания).

Изменено в версии 3.2: Параметр timeout новый.

Изменено в версии 3.2: Получение блокировки теперь может быть прервано сигналами на POSIX, если это поддерживает основная реализация многопоточности.

release()

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

Если блокировка заблокирована, сбросьте её в состояние «разблокированная» и верните значение. Если какие-либо другие потоки заблокированы, ожидая, пока блокировка станет разблокированной, разрешите точно одному из них продолжить работу.

Если вызов производится на разблокированной блокировке, возникает RuntimeError.

Значение не возвращается.

locked()

Возвращает true, если блокировка получена.

Объекты рекурсивных блокировок

Рекурсивная блокировка — это примитив синхронизации, который может быть получен несколько раз одним и тем же потоком. Внутренне он использует понятия «владеющий поток» и «уровень рекурсии» в дополнение к состоянию «заблокирована/разблокирована», используемому примитивными блокировками. В состоянии «заблокирована» некоторый поток владеет блокировкой; в состоянии «разблокирована» ни один поток не владеет ею.

Чтобы заблокировать блокировку, поток вызывает свой метод acquire(); это возвращает значение, когда поток владеет блокировкой. Чтобы разблокировать блокировку, поток вызывает метод release(). Вызовы acquire()/release() могут быть вложенными; только последний вызов release() (вызов release() самого внешнего пары) сбрасывает блокировку в состояние «разблокирована» и позволяет другому потоку, заблокированному в acquire(), продолжить работу.

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

class threading.RLock

Этот класс реализует объекты рекурсивных блокировок. Рекурсивная блокировка должна быть освобождена потоком, который её получил. После того, как поток получил рекурсивную блокировку, тот же поток может получить её снова без блокировки; поток должен освободить её один раз для каждого раза, когда он её получил.

Обратите внимание, что RLock на самом деле является функцией-фабрикой, которая возвращает экземпляр самой эффективной версии конкретного класса RLock, поддерживаемого платформой.

acquire(blocking=True, timeout=-1)

Получение блокировки, блокирующая или неблокирующая.

При вызове без аргументов: если этот поток уже владеет блокировкой, увеличить уровень рекурсии на единицу и вернуть значение немедленно. В противном случае, если другая нить владеет блокировкой, заблокировать поток до тех пор, пока блокировка не будет разблокирована. Как только блокировка разблокирована (не принадлежит ни одному потоку), захватить владение, установить уровень рекурсии в единицу и вернуть значение. Если более одного потока заблокированы, ожидая, пока блокировка не будет разблокирована, только один сможет захватить владение блокировкой. В этом случае значение не возвращается.

При вызове с аргументом blocking установленным в true, выполнить то же самое, что и при вызове без аргументов, и вернуть True.

При вызове с аргументом blocking установленным в false, не блокировать. Если вызов без аргумента заблокирует, вернуть False немедленно; в противном случае выполнить то же самое, что и при вызове без аргументов, и вернуть True.

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

Изменено в версии 3.2: Параметр timeout новый.

release()

Освобождение блокировки, уменьшая уровень рекурсии. Если после уменьшения он равен нулю, сбросить блокировку в состояние «разблокированная» (ни один поток не владеет ею) и если какие-либо другие потоки заблокированы, ожидая, пока блокировка станет разблокированной, разрешить точно одному из них продолжить работу. Если после уменьшения уровень рекурсии по-прежнему не равен нулю, блокировка остаётся заблокированной и принадлежит вызывающему потоку.

Вызывайте этот метод только тогда, когда вызывающий поток владеет блокировкой. Если метод вызван, когда блокировка разблокирована, возникает RuntimeError.

Значение не возвращается.

END_OF_DOCUMENT_MARKER

Объекты условия

Переменная условия всегда связана с каким-либо замком; этот замок может быть передан или будет создан по умолчанию. Передача замка полезна, когда несколько переменных условия должны использовать один и тот же замок. Замок является частью объекта условия: вам не нужно отслеживать его отдельно.

Переменная условия подчиняется протоколу управления контекстом: использование оператора with приобретает связанный замок на время выполнения заключённого блока. Методы acquire() и release() также вызывают соответствующие методы связанного замка.

Другие методы должны вызываться, удерживая связанный замок. Метод wait() освобождает замок и затем блокируется, пока другой поток не разбудит его, вызвав notify() или notify_all(). После пробуждения wait() повторно приобретает замок и возвращает значение.

Также можно указать таймаут.

Метод notify() разбуживает один из потоков, ждущих переменной условия, если такие потоки есть. Метод notify_all() разбуживает все потоки, ждущие переменной условия.

Примечание: методы notify() и notify_all() не освобождают замок; это означает, что разбуженный поток или потоки не сразу вернутся из своего вызова wait(), а только тогда, когда поток, вызвавший notify() или notify_all(), окончательно передаст владение замком.

Типичный стиль программирования с использованием переменных условия использует замок для синхронизации доступа к некоторому общему состоянию; потоки, заинтересованные в определённом изменении состояния, вызывают wait() многократно, пока не увидят желаемое состояние, в то время как потоки, изменяющие состояние, вызывают notify() или notify_all(), когда они изменяют состояние таким образом, что оно потенциально может быть желаемым состоянием для одного из ожидающих потоков. Например, следующий код демонстрирует общую ситуацию «производитель-потребитель» с неограниченной ёмкостью буфера:

# Consume one item
with cv:
    while not an_item_is_available():
        cv.wait()
    get_an_available_item()

# Produce one item
with cv:
    make_an_item_available()
    cv.notify()

Цикл проверки условия приложения while необходим, потому что wait() может возвратиться через произвольно длительное время, и условие, которое спровоцировало вызов notify(), может больше не выполняться. Это присуще многопоточной программированию. Метод wait_for() может использоваться для автоматизации проверки условия и упрощения вычисления таймаутов:

# Consume an item
with cv:
    cv.wait_for(an_item_is_available)
    get_an_available_item()

Чтобы выбрать между notify() и notify_all(), необходимо учитывать, может ли одно изменение состояния быть интересным только для одного или нескольких ожидающих потоков. Например, в типичной ситуации «производитель-потребитель» добавление одного элемента в буфер требует разбудить только один поток потребителя.

class threading.Condition(lock=None)

Этот класс реализует объекты переменных условия. Переменная условия позволяет одному или нескольким потокам ожидать, пока другой поток не уведомит их.

Если аргумент lock задан и не None, он должен быть объектом Lock или RLock, и он используется в качестве базового замка. В противном случае создаётся новый объект RLock и используется в качестве базового замка.

Изменено в версии 3.3: изменено с фабричной функции на класс.

acquire(*args)

Приобретение базового замка. Этот метод вызывает соответствующий метод базового замка; возвращаемое значение — это то, что возвращает этот метод.

release()

Освобождение базового замка. Этот метод вызывает соответствующий метод базового замка; возвращаемое значение отсутствует.

wait(timeout=None)

Ожидание уведомления или таймаута. Если вызывающий поток не приобрел замок, когда этот метод был вызван, возникает RuntimeError.

Этот метод освобождает базовый замок, а затем блокируется, пока его не разбудит вызов notify() или notify_all() для той же переменной условия в другом потоке или пока не наступит указанный таймаут. После пробуждения или таймаута он повторно приобретает замок и возвращает значение.

Если аргумент timeout присутствует и не None, он должен быть числом с плавающей точкой, указывающим таймаут операции в секундах (или дробных частях секунды).

Когда базовый замок является RLock, он не освобождается с помощью метода release(), так как это может не разблокировать замок при многократном рекурсивном приобретении. Вместо этого используется внутренний интерфейс класса RLock, который действительно разблокирует его, даже если он был приобретён несколько раз рекурсивно. Затем используется другой внутренний интерфейс для восстановления уровня рекурсии, когда замок повторно приобретается.

Возвращаемое значение — это True , если не истекло время ожидания, в противном случае — False.

Изменено в версии 3.2: Ранее метод всегда возвращал None.

wait_for(predicate, timeout=None)

Ожидание, пока условие не станет истинным. predicate должен быть вызываемым объектом, результат которого будет интерпретироваться как булево значение. Может быть предоставлен таймаут, задающий максимальное время ожидания.

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

Без учёта таймаута, вызов этого метода примерно эквивалентен записи:

while not predicate():
    cv.wait()

Поэтому применяются те же правила, что и для wait(): замок должен быть удерживаем при вызове и повторно приобретается при возврате. Предикат вычисляется при удерживании замка.

Введено в версии 3.2.

notify(n=1)

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

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

Текущая реализация разбуживает ровно n потоков, если ожидающих потоков по крайней мере n. Однако нельзя полагаться на это поведение. Будущая оптимизированная реализация может иногда разбудить больше, чем n потоков.

Примечание: разбуженный поток не фактически возвращается из своего вызова wait(), пока не сможет повторно приобрести замок. Так как notify() не освобождает замок, его вызывающий поток должен это сделать.

notify_all()

Разбудить все потоки, ожидающие этого условия. Этот метод действует как notify(), но разбуживает все ожидающие потоки вместо одного. Если вызывающий поток не приобрел замок, когда этот метод был вызван, возникает RuntimeError.

Объекты семафоров

Это одна из старейших примитивных операций синхронизации в истории компьютерной науки, изобретённая ранним голландским компьютерным учёным Эдсгером В. Дейкстрой (он использовал имена P() и V() вместо acquire() и release()).

Семафор управляет внутренней счётной переменной, которая уменьшается при каждом вызове acquire() и увеличивается при каждом вызове release(). Счётчик никогда не может опуститься ниже нуля; когда acquire() обнаруживает, что он равен нулю, он блокируется, ожидая, пока какая-то другая нить не вызовет release().

Семафоры также поддерживают протокол управления контекстом протокол управления контекстом.

class threading.Semaphore(value=1)

Этот класс реализует объекты семафоров. Семафор управляет атомным счётчиком, представляющим количество вызовов release() минус количество вызовов acquire(), плюс начальное значение. Метод acquire() блокируется при необходимости, пока не сможет вернуть значение, не сделав счётчик отрицательным. Если не указано, значение по умолчанию равно 1.

Необязательный аргумент задаёт начальное значение для внутреннего счётчика; по умолчанию он равен 1. Если заданное значение меньше 0, генерируется ValueError.

Изменено в версии 3.3: изменено с фабричной функции на класс.

acquire(blocking=True, timeout=None)

Получить семафор.

При вызове без аргументов:

  • Если внутренний счётчик больше нуля при входе, уменьшите его на единицу и верните True немедленно.
  • Если внутренний счётчик равен нулю при входе, заблокируйтесь, ожидая пробуждения от вызова release(). После пробуждения (и счётчик больше 0), уменьшите счётчик на 1 и верните True. Ровно одна нить будет пробуждена каждым вызовом release(). Порядок пробуждения нитей не гарантируется.

При вызове с блокировкой, установленной в false, не блокироваться. Если вызов без аргумента заблокируется, вернуть False немедленно; в противном случае сделайте то же, что и при вызове без аргументов, и верните True.

При вызове с timeout отличным от None, он будет блокироваться не более чем timeout секунд. Если acquire не завершится успешно в этом интервале, верните False. В противном случае верните True.

Изменено в версии 3.2: Параметр timeout новый.

release()

Освободить семафор, увеличив внутренний счётчик на единицу. Если он был равен нулю при входе и другая нить ждёт, пока он снова не станет больше нуля, разбудить эту нить.

class threading.BoundedSemaphore(value=1)

Класс, реализующий объекты ограниченных семафоров. Ограниченный семафор проверяет, чтобы его текущее значение не превышало его начального значения. Если превысит, генерируется ValueError. В большинстве случаев семафоры используются для защиты ресурсов с ограниченной ёмкостью. Если семафор слишком много раз освобождается, это указывает на ошибку. Если не задано, значение по умолчанию равно 1.

Изменено в версии 3.3: изменено с фабричной функции на класс.

Semaphore Пример

Семафоры часто используются для защиты ресурсов с ограниченной ёмкостью, например, сервера базы данных. В любой ситуации, где размер ресурса фиксирован, следует использовать ограниченный семафор. Перед запуском нитей-работников основная нить инициализирует семафор:

maxconnections = 5
# ...
pool_sema = BoundedSemaphore(value=maxconnections)

После запуска нити-работники вызывают методы acquire и release семафора, когда им нужно подключиться к серверу:

with pool_sema:
    conn = connectdb()
    try:
        # ... use connection ...
    finally:
        conn.close()

Использование ограниченного семафора уменьшает вероятность того, что ошибка программирования, которая приводит к тому, что семафор освобождается больше, чем приобретается, останется незамеченной.

Объекты событий

Это один из самых простых механизмов связи между нитями: одна нить сигнализирует о событии, а другие нити ждут его.

Объект события управляет внутренней меткой, которая может быть установлена в true с помощью метода set() и сброшена в false с помощью метода clear(). Метод wait() блокируется, пока метка не станет равной true.

class threading.Event

Класс, реализующий объекты событий. Событие управляет флагом, который может быть установлен в true с помощью метода set() и сброшен в false с помощью метода clear(). Метод wait() блокируется, пока флаг не станет равным true. Флаг изначально равен false.

Изменено в версии 3.3: изменено с фабричной функции на класс.

is_set()

Возвращает True тогда и только тогда, когда внутренний флаг равен true.

set()

Установить внутренний флаг в значение true. Все нити, ожидающие, пока он станет true, будут разбужены. Нити, которые вызывают wait(), когда флаг уже true, не будут блокироваться.

clear()

Сбросить внутренний флаг в false. Впоследствии нити, вызывающие wait(), будут блокироваться, пока set() не вызовет установку внутреннего флага в true снова.

wait(timeout=None)

Заблокироваться, пока внутренний флаг не станет равным true. Если внутренний флаг равен true при входе, вернуть значение немедленно. В противном случае заблокироваться, ожидая, пока другая нить не вызовет set() для установки флага в true или пока не наступит необязательный таймаут.

Когда параметр timeout присутствует и не равен None, он должен быть числом с плавающей точкой, определяющим таймаут операции в секундах (или долях секунды).

Этот метод возвращает True тогда и только тогда, когда внутренний флаг был установлен в true, либо до вызова wait, либо после его начала, поэтому он всегда вернёт True, за исключением случая, когда задан таймаут, и операция его превысила.

Изменено в версии 3.1: Ранее метод всегда возвращал None.

Объекты таймеров

Этот класс представляет собой действие, которое должно быть выполнено только после того, как пройдёт определённое количество времени — таймер. Timer является подклассом Thread и, таким образом, также служит примером создания пользовательских нитей.

Таймеры запускаются, как и нити, вызовом их метода start(). Таймер можно остановить (до начала его действия) вызовом метода cancel(). Интервал ожидания таймера перед выполнением его действия может отличаться от указанного пользователем.

Например:

def hello():
    print("hello, world")

t = Timer(30.0, hello)
t.start()  # after 30 seconds, "hello, world" will be printed
class threading.Timer(interval, function, args=None, kwargs=None)

Создать таймер, который выполнит функцию с аргументами args и ключевыми аргументами kwargs, через interval секунд. Если args равно None (по умолчанию), будет использоваться пустой список. Если kwargs равно None (по умолчанию), будет использоваться пустой словарь.

Изменено в версии 3.3: изменено с фабричной функции на класс.

cancel()

Остановить таймер и отменить выполнение действия таймера. Это сработает только если таймер всё ещё находится в стадии ожидания.

END_OF_DOCUMENT_MARKER

Объекты барьеров

Новая версия с 3.2.

Этот класс предоставляет простую синхронизационную примитив для использования фиксированным числом потоков, которые должны ждать друг друга. Каждый из потоков пытается пройти барьер, вызвав метод wait(), и будет заблокирован, пока все потоки не выполнят свои вызовы wait(). В этот момент потоки освобождаются одновременно.

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

В качестве примера, вот простой способ синхронизации потока клиента и сервера:

b = Barrier(2, timeout=5)

def server():
    start_server()
    b.wait()
    while True:
        connection = accept_connection()
        process_server_connection(connection)

def client():
    b.wait()
    while True:
        connection = make_connection()
        process_client_connection(connection)
class threading.Barrier(parties, action=None, timeout=None)

Создайте объект барьера для parties числа потоков. Если указано, action — это вызываемая функция, которая будет вызвана одним из потоков, когда они будут освобождены. timeout — это значение таймаута по умолчанию, если он не указан для метода wait().

wait(timeout=None)

Пройди барьер. Когда все потоки, участвующие в барьере, вызовут эту функцию, все они будут одновременно освобождены. Если указан timeout, он используется вместо любого, который был передан в конструктор класса.

Возвращаемое значение — целое число от 0 до parties – 1, разное для каждого потока. Это можно использовать для выбора потока, который выполнит некоторую специальную работу, например:

i = barrier.wait()
if i == 0:
    # Only one thread needs to print this
    print("passed the barrier")

Если в конструктор был передан action, один из потоков его вызовет перед освобождением. Если при этом возникнет ошибка, барьер переводится в состояние «сломанный».

Если вызов превысит время ожидания, барьер переводится в состояние «сломанный».

Этот метод может вызвать исключение BrokenBarrierError, если барьер сломан или сброшен, пока поток ожидает.

reset()

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

Обратите внимание, что для использования этой функции могут потребоваться внешние механизмы синхронизации, если состояние других потоков неизвестно. Если барьер сломан, лучше просто оставить его и создать новый.

abort()

Переводит барьер в состояние «сломанный». Это приводит к тому, что любые активные или будущие вызовы wait() завершаются с ошибкой BrokenBarrierError. Используйте это, например, если один из потоков должен прервать работу, чтобы избежать тупика в приложении.

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

parties

Число потоков, необходимых для прохождения барьера.

n_waiting

Число потоков, которые в данный момент ожидают в барьере.

broken

Булево значение, которое True, если барьер в состоянии «сломанный».

exception threading.BrokenBarrierError

Это исключение, подкласс RuntimeError, возникает при сбросе или разрыве объекта Barrier.

Использование блокировок, условий и семафоров в операторе with

Все объекты, предоставленные этим модулем, имеющие методы acquire() и release(), могут использоваться в качестве менеджеров контекста для оператора with. Метод acquire() вызывается при входе в блок, а release() — при выходе из блока. Следовательно, следующий фрагмент:

with some_lock:
    # do something...

эквивалентен:

some_lock.acquire()
try:
    # do something...
finally:
    some_lock.release()

В настоящее время объекты Lock, RLock, Condition, Semaphore и BoundedSemaphore могут использоваться в качестве менеджеров контекста оператора with.

© 2001–2022 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.8/library/threading.html

Spec-Zone.ru

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