Перейти к содержимому

Для чего нужен threading lock

  • автор:

Для чего нужен lock в python? Как работает данный пример кода?

Vindicar

Это объясняется тем, что в базовом питоне потоки не вполне честные — они конкурируют за global interpreter lock, так что код выполняется всё равно поочерёдно. Так что многопоточность в питоне полезна с точки зрения распараллеливания, но не ускорения. ЕМНИП, есть реализации питона, в которых нет этой GIL problem.
Но нужно иметь ввиду, что этот GIL блокирует только элементарные операции (как в твоём примере), тогда как явное использование lock может накрывать целые блоки кода, состоящие из нескольких операций с защищаемым ресурсом.

Вот тебе пример:

import threading import time class Data: def __init__(self): self.x: int = 0 self.y: int = 0 do_sleep = False run = True def reader(d: Data): while run: x, y = d.x, d.y # по идее это условие не должно выполниться никогда if (x != 0) != (y != 0): print(f'Got x= and y=') else: print(f'OK ', end='\x08\x08\x08\x08') def writer(d: Data): while run: if d.x == 0: d.x = 1 if do_sleep: pass d.y = 1 else: d.x = 0 if do_sleep: pass d.y = 0 do_sleep = False instance = Data() reader_thread = threading.Thread(target=reader, args=(instance,), daemon=True) writer_thread = threading.Thread(target=writer, args=(instance,), daemon=True) reader_thread.start() writer_thread.start() try: input() finally: run = False reader_thread.join() writer_thread.join()

На моей машине, если if do_sleep: pass закомментировать, то в консоли высвечивается только OK — иными словами, присваивание двух полей выполняется достаточно быстро, чтобы поток не успел переключиться в промежутке. Как следствие, reader() всегда видит либо x=0 y=0, либо x=1 y=1.
Но если if do_sleep: pass оставить, то выполнение тела цикла замедляется достаточно, чтобы поток успел переключиться — и, как следствие, reader() начинает видеть структуру данных Data в неконсистентном состоянии, когда x=0 y=1 или когда x=1 y=0.
И вот чтобы не гадать «успеет — не успеет», нужно в таких случаях защищать связные серии обращений к структуре с помощью мьютекса, ну или в питоновских терминах — Lock.

Класс RLock() модуля threading в Python

Класс RLock() модуля threading реализует объекты реентерабельной (повторной) блокировки.

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

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

Повторная блокировка потока — это примитив синхронизации, который может быть получен несколько раз одним и тем же потоком. Внутри он использует концепции «владеющего потоком» и «уровня рекурсии» в дополнение к locked / unlocked состоянию, используемому примитивными блокировками threading.Lock . В заблокированном locked состоянии какой-то поток владеет блокировкой, в разблокированном unlocked состоянии ни один поток не владеет им.

Чтобы включить блокировку, поток вызывает свой метод RLock.acquire() , он возвращает результат своему экземпляру, когда поток владеет блокировкой. Чтобы снять блокировку, поток вызывает свой метод RLock.release() .

Пары вызовов RLock.acquire() / RLock.release() могут быть вложенными, только последний RLock.release() ( release() самой внешней пары) сбрасывает режим locked на unlocked и позволяет продолжить работу другому потоку, заблокированному в RLock.acquire() .

Класс threading.RLock() также поддерживают протокол управления контекстом.

Методы объекта threading.RLock .

  • RLock.acquire() устанавливает блокировку,
  • RLock.release() снимает блокировку,
  • Пример работы повторной блокировки threading.RLock() ,
RLock.acquire(blocking=True, timeout=-1) :

Метод RLock.acquire() устанавливает блокировку, блокирующую или неблокирующую..

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

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

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

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

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

RLock.release() :

Метод RLock.release() снимает блокировку, уменьшив уровень рекурсии. Если после декремента он равен нулю, то сбрасывает блокировку на unlocked (т.е. блокировка не принадлежащую ни одному потоку), и если какие-либо другие потоки заблокированы, ожидая разблокировки, то разрешит выполнение ровно одному из них. Если после декремента уровень рекурсии все еще отличен от нуля, то блокировка остается locked и принадлежит вызывающему потоку.

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

Возвращаемого значения нет.

Пример работы повторной блокировки threading.RLock() .

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

В этом случае второму вызову Lock.acquire() необходимо дать нулевой тайм-аут, чтобы предотвратить его блокировку, потому что блокировка была уже получена первым вызовом.

>>> import threading >>> lock = threading.Lock() >>> lock # # первая блокировка >>> 'First try :', lock.acquire() # ('First try :', True) >>> lock # # вторая блокировка 'Second try:', lock.acquire(timeout=0) # ('Second try:', False) 

