【问题标题】:Unable to .get() from multiprocessing.Queue无法从 multiprocessing.Queue 获取 .get()
【发布时间】:2014-12-10 05:14:29
【问题描述】:

我正在构建一个 Web 应用程序来处理大约 60,000 个(并且还在增长的)大文件,执行一些分析并返回一个需要用户验证的“最佳猜测”。这些文件将按类别细化以避免加载每个文件,但我仍然面临一次可能需要处理 1000 多个文件的情况。

这些是大型文件,每个文件最多可能需要 8-9 秒来处理,在 1000 多个文件的情况下,让用户在审查之间等待 8 秒或在文件处理之前等待 2 小时以上是不切实际的。

为了克服这个问题,我决定使用多处理来生成多个工作人员,每个工作人员将从文件队列中挑选、处理它们并插入到输出队列中。我有另一种方法,它基本上轮询输出队列中的项目,然后在可用时将它们流式传输到客户端。

这很有效,直到队列任意停止返回项目的一部分。我们在环境中将 gevent 与 Django 和 uwsgi 一起使用,我知道在 gevent 的上下文中通过多处理创建子进程会在子进程中产生不希望的事件循环状态。在分叉之前产生的小绿叶在孩子中复制。因此,我决定使用gipc 来协助处理子进程。

我的代码的示例(我无法发布我的实际代码):

import multiprocessing
import gipc
from item import Item

MAX_WORKERS = 10

class ProcessFiles(object):

    def __init__(self):
        self.input_queue = multiprocessing.Queue()
        self.output_queue = multiprocessing.Queue()
        self.file_count = 0

    def query_for_results(self):
        # Query db for records of files to process.
        # Return results and set self.file_count equal to
        # the number of records returned.
        pass

    # The subprocess.
    def worker(self):
        # Chisel away at the input queue until no items remain.
        while True:
            if self.no_items_remain():
                return

            item = self.input_queue.get(item)
            item.process()
            self.output_queue.put(item)

    def start(self):
        # Get results and store in Queue for processing
        results = self.query_for_results()
        for result in results:
             item = Item(result)
             self.input_queue.put(item)

        # Spawn workers to process files.
        for _ in xrange(MAX_WORKERS):
            process = gipc.start_process(self.worker)

        # Poll for items to send to client.
        return self.get_processed_items()

    def get_processed_items(self):

        # Wait for the output queue to hold at least 1 item.
        # When an item becomes available, yield it to client.
        count = 0
        while count != self.file_count:
            #item = self._get_processed_item()
            # Debugging:
            try:
                item = self.output_queue.get(timeout=1)
            except:
                print '\tError fetching processed item. Retrying...'
                continue

            if item:
                print 'QUEUE COUNT: {}'.format(self.output_queue.qsize())
                count += 1
                yield item
        yield 'end'

我希望输出在处理并产生一个项目后显示队列的当前计数:

QUEUE COUNT: 999
QUEUE COUNT: 998
QUEUE COUNT: 997
QUEUE COUNT: 996
...
QUEUE COUNT: 4
QUEUE COUNT: 3
QUEUE COUNT: 2
QUEUE COUNT: 1

但是,该脚本在失败之前只能生成一些项目:

QUEUE COUNT: 999
QUEUE COUNT: 998
QUEUE COUNT: 997
QUEUE COUNT: 996
    Error fetching processed item. Retrying...
    Error fetching processed item. Retrying...
    Error fetching processed item. Retrying...
    Error fetching processed item. Retrying...
    Error fetching processed item. Retrying...
    Error fetching processed item. Retrying...
    ...

我的问题是:到底发生了什么?为什么我不能从队列中get?我怎样才能退回我期望的物品并避免这种情况?

【问题讨论】:

    标签: python queue multiprocessing gevent gipc


    【解决方案1】:

    当您无法获得物品时引发的实际异常是什么?您盲目地捕获所有可能引发的异常。此外,为什么不直接使用get 而没有超时呢?您立即重试,无需执行任何其他操作。可能只是调用获取一个块,直到一个项目准备好。

    关于这个问题,我认为正在发生的事情是gipc 正在关闭与您的队列关联的管道,从而破坏了队列。我预计会抛出OSError 而不是queue.Empty。有关详细信息,请参阅此bug report。

    作为替代方案,您可以使用进程池,在任何gevent 发生之前启动池(这意味着您不必担心事件循环问题)。使用imap_unordered 将作业提交到池中,您应该没问题。

    你的启动函数看起来像:

    def start(self):
        results = self.query_for_results()
        return self.pool.imap_unordered(self.worker, results, 
            chunksize=len(results) // self.num_procs_in_pool)
    
    @staticmethod
    def worker(item):
        item.process()
        return item
    

    【讨论】:

    • 感谢您的回答,但是因为这是一个 Web 应用程序,并且因为我们在 Django 的上下文中使用带有 uwsgi 的 gevent 来拥有多个 django 工作人员,所以 gevent 循环在 Django 服务器启动时启动.当我的视图被实例化时,循环已经运行了很长一段时间,这意味着在 gevent 循环开始之前创建一个工作池是不可能的。
    • 要评论这个错误,它实际上是一个Queue.Empty 异常,因为它试图在队列中放入任何内容之前从队列中获取。
    猜你喜欢
    • 2018-03-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-12-02
    • 1970-01-01
    • 1970-01-01
    • 2016-11-14
    • 2017-07-16
    相关资源
    最近更新 更多