Condition не успевает проснуться после notify
Изучаю синхронизацию потоков в Python и в первом же примере столкнулся с проблемой. Обычный пример с производителем и покупателем
import threading
from time import sleep·
products = []
condition = threading.Condition()
class Customer(threading.Thread):
def run(self):
while True:
with condition:
while len(products) == 0:
condition.wait()
product = products.pop()
print('Get product %i ' % product)
class Producer(threading.Thread):
def run(self):
product = 0
while True:
product += 1
with condition:
sleep(2) # "make" product
products.append(product)
print('Append %i ' % product)
condition.notify()
customer = Customer()
producer = Producer()
customer.start()
producer.start()
Я ожидаю следующую логику работы кода: Customer начинает ждать, и когда Producer вызывает notify и выходит из контекста (with condition), Customer должен проснуться и захватить condition.
Однако этого не происходит. Producer захватывает condition повторно быстрее, и Customer не просыпается. Если немного помешать ему (Producer) захватывать condition, например добавить sleep(0.001) перед входом в контекст, то все начинает работать логично: Customer успевает проснуться и забрать товар. Таким образом получается состояние гонки, в которой Producer обычно выигрывает. Как грамотно избежать такого состояния? Если я например хочу, чтобы Customer забирал товар сразу после его появления в Producer.
UPD. Исправил опечатку (не было отступа)
with condition:
condition.wait_for(lambda : len(products) > 0)
-> product = products.pop()
Ответы (3 шт):
На самом деле тут гонки нет, Customer просыпается столько же раз, сколько было notify. Когда выплюнет результат зависит похоже от того в каком треде сейчас GIL.
Исправить можно если захватывать кондишн только в момент передачи результата:
import threading
from time import sleep
products = []
condition = threading.Condition()
from queue import Queue
q = Queue()
class Customer(threading.Thread):
def run(self):
while True:
with condition:
condition.wait_for(lambda : len(products) > 0)
product = products.pop()
q.put('Get product %i ' % product)
class Producer(threading.Thread):
def run(self):
product = 0
while True:
product += 1
sleep(2) # "make" product
with condition:
products.append(product)
q.put('Append %i ' % product)
condition.notify()
customer = Customer()
producer = Producer()
customer.start()
producer.start()
while True:
print(q.get())
Состояние под with синхронизируется и похоже не дает прыгать GIL.
Вы неправильно используете condition. Объект Condition является просто усовершенствованным вариантом объекта Event. Он тоже работает как коммуникатор между потоками и может применяться для уведомления notify() других потоков об изменении состояния программы. Например, его можно использовать для сигнализации доступности ресурса. Другие потоки также должны получать условие acquire() (и, следовательно, связанное с ним блокирование) до ожидания wait() для удовлетворения условия. Кроме того, поток должен освободить release() по условию Condition после завершения связанных с ним действий, так что другие потоки могут получить условие для своих целей.
Ваш код будет работать, если использовать указанное выше:
import threading
from time import sleep
products = []
condition = threading.Condition()
class Customer(threading.Thread):
def run(self):
while True:
condition.acquire()
condition.wait()
product = products.pop()
print('Get product %i ' % product)
condition.release()
class Producer(threading.Thread):
def run(self):
product = 0
while True:
product += 1
sleep(2) # "make" product
condition.acquire()
products.append(product)
condition.notify()
print('Append %i ' % product)
condition.release()
customer = Customer()
producer = Producer()
customer.start()
producer.start()
Вывод:
Append 1
Get product 1
Append 2
Get product 2
Append 3
Get product 3
Я ожидаю следующую логику работы кода: Customer начинает ждать, и когда Producer вызывает notify и выходит из контекста (with condition), Customer должен проснуться и захватить condition
Это необоснованное ожидание. notify всего лишь выводит поток из состояния блокировки, в которое он ушел после wait. Это не значит, что этот поток сразу же гарантированно получит свое время и будет первый в очереди на выполнение. Абсолютно нет, поток Producer-а будет работать и успешно захватывать lock ассоциированный с Condition пока не кончится его (Producer-а) квант времени (раньше это было 100 мс или до IO операции).
Более того, даже когда квант времени кончится, Consumer будет на общих основаниях конкурировать со всеми существующими потоками за право выполнятся, и может быть, что это право он получит не скоро (если готовых выполнятся потоков много).
И что еще усложняет тут ситуацию. Представьте, что квант времени Producer-а закончился и Consumer наконец-то может получить свой квант. Но если в этом момент блокировка захвачена Producer-ом (т.е. он находится в блоке with), то Consumer обратно уйдет в сон, а Producer получит новый квант времени. Это следствие того, что блокировка удерживается дольше чем нужно (см. ниже об этом).
Если вы хотите, чтобы получатель гарантированно обрабатывал сообщения по мере их создания, то вам именно что нужно исскуственно притормозить продюсера. Это называется backpressure - а именно, возможность Consumer-у просигнализировать Producer-у, что обработка не успевает за генерацией.
Есть много способов как это сделать. Можно смотреть на длину очереди, если длинная (или даже не пустая), то ждать пока это не изменится. Для этого может понадобиться отдельный Event или Condition, только уже в другую сторону - от Consumer-а к Producer-у.
И еще одно замечание. В данном примере condition используется для двух вещей:
- синхронизации доступа к
productsиз разных потоков - для отсылки уведомления
Удерживать блокировку ассоциированную с condition нужно только, когда идет модификация или чтение products и посылка уведомления. В другое время это только будет тормозить Consumer-a без необходимости.