В ситуации, когда отдельный код из одного и того же потока должен “повторно получить” блокировку, необходимо использовать объекты threading.RLock .

Единственным изменением в коде из предыдущего примера была замена объекта Rlock на Lock.

>>> import threading >>> rlock = threading.RLock() # первая блокировка >>> 'First try :', rlock.acquire() # ('First try :', True) # вторая блокировка >>>'Second try:', rlock.acquire(timeout=0) # ('Second try:', True) # смотрим состояние - счетчик повторных блокировок count=2 >>> rlock # # разблокируем 2 раза >>> rlock.release() >>> rlock.release() # смотрим состояние - счетчик повторных блокировок count=0 >>> rlock # # пробуем разблокировать уже разблокированное состояние >>> rlock.release() # Traceback (most recent call last): # File "", line 1, in # RuntimeError: cannot release un-acquired lock 
  • КРАТКИЙ ОБЗОР МАТЕРИАЛА.
  • Получение общих сведений о потоках, модуль threading
  • Класс Thread() модуля threading
  • Класс local() модуля threading
  • Класс Event() модуля threading
  • Класс Lock() модуля threading
  • Класс RLock() модуля threading
  • Класс Condition() модуля threading
  • Класс Semaphore() модуля threading
  • Класс Timer() модуля threading
  • Класс Barrier() модуля threading
  • Протокол управления контекстом в модуле threading
  • Трассировка и профилирование потоков модулем threading
  • Как перезапускать потоки?

Нужен ли класс threading.Lock?

Вопрос вот в чем.
Если у питона есть GIL который блокирует доступ разных потоков к одному и тому же участку памяти, что собственно является одним из якорей в производетельности, то зачем нужен класс threading.Lock, ведь потоки уже заблокированы GIL-ом?

P.S.
Наверное было бы круто, если у кого-то завалялись статьи на эту тему

  • Вопрос задан более двух лет назад
  • 364 просмотра

Комментировать
Решения вопроса 2

Скопируй и запусти две разных версии кода.

Без Lock

from threading import * def work(i): for _ in range(100): print(f"hello i'm a thread #") t1 = Thread(target=work, args=(1,)) t2 = Thread(target=work, args=(2,)) t1.start() t2.start() t1.join() t2.join()

С Lock

from threading import * lock = Lock() def work(i): for _ in range(100): with lock: print(f"hello i'm a thread #") t1 = Thread(target=work, args=(1,)) t2 = Thread(target=work, args=(2,)) t1.start() t2.start() t1.join() t2.join()

Как можешь заметить, первый вариант иногда печатает две строки на одной, а иногда печатает пустые строки

Без Lock

from threading import * from time import sleep class GlobalState: def __init__(self, x): self.x = x def set_x(self, x): self.x = x def reader(state: GlobalState): if state.x % 2 == 0: sleep(0.01) # simulate OS context switch print(f" is even") else: print(f" is odd") def changer(state: GlobalState): state.set_x(state.x + 1) state = GlobalState(2) t1 = Thread(target=reader, args=(state,)) t2 = Thread(target=changer, args=(state,)) t1.start() t2.start() t1.join() t2.join()

С Lock

from threading import * from time import sleep class GlobalState: def __init__(self, x): self.x = x self.lock = Lock() def set_x(self, x): self.x = x def reader(state: GlobalState): with state.lock: if state.x % 2 == 0: sleep(0.01) # simulate OS context switch print(f" is even") else: print(f" is odd") def changer(state: GlobalState): with state.lock: state.set_x(state.x + 1) state = GlobalState(2) t1 = Thread(target=reader, args=(state,)) t2 = Thread(target=changer, args=(state,)) t1.start() t2.start() t1.join() t2.join()

Ну и совсем упоротый пример для тех, кто говорит, что list — threadsafe (что фактически является истиной, но логически не всегда) и не нужно использовать Lock:

Открыть

from threading import * from random import * class GlobalState: def __init__(self): self.x = [] def do_something_changing(self): if random() < 0.5: self.x.append(1) elif self.x: self.x.pop() def reader(state: GlobalState): for _ in range(10000000): if len(state.x) % 2 == 0: if len(state.x) % 2 != 0: # wtf how it's possible? print(f"length was even before context switch") def changer(state: GlobalState): for _ in range(10000000): state.do_something_changing() state = GlobalState() t1 = Thread(target=reader, args=(state,)) t2 = Thread(target=changer, args=(state,)) t1.start() t2.start() t1.join() t2.join()

