【问题标题】:Inheritance in iterable implementation of python's multiprocessing.Queuepython的multiprocessing.Queue的可迭代实现中的继承
【发布时间】:2019-12-04 01:57:12
【问题描述】:

我发现 python 的 multiprocessing.Queue 缺少默认实现,因为它不像任何其他集合一样可迭代。所以我开始努力创建它的“子类”,添加功能。从下面的代码中可以看出,它不是正确的子类,因为multiprocess.Queue 本身不是直接类,而是工厂函数,真正的底层类是multiprocess.queues.Queue。我没有理解也没有努力去模仿工厂函数,以便我可以正确地从类继承,所以我只是让新类从工厂创建它自己的实例并将其视为超类。这是代码;

from multiprocessing import Queue, Value, Lock
import queue

class QueueClosed(Exception):
    pass

class IterableQueue:
    def __init__(self, maxsize=0):
        self.closed = Value('b', False)
        self.close_lock = Lock()
        self.queue = Queue(maxsize)

    def close(self):
        with self.close_lock:
            self.closed.value = True
            self.queue.close()

    def put(self, elem, block=True, timeout=None):
        with self.close_lock:
            if self.closed.value:
                raise QueueClosed()
            else:
                self.queue.put(elem, block, timeout)

    def put_nowait(self, elem):
        self.put(elem, False)

    def get(self, block=True):
        if not block:
            return self.queue.get_nowait()
        elif self.closed.value:
            try:
                return self.queue.get_nowait()
            except queue.Empty:
                return None
        else:
            val = None
            while not self.closed.value:
                try:
                    val = self.queue.get_nowait()
                    break
                except queue.Empty:
                    pass
            return val

    def get_nowait(self):
        return self.queue.get_nowait()

    def join_thread(self):
        return self.queue.join_thread()

    def __iter__(self):
        return self

    def __next__(self):
        val = self.get()
        if val == None:
            raise StopIteration()
        else:
            return val

    def __enter__(self):
        return self

    def __exit__(self, *args):
        self.close()

这让我可以像普通的 multiprocessing.Queue 一样实例化一个 IterableQueue 对象,像往常一样将元素放入其中,然后在子消费者内部,像这样简单地循环它;

from iterable_queue import IterableQueue
from multiprocessing import Process, cpu_count
import os

def fib(n):
    if n < 2:
        return n
    return fib(n-1) + fib(n-2)

def consumer(queue):
    print(f"[{os.getpid()}] Consuming")
    for i in queue:
        print(f"[{os.getpid()}] < {i}")
        n = fib(i)
        print(f"[{os.getpid()}] {i} > {n}")
    print(f"[{os.getpid()}] Closing")

def producer():
    print("Enqueueing")
    with IterableQueue() as queue:
        procs = [Process(target=consumer, args=(queue,)) for _ in range(cpu_count())]
        [p.start() for p in procs]
        [queue.put(i) for i in range(36)]
    print("Finished")

if __name__ == "__main__":
    producer()

它几乎可以无缝运行;一旦队列关闭,消费者退出循环,但只有在耗尽所有剩余元素之后。但是,我对缺少继承方法感到不满意。为了模仿实际的继承行为,我尝试将以下元函数调用添加到类中;

def __getattr__(self, name):
    if name in self.__dict__:
        return self.__dict__[name]
    else:
        return self.queue.__getattr__[name]

但是,当 IterableQueue 类的实例在子 multiprocessing.Process 线程中被操作时,这会失败,因为该类的 __dict__ 属性未保留在其中。我试图通过用multiprocessing.Manager().dict() 替换类的默认__dict__ 来以一种骇人听闻的方式解决这个问题,就像这样;

def __init__(self, maxsize=0):
    self.closed = Value('b', False)
    self.close_lock = Lock()
    self.queue = Queue(maxsize)
    self.__dict__ = Manager().dict(self.__dict__)

但是在这样做时,我收到了一条错误消息,指出 RuntimeError: Synchronized objects should only be shared between processes through inheritance。所以我的问题是,我应该如何正确地从 Queue 类继承,以便子类继承对它所有属性的访问?此外,当队列为空但未关闭时,消费者都处于忙碌的循环中,而不是真正的 IO 块,占用了宝贵的 cpu 资源。如果您对我在使用此代码时可能遇到的并发和竞争条件问题有任何建议,或者我可能如何解决繁忙循环问题,我也愿意接受其中的建议。


基于 MisterMiyagi 提供的代码,我创建了这个通用的IterableQueue 类,它可以接受任意输入,正确阻塞,并且不会在队列关闭时挂起;

from multiprocessing.queues import Queue
from multiprocessing import get_context

class QueueClosed(Exception):
    pass

class IterableQueue(Queue):
    def __init__(self, maxsize=0, *, ctx=None):
        super().__init__(
            maxsize=maxsize,
            ctx=ctx if ctx is not None else get_context()
        )

    def close(self):
        super().put((None, False))
        super().close()

    def __iter__(self):
        return self

    def __next__(self):
        try:
            return self.get()
        except QueueClosed:
            raise StopIteration

    def get(self, *args, **kwargs):
        result, is_open = super().get(*args, **kwargs)
        if not is_open:
            super().put((None, False))
            raise QueueClosed
        return result

    def put(self, val, *args, **kwargs):
        super().put((val, True), *args, **kwargs)

    def __enter__(self):
        return self

    def __exit__(self, *args):
        self.close()

