【问题标题】:What data structure should I be using to store results returned by processes running in parallel?我应该使用什么数据结构来存储并行运行的进程返回的结果?
【发布时间】:2014-05-26 23:42:30
【问题描述】:

我正在使用multiprocessing 模块并行运行一个函数。我并行运行的进程本身具有并行运行的子进程,因此我不能使用 pool 类,除了最低级别的子进程(即那些不再创建子进程的子进程),因为守护进程可能不会创造孩子。

因此,我使用Process,并手动管理运行和加入进程。

今天,我花了很长时间尝试调试我的代码,因为它似乎挂了,但我不明白在哪里。经过一些调试,我最终发现当我的一个函数试图将数据放入其中时,我用来存储结果的 multiprorcessing.Queue 对象无限期地阻塞。我还不确定为什么它会无限期地阻塞,但我已经确认这是问题,因为删除 put 命令允许继续执行(虽然,我没有得到任何数据)。

multiprorcessing.Queue 是用于存储并行运行的函数返回的信息的正确对象吗?

【问题讨论】:

    标签: python-2.7 multiprocessing


    【解决方案1】:

    多进程系统难以调试。我建议彻底调试所有底层函数,然后添加一点 multiproc 来品尝。

    因为 multiproc 可能会令人困惑,我建议尽早并经常记录。如果你记录太多,很容易去掉绒毛。但是,如果您在日志中看不到奇怪的极端情况,那么事情可能会......困难:)

    worker 用它的名字记录,而 parent 记录为“MainProcess”。正常的东西被记录为 INFO,可怕的问题被记录为 ERROR。

    为了模拟“不能将东西放入我的队列”错误,我创建了输出 Queue 以仅容纳一项。还有其他代码会监视它,并专门记录它。 (将 if 0 更改为 if 1 以获取完整的代码以运行。)

    玩得开心!

    import logging, multiprocessing, Queue
    
    def myproc(arg):
        return arg*2
    
    def worker(inqueue, outqueue):
        mylog = multiprocessing.get_logger()
        mylog.info('start')
        for job in iter(inqueue.get, 'STOP'):
            mylog.info('got %s', job)
            try:
                outqueue.put( myproc(job), timeout=1 )
            except Queue.Full:
                mylog.error('queue full!')
    
        mylog.info('done')
    
    
    logger = multiprocessing.log_to_stderr(
        level=logging.INFO,
    )
    logger.info('setup')
    
    inqueue, outqueue = multiprocessing.Queue(), multiprocessing.Queue()
    if 1:                           # debug 'queue full!' issues
        outqueue = multiprocessing.Queue(maxsize=1)
    # prefill with 3 jobs
    for num in range(3):
        inqueue.put(num)
    # signal end of jobs
    inqueue.put('STOP')
    
    worker_p = multiprocessing.Process(
        target=worker, args=(inqueue, outqueue),
        name='worker',
    )
    worker_p.start()
    
    worker_p.join()
    
    logger.info('done')
    

    示例运行:

    [INFO/MainProcess] setup
    [INFO/worker] child process calling self.run()
    [INFO/worker] start
    [INFO/worker] got 0
    [INFO/worker] got 1
    [ERROR/worker] queue full!
    [INFO/worker] got 2
    [ERROR/worker] queue full!
    [INFO/worker] done
    [INFO/worker] process shutting down
    [INFO/worker] process exiting with exitcode 0
    [INFO/MainProcess] done
    [INFO/MainProcess] process shutting down
    

    【讨论】:

    • 亲爱的 shavenwarthog -- 非常感谢您的回答!您关于使用 logging 函数的提示似乎特别有用。顺便说一句,为了调试我上面提到的问题,我最终做了一些类似(但更混乱)的事情。现在我知道我的队列被阻塞了(而不是因为它已满!)当我尝试将某些东西放到它上面时——你仍然建议我使用队列来存储 multiproc 函数的返回值吗?
    • P.S.我知道这不是 Queue.Full 错误,因为我没有为队列指定最大大小。
    • 不客气!我发现记录一切对于跟踪奇怪的错误非常有价值。现在想起来,Queue 不适合存储数据。它有一些开销,因为它被监视更改并且在推送/弹出时通知其他进程。考虑另一个使用数据然后将其写入文件或对其进行汇总或将其粘贴到list 中的工作人员。
    • shavenwarthog,我使用“经理名单”取得了不错的成绩;即我创建了一个manager object 并用它来生成一个我传递的列表!所以,我认为为了做我想做的事情,如果用例超级简单,最好使用shared memory objects,或者如果要处理的数据对象更复杂,最好使用管理器.
    猜你喜欢
    • 1970-01-01
    • 2013-04-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-04-25
    • 2011-04-19
    • 1970-01-01
    相关资源
    最近更新 更多