И напоследок. Хватит программировать (или пытаться) на тредах. Это сложно и никому не нужно. Давно существуют куда более удачные реализации использования всех ядер процессора (csp например в golang). А если треды используются для IO (а в питоне они в 99.9% используются именно для IO), то давно есть и довольно юзабельный asyncio.

Python в три ручья (часть 2). Блокировки

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

Потоки стремятся к ресурсам, с которыми должны работать. И когда к одному и тому же ресурсу обращается несколько потоков, возникает конфликт. Как его предотвратить?

Потоки нельзя в любой момент напрямую остановить или завершить: на то они и потоки. Но можно на их пути поставить дамбу — блокировку. Она пропустит только один поток, а остальные временно удержит. Так вы исключите конфликт.

Когда поток A выполняет операцию с общими ресурсами, а поток Б не может вмешаться в нее до завершения — говорят, что такая операция атомарна. Залогом потокобезопасности как раз выступает атомарность — непрерывность, неделимость операции.

Простая блокировка в Python

Взаимоисключение (mutual exception, кратко — mutex) — простейшая блокировка, которая на время работы потока с ресурсом закрывает последний от других обращений. Реализуют это с помощью класса Lock.

import threading mutex = threading.Lock()

Мы создали блокировку с именем mutex, но могли бы назвать её lock или иначе. Теперь её можно ставить и снимать методами .acquire() и .release():

resource = 0 def thread_safe_function(): global resource for i in range(1000000): mutex.acquire() # Делаем что-то с переменной resource mutex.release()

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

С блокировками и без. Пример–сравнение

Что происходит, когда два потока бьются за ресурсы, и как при этом сохранить целостность данных? Разберёмся на практике.

Возьмём простейшие операции инкремента и декремента (увеличения и уменьшения числа). В роли общих ресурсов выступят глобальные числовые переменные: назовём их protected_resource и unprotected_resource. К каждой обратятся по два потока: один будет в цикле увеличивать значение с 0 до 50 000, другой — уменьшать до 0. Первую переменную обработаем с блокировками, а вторую — без.

import threading protected_resource = 0 unprotected_resource = 0 NUM = 50000 mutex = threading.Lock() # Потокобезопасный инкремент def safe_plus(): global protected_resource for i in range(NUM): # Ставим блокировку mutex.acquire() protected_resource += 1 mutex.release() # Потокобезопасный декремент def safe_minus(): global protected_resource for i in range(NUM): mutex.acquire() protected_resource -= 1 mutex.release() # То же, но без блокировки def risky_plus(): global unprotected_resource for i in range(NUM): unprotected_resource += 1 def risky_minus(): global unprotected_resource for i in range(NUM): unprotected_resource -= 1 

В названия потокобезопасных функций мы поставили префикс safe_, а небезопасных — risky_.

Создадим 4 потока, которые будут выполнять функции с блокировками и без:

thread1 = threading.Thread(target = safe_plus) thread2 = threading.Thread(target = safe_minus) thread3 = threading.Thread(target = risky_plus) thread4 = threading.Thread(target = risky_minus) thread1.start() thread2.start() thread3.start() thread4.start() thread1.join() thread2.join() thread3.join() thread4.join() print ("Результат при работе с блокировкой %s" % protected_resource) print ("Результат без блокировки %s" % unprotected_resource)

Запускаем код несколько раз подряд и видим, что полученное без блокировки значение меняется случайным образом. При использовании блокировки всё работает последовательно: сначала значение растёт, затем — уменьшается, и в итоге получаем 0. А потоки thread3 и thread4 работают без блокировки и наперебой обращаются к глобальной переменной. Каждый выполняет столько операций своего цикла, сколько успевает за время активности. Поэтому при каждом запуске получаем случайные числа.

Как избежать взаимных блокировок?

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

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

Пишите код так, чтобы блокировки снимались, даже если функция выбрасывает исключение и завершает работу нештатно. Подстраховаться можно с помощью конструкции try-except-finally:

try: mutex.acquire() # Ваш код. except SomethingGoesWrong: # Обрабатываем исключения finally: # Ещё код mutex.release()

Другие инструменты синхронизации в Python

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

Семафоры (Semaphore)

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

Значение счётчика уменьшается с каждым новым вызовом acquire(), то есть с подключением к ресурсу новых потоков. Когда ресурс высвобождается, значение возрастает. При нулевом значении счётчика работа потока останавливается, пока другой поток не вызовет метод release(). По умолчанию значение счётчика равно 1.

s = Semaphore(5) # В скобках при необходимости указывают стартовое значение счётчика 

Можно создать «ограниченный семафор» конструктором BoundedSemaphore().

События (Event)

