Нередко в потоках используются некоторые разделяемые ресурсы, общие для всей программы. Это могут быть общие переменные, файлы, другие ресурсы. Когда несколько потоков одновременно пытаются изменить одну и ту же переменную или структуру данных, возникает опасная ситуация, называемая состоянием гонки (race condition). Для решения проблемы в языке Python применяется такой механизм блокировок как класс threading.Lock. Этот класс реализует объекты примитивной блокировки. Как только поток захватывает блокировку, последующие попытки захвата приводят к блокировке потока до тех пор, пока блокировка не будет освобождена.
Объект блокировки может находиться в одном из двух состояний: "захвачено" (locked) или "свободно" (unlocked). После создания объект блокировки изначально находится в состоянии "свободно".
Объект блокировки имеет два основных метода: acquire() и release(). Если блокировка свободна, метод acquire() переводит ее
в состояние "захвачено" (locked) и немедленно завершает выполнение. Если блокировка захвачена, метод acquire() блокирует поток до тех пор, пока вызов release()
в другом потоке не переведет блокировку в состояние "свободно". После этого acquire() снова переводит ее в состояние "захвачено" и завершает выполнение.
Метод release() следует вызывать только тогда, когда блокировка захвачена - он переводит ее в состояние "свободно" и немедленно завершает выполнение.
Попытка освободить свободную блокировку приводит к возникновению исключения RuntimeError.
Если в ожидании перехода блокировки в состояние "свободно" методом acquire() заблокировано несколько потоков, то после того, как вызов release()
переведет блокировку в состояние "свободно", выполнение продолжит только один из них. Выбор потока, который продолжит работу, не определен и может зависеть от конкретной реализации.
Прежде чем переходить к рассмотрению блокировок, сначала посмотрим на саму проблему. Допустим, у нас есть функционал банковского счета, с которого два потока одновременно пытаются списать деньги. Они считывают один баланс, уменьшают его и записывают обратно. Из-за того, что операционная система переключает потоки в случайный момент, один поток может затереть изменения другого. Например, рассмотрим следующую программку:
import threading
import time
x = 0 # общий ресурс
threads = [] # список потоков
def increase_x():
global x
x = 1 # устанавливаем новое значение для общего ресурса
thread_name = threading.current_thread().name # имя текущего потока
for _ in range(1, 4):
print(f"{thread_name}: {x}")
x += 1 # изменение общего ресурса
time.sleep(0.5) # имитация работы - полминуты
# запускаем три потока
for i in range(1, 4):
my_thread = threading.Thread(target=increase_x, name=f"Поток {i}")
threads.append(my_thread)
my_thread.start() # запускаем потоки
# ожидаем завершения потоков
for my_thread in threads:
my_thread.join()
Здесь у нас запускаются три потокоа, которые вызывают функцию increase_x() и которые работают с общей глобальной переменной x. И мы предполагаем, что функция выведет все значения x от 1 до 3.
И так для каждого потока. Однако в реальности в процессе работы будет происходить переключение между потоками, и значение переменной x становится непредсказуемым. Например, в моем случае я получил
следующий консольный вывод (он может в каждом конкретном случае различаться):
Поток 1: 1 Поток 2: 1 Поток 3: 1 Поток 2: 2 Поток 1: 2 Поток 3: 3 Поток 2: 5 Поток 1: 6 Поток 3: 7
В данном случае в качестве общегно ресурса применяется обычная переменная, но в реальности это может быть файл, сокет для отправки по сети, какой-то другой объект.
Решение проблемы состоит в том, чтобы синхронизировать потоки и ограничить доступ к разделяемым ресурсам на время их использования каким-нибудь потоком. Именно для этого Python предоставляет инструмент Lock (замок/блокировка). Его принцип работы прост:
Поток перед изменением переменной блокирует доступ к ней с помощью метода acquire().
Если другой поток тоже пытается изменить данные, он видит блокировку и ждет, пока она будет снята.
Первый поток меняет данные и выполняет разблокировку методом release()
Применим Lock:
import threading
import time
x = 0 # общий ресурс
x_lock = threading.Lock() # Создаем объект блокировки
threads = [] # список потоков
def increase_x():
global x
thread_name = threading.current_thread().name # имя текущего потока
x_lock.acquire() # Блокируем доступ для других потоков
try:
x = 1 # устанавливаем новое значение для общего ресурса
for _ in range(1, 4):
print(f"{thread_name}: {x}")
x += 1 # изменение общего ресурса
time.sleep(0.5) # имитация работы
finally:
x_lock.release() # обязательно разблокируем, даже если произошла ошибка
# запускаем три потока
for i in range(1, 4):
my_thread = threading.Thread(target=increase_x, name=f"Поток {i}")
threads.append(my_thread)
my_thread.start() # запускаем потоки
# ожидаем завершения потоков
for my_thread in threads:
my_thread.join()
Здесь для блокировки применяется объект x_lock - объект класса threading.Thread, в данном случае это переменная locker. И когда выполнение функции increase_x() в потоке
доходит до вызова x_lock.acquire(), объект x_lock блокируется, и на время его блокировки монопольный доступ к последующему коду вплоть до вызова x_lock.release() имеет только один поток. После выполнения x_lock.release() , объект x_lock освобождается и становится доступным для других потоков.
И в этом случае консольный вывод будет более упорядоченным:
Поток 1: 1 Поток 1: 2 Поток 1: 3 Поток 2: 1 Поток 2: 2 Поток 2: 3 Поток 3: 1 Поток 3: 2 Поток 3: 3
Чтобы не писать конструкции try...finally и случайно не забыть вызвать release() (что намертво заблокирует программу), в Python принято использовать менеджер контекста
with. Он автоматически заблокирует объект блокировки при входе и разблокирует его при выходе:
import threading
import time
x = 0 # общий ресурс
x_lock = threading.Lock() # Создаем объект блокировки
threads = [] # список потоков
def increase_x():
global x
thread_name = threading.current_thread().name # имя текущего потока
# Менеджер with сам вызовет acquire() и release()
with x_lock:
x = 1 # устанавливаем новое значение для общего ресурса
for _ in range(1, 4):
print(f"{thread_name}: {x}")
x += 1 # изменение общего ресурса
time.sleep(0.5) # имитация работы
# запускаем три потока
for i in range(1, 4):
my_thread = threading.Thread(target=increase_x, name=f"Поток {i}")
threads.append(my_thread)
my_thread.start() # запускаем потоки
# ожидаем завершения потоков
for my_thread in threads:
my_thread.join()
При работе с Lock следует учитывать, что хотя блокировки и гарантируют безопасность данных, но за это приходится платить скоростью. Пока один поток держит блокировку, остальные простаивают.
Поэтому лучше блокировать Lock на как можно меньшее время, например, на время непосредственного измненения общей переменной или обновления общего ресурса. И лучше не оборачивать в with весь тяжелый код (например, скачивание файла из сети и его запись). Иначе программа станет работать медленно, как в один поток.