【问题讨论】:

  • 您需要对IterableQueue 实例进行哪些操作?您是否在一个进程中设置属性并在另一个进程中读取它们?这背后的用例是什么?
  • 我主要不是更简化的生产者-消费者模式方法。我可以启动一个消费者线程,让它在队列上的一个可迭代循环中完成整个操作,只要提供了元素,它就会继续使用它们。然后一旦完成,它会自动退出。

标签: python-3.x multiprocessing queue subclass iterable


【解决方案1】:

multiprocess.Queue 包装器仅用于 use the default context

def Queue(self, maxsize=0):
    '''Returns a queue object'''
    from .queues import Queue
    return Queue(maxsize, ctx=self.get_context())

继承时,您可以在__init__ 方法中复制它。这允许您继承整个 Queue 行为。您只需要添加迭代器方法:

from multiprocessing.queues import Queue
from multiprocessing import get_context


class IterableQueue(Queue):
    """
    ``multiprocessing.Queue`` that can be iterated to ``get`` values

    :param sentinel: signal that no more items will be received
    """
    def __init__(self, maxsize=0, *, ctx=None, sentinel=None):
        self.sentinel = sentinel
        super().__init__(
            maxsize=maxsize,
            ctx=ctx if ctx is not None else get_context()
        )

    def close(self):
        self.put(self.sentinel)
        # wait until buffer is flushed...
        while self._buffer:
            time.sleep(0.01)
        # before shutting down the sender
        super().close()

    def __iter__(self):
        return self

    def __next__(self):
        result = self.get()
        if result == self.sentinel:
            # re-queue sentinel for other listeners
            self.put(result)
            raise StopIteration
        return result

请注意,表示队列结束的sentinel 是通过相等性比较的,因为身份不会跨进程保留。常用的queue.Queue sentinel object() 不能正常工作。

【讨论】:

  • 我将您的代码与我在问题中的示例生产者-消费者代码一起使用作为演示,当消费者到达队列末尾时,他们挂起而不是退出。 Queue.close() 上没有引发异常。
  • @Maurdekye 我得检查一下。有更新版本时会通知您。
  • 这会处理上述有多个消费者在队列上监听的情况吗?
  • 我可以更改代码以通过使用元组来强制交互,允许在 put/get 上给出一对值和isclosed,允许任意哨兵。但除此之外,这似乎解决了我之前遇到的所有问题。我会给它一个检查,确认后我会选择你的答案。
  • 等等,别在意最后一条评论;我想我上次设置的程序错误,并且破坏了一些功能。这次成功了。
【解决方案2】:

遍历用于线程/进程之间通信的队列没有任何意义。

消费和迭代本质上是两件完全不同的事情,这就是为什么需要同步的队列不允许您对其自身进行迭代。

您可以使用iter() 函数让您在收到的项目进入时对其进行迭代:

for item in iter(queue.get, None):
    print(item)

一旦收到None,它将停止迭代,但您可以在此处放置任何表明退出条件的内容,特别是如果 None 可能是您队列中的有效值。

【讨论】:

  • 这并没有真正回答我的问题,它只是告诉我我的代码是错误的,不应该被允许工作。尽管如此,它做得非常好。
  • 我的解决方案更复杂,因为它处理了消费者被阻塞在队列中等待,然后队列关闭的情况。一般情况下消费者是无从知晓的,只能等到永远。我的代码将消费者置于一个忙碌的循环中,不断检查队列是否已关闭,如果是则中断。
  • 基本上,如果我可以从生产者进程中进行一些调用,立即在所有被阻止的消费者中引发一些异常,并在所有未来尝试get 的消费者中引发相同的异常,我的问题将得到解决.但我不知道如何让它发挥作用,或者是否有可能。
  • 如果它运行良好,那么就没有问题,也没有什么我可以添加的,看起来你希望队列为你完成所有繁重的工作,而不是塞满你的将自己的代码放入队列中,这与您最初期望的迭代无关,您应该只编写自己的代码来执行您希望它执行的操作,而不是将其隐藏在“for循环”后面,就像您看到了,几乎无法调试。
  • 一方面,“把自己的代码塞进队列,与迭代无关”,正是我想要做的,因为这就是代码的封装和模块化的概念是关于。重点是创建一个统一的队列类,当我需要使用多进程队列时,它更易于使用。出于模糊、难以描述的原因说我“不应该”不会让你有任何收获。它工作得很好,是的;在最低限度的可接受水平,我想改进。你读过这篇文章吗?
猜你喜欢
  • 2014-07-18
  • 1970-01-01
  • 2019-09-14
  • 1970-01-01
  • 1970-01-01
  • 2017-12-22
  • 1970-01-01
  • 2016-04-22
  • 2015-04-13
相关资源
最近更新 更多