Событие — сигнал от одного потока другим. Если событие возникло — ставят флаг методом .set(), а после обработки события — снимают с помощью .clear(). Пока флага нет, ресурс заблокирован. Ждать события могут один или несколько потоков. Важную роль играет wait(): если флаг установлен, этот метод спокойно отдаёт управление ресурсом; если нет — блокирует его на заданное время или до установки флага одним из потоков.

e = threading.Event() def event_manager(): # Ждём, когда кто-нибудь захватит флаг e.wait() . # Ставим флаг e.set() # Работаем с ресурсом . # Снимаем флаг и ждём нового e.clear()

Если нужно задать время ожидания, его пишут в секундах, в виде числа с плавающей запятой. Например: e.wait(3,0).

Метод is_set() проверяет, активно ли событие. Важно следить, чтобы события попадали в поле зрения потоков-потребителей сразу после появления. Иначе работа зависящих от события потоков нарушится.

Рекурсивная блокировка (RLock)

Такая блокировка позволяет одному потоку захватывать ресурс несколько раз, но блокирует все остальные потоки. Это полезно, когда вы используете вложенные функции, каждая из которых тоже применяет блокировку. Число вложенных .acquire() и .release() не даст интерпретатору запутаться, сколько раз поток имеет право захватывать ресурс, а когда блокировку надо снять полностью. Механизм основан на классе RLock:

import threading, random counter = 0 re_mutex = threading.RLock() def step_one(): global counter re_mutex.acquire() counter = random.randint(1,100) print("Random number %s" % counter) re_mutex.release() def step_two(): global counter re_mutex.acquire() counter *= 2 print("Doubled = %s" % counter) re_mutex.release() def walkthrough(): re_mutex.acquire() try: step_one() step_two() finally: re_mutex.release() t = threading.Thread(target = walkthrough) t2 = threading.Thread(target = walkthrough) t.start() t2.start() t.join() t2.join()

Запустите это и проверьте результат: арифметика должна быть верна.

Теперь попробуйте убрать блокировку внутри walkthrough:

def walkthrough(): step_one() step_two()

Ещё раз запустите код — порядок действий нарушится. Программа умножит на 2 только второе случайное число, а затем удвоит полученное произведение.

Переменные состояния (Condition)

Переменная состояния — усложнённый вариант события (Event). Через Condition на ресурс ставят блокировку нужного типа, и она работает, пока не произойдёт ожидаемое потоками изменение. Как только это случается, один или несколько потоков разблокируются. Оповестить потоки о событии можно методами:

  • notify() — для одного потока;
  • notifyAll() — для всех ожидающих потоков.

Это выглядит так:

# Создаём рекурсивную блокировку mutex = threading.RLock() # Создаём переменную состояния и связываем с блокировкой cond = threading.Condition(mutex) # Поток-потребитель ждёт свободного ресурса и захватывает его def consumer(): while True: cond.acquire() while not resourse_free(): cond.wait() get_free_resource() cond.release() # Поток-производитель разблокирует ресурс и уведомляет об этом потребителя def producer(): while True: cond.acquire() unblock_resource() # Сигналим потоку: "Налетай на новые данные!" cond.notify() cond.release()

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

cond = threading.Condition(mutex) another_cond = threading.Condition(mutex)

Компактные блокировки с with

При множестве участков с блокировками каждый раз прописывать «захват» и «высвобождение» утомительно. Сократить код поможет конструкция с оператором with. Она использует менеджер контекста, который позволяет сначала подготовить приложение к выполнению фрагмента кода, а затем гарантированно освободить задействованные ресурсы.

Чтобы понять дальнейший материал, кратко разберем работу with, хотя это и не про блокировки. У класса, который мы собираемся использовать с with, должно быть два метода:

  • «Предисловие» — метод __enter__(). Здесь можно ставить блокировку и прописывать другие настройки;
  • «Послесловие» — метод __exit__(). Он срабатывает, когда все инструкции выполнены или работа блока прервана. Здесь можно снять блокировку и/или предусмотреть реакцию на исключения, которые могут быть выброшены.

Удача! У нашего целевого класса Lock эти два метода уже прописаны. Поэтому любой экземпляр объекта Lock можно использовать с with без дополнительных настроек.

Отредактируем функцию из примера с инкрементом. Поставим блокировку, которая сама снимется, как только управляющий поток выйдет за пределы with-блока:

def safe_plus(): global protected_resource for i in range(NUM): with mutex: protected_resource += 1 # И никаких acquire-release! 

Добавить комментарий

Ваш адрес email не будет опубликован. Обязательные поля помечены *

https://alkogolizm.vyvod-iz-zapoya-v-stacionare-samara11.ru/