Spec-Zone.ru › Python 3.9

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 (используется платформа или настроенное значение по умолчанию), либо положительным целым числом, не менее 32768 (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 они представляют собой отдельные объекты. Класс Thread Python поддерживает подмножество поведения класса 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. Это не демонический поток.

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

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

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

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

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

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

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

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

Если не 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) присвоенный потоку ОС (ядро). Его значение может быть использовано для уникальной идентификации данного потока в системе (до завершения потока, после чего значение может быть переиспользовано ОС).

Примечание

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

Доступность: Требует функции 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() блокируется, если необходимо, пока не сможет вернуть значение без создания отрицательного значения счётной переменной. Если не указано, значение по умолчанию для value равно 1.

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

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

acquire(blocking=True, timeout=None)

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

Когда вызывается без аргументов:

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

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

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

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

release(n=1)

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

Изменено в версии 3.9: Добавлен параметр n для одновременного освобождения нескольких ожидающих нитей.

class threading.BoundedSemaphore(value=1)

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

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

Semaphore Пример

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

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

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

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

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

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

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

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

class threading.Event

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

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

is_set()

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

set()

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

clear()

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

wait(timeout=None)

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

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

Этот метод возвращает True тогда и только тогда, когда внутренний флаг был установлен в истинное значение, либо до вызова ожидания, либо после его начала, поэтому он всегда вернёт 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)

Создаёт таймер, который будет запускать function с аргументами 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.9/library/threading.html

Spec-Zone.ru

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