【问题标题】:How to create a joinable semaphore in python?如何在 python 中创建可连接的信号量?
【发布时间】: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 &lt; self._initial_value 没有任何保护,它有风险。 当我在计算self._value &lt; 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 &lt; self._initial_value 并等待_empty 事件并在等待它之前释放锁以避免死锁?

非常感谢

【问题讨论】:

  • 能不能直接用自己的threading.Lock包裹计数器的读写操作?

标签: python python-3.x multithreading semaphore python-multithreading


【解决方案1】:

这是我为这个问题找到的一种方法,我想从基本 Semaphore 类继承,但找不到解决问题的简单解决方案,所以我决定编辑基本 BoundSemaphore 类:

from threading import Event, Condition, Lock
from time import monotonic as _time


class JoinSemaphore:

    def __init__(self, value=1):
        if value < 0:
            raise ValueError("semaphore initial value must be >= 0")
        self._cond = Condition(Lock())
        self._value = value
        self._initial_value = value
        self._empty = Event()
        self._empty.set()

    def acquire(self, blocking=True, timeout=None):
        if not blocking and timeout is not None:
            raise ValueError("can't specify timeout for non-blocking acquire")
        rc = False
        endtime = None
        with self._cond:
            while self._value == 0:
                if not blocking:
                    break
                if timeout is not None:
                    if endtime is None:
                        endtime = _time() + timeout
                    else:
                        timeout = endtime - _time()
                        if timeout <= 0:
                            break
                self._cond.wait(timeout)
            else:
                self._empty.clear()
                self._value -= 1
                rc = True
        return rc

    __enter__ = acquire

    def join(self, timeout=None):
        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

差异是:

class JoinSemaphore:

    def __init__(self, value=1):
        ...
        self._empty = Event()
        self._empty.set()

    def acquire(self, blocking=True, timeout=None):
        ...
        with self._cond:
            while ...:
        ...
            else:
                self._empty.clear()
        ...

    def join(self, timeout=None):
        self._empty.wait(timeout)

    def release(self):
        ...
            elif self._value == self._initial_value - 1:
                self._empty.set()
        ...

    def acquired(self):
        with self._cond:
            return self._initial_value - self._value

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-05-30
    • 2015-04-24
    • 1970-01-01
    • 1970-01-01
    • 2017-03-26
    • 1970-01-01
    • 2015-11-16
    相关资源
    最近更新 更多