【发布时间】:2020-05-12 12:45:31
【问题描述】:
在某些用例中,我需要等待所有已经创建的线程完成并根据它们的结果做出一些决定,看看我们是否需要进一步移动 - 没有ThreadPoolExecutor.shutdown()。
我是这样实现的:
from threading import BoundedSemaphore, Event
class JoinSemaphore(BoundedSemaphore):
def __init__(self, value=1):
super().__init__(value)
self._empty = Event()
def join(self, timeout=None):
if self._value < self._initial_value:
self._empty.wait(timeout)
def release(self):
with self._cond:
if self._value >= self._initial_value:
raise ValueError("Semaphore released too many times")
elif self._value == self._initial_value - 1:
self._empty.set()
self._value += 1
self._cond.notify()
def acquired(self):
with self._cond:
return self._initial_value - self._value
在这里我计算self._value < self._initial_value 没有任何保护,它有风险。
当我在计算self._value < self._initial_value 时编写如下所示的join() 函数以防止不必要的更改时,当主线程加入信号量时,我将面临死锁,此时其他线程并没有获得release() 锁定,因为主线程已经累积它,但主线程仍在等待其他线程。所以这个实现是不正确的。
def join(self, timeout=None):
with self._cond:
if self._value < self._initial_value:
self._empty.wait(timeout)
在第三个实现中,我不能保证当我想wait() 处理_empty 事件时,其他线程不可能提交他们的结果并释放锁。
def join(self, timeout=None):
with self._cond:
if self._value == self._initial_value:
return
self._empty.wait(timeout)
问题是:
如何正确使用锁计算self._value < self._initial_value 并等待_empty 事件并在等待它之前释放锁以避免死锁?
非常感谢
【问题讨论】:
-
能不能直接用自己的
threading.Lock包裹计数器的读写操作?
标签: python python-3.x multithreading semaphore python-multithreading