【问题标题】:Python multiprocessing pool map_async freezesPython多处理池map_async冻结
【发布时间】:2018-05-03 19:05:36
【问题描述】:

我有一个包含 80,000 个字符串的列表,我正在通过话语解析器运行,为了提高这个过程的速度,我一直在尝试使用 python 多处理包。

解析器代码需要 python 2.7,我目前正在使用字符串子集在 2 核 Ubuntu 机器上运行它。对于短列表,即 20,该进程在两个内核上运行都没有问题,但是如果我运行大约 100 个字符串的列表,两个工作人员将在不同点冻结(因此在某些情况下,工作人员 1 直到几分钟才会停止在工人 2) 之后。这发生在所有字符串完成并返回任何内容之前。每次使用相同的映射函数时内核在同一点停止,但如果我尝试不同的映射函数,这些点会有所不同,即 map vs map_async vs imap。

我尝试删除那些索引处的字符串,这没有任何影响,并且这些字符串在较短的列表中运行良好。根据我包含的打印语句,当进程似乎冻结时,当前字符串的当前迭代似乎已经完成,它只是不会移动到下一个字符串。大约需要一个小时的运行时间才能到达两个工人都被冻结的地方,而我无法在更短的时间内重现该问题。涉及多处理命令的代码是:

def main(initial_file, chunksize = 2):
    entered_file = pd.read_csv(initial_file)
    entered_file = entered_file.ix[:, 0].tolist()

    pool = multiprocessing.Pool()

    result = pool.map_async(discourse_process, entered_file, chunksize = chunksize)

    pool.close()
    pool.join()

    with open("final_results.csv", 'w') as file:
        writer = csv.writer(file)
        for listitem in result.get():
            writer.writerow([listitem[0], listitem[1]])

if __name__ == '__main__':
    main(sys.argv[1])

当我使用 Ctrl-C 停止进程时(并不总是有效),我收到的错误消息是:

^CTraceback (most recent call last):
  File "Combined_Script.py", line 94, in <module>
    main(sys.argv[1])
  File "Combined_Script.py", line 85, in main
    pool.join()
  File "/usr/lib/python2.7/multiprocessing/pool.py", line 474, in join
    p.join()
  File "/usr/lib/python2.7/multiprocessing/process.py", line 145, in join
    res = self._popen.wait(timeout)
  File "/usr/lib/python2.7/multiprocessing/forking.py", line 154, in wait
    return self.poll(0)
  File "/usr/lib/python2.7/multiprocessing/forking.py", line 135, in poll
    pid, sts = os.waitpid(self.pid, flag)
KeyboardInterrupt
Process PoolWorker-1:
Traceback (most recent call last):
  File "/usr/lib/python2.7/multiprocessing/process.py", line 258, in _bootstrap
    self.run()
  File "/usr/lib/python2.7/multiprocessing/process.py", line 114, in run
    self._target(*self._args, **self._kwargs)
  File "/usr/lib/python2.7/multiprocessing/pool.py", line 117, in worker
    put((job, i, result))
  File "/usr/lib/python2.7/multiprocessing/queues.py", line 390, in put
    wacquire()
KeyboardInterrupt
^CProcess PoolWorker-2:
Traceback (most recent call last):
  File "/usr/lib/python2.7/multiprocessing/process.py", line 258, in _bootstrap
    self.run()
  File "/usr/lib/python2.7/multiprocessing/process.py", line 114, in run
    self._target(*self._args, **self._kwargs)
  File "/usr/lib/python2.7/multiprocessing/pool.py", line 117, in worker
    put((job, i, result))
  File "/usr/lib/python2.7/multiprocessing/queues.py", line 392, in put
    return send(obj)
KeyboardInterrupt
Error in atexit._run_exitfuncs:
Traceback (most recent call last):
  File "/usr/lib/python2.7/atexit.py", line 24, in _run_exitfuncs
    func(*targs, **kargs)
  File "/usr/lib/python2.7/multiprocessing/util.py", line 305, in _exit_function
    _run_finalizers(0)
  File "/usr/lib/python2.7/multiprocessing/util.py", line 274, in _run_finalizers
    finalizer()
  File "/usr/lib/python2.7/multiprocessing/util.py", line 207, in __call__
    res = self._callback(*self._args, **self._kwargs)
  File "/usr/lib/python2.7/multiprocessing/pool.py", line 500, in _terminate_pool
    outqueue.put(None)                  # sentinel
  File "/usr/lib/python2.7/multiprocessing/queues.py", line 390, in put
    wacquire()
KeyboardInterrupt
Error in sys.exitfunc:
Traceback (most recent call last):
  File "/usr/lib/python2.7/atexit.py", line 24, in _run_exitfuncs
    func(*targs, **kargs)
  File "/usr/lib/python2.7/multiprocessing/util.py", line 305, in _exit_function
    _run_finalizers(0)
  File "/usr/lib/python2.7/multiprocessing/util.py", line 274, in _run_finalizers
    finalizer()
  File "/usr/lib/python2.7/multiprocessing/util.py", line 207, in __call__
    res = self._callback(*self._args, **self._kwargs)
  File "/usr/lib/python2.7/multiprocessing/pool.py", line 500, in _terminate_pool
    outqueue.put(None)                  # sentinel
  File "/usr/lib/python2.7/multiprocessing/queues.py", line 390, in put
    wacquire()
KeyboardInterrupt

当我使用 htop 在另一个命令窗口中查看内存时,一旦工作人员冻结,内存就会低于 3%。这是我第一次尝试并行处理,我不确定我还缺少什么?

【问题讨论】:

    标签: python python-2.7 multiprocessing python-multiprocessing


    【解决方案1】:

    您可以为您的流程定义一个时间来返回结果,否则会引发错误:

    try:
        result.get(timeout = 1)
    except multiprocessing.TimeoutError:
        print("Error while retrieving the result")
    

    您还可以验证您的流程是否成功

    import time
    while True:
        try:
            result.succesful()
        except Exception:
            print("Result is not yet succesful")
        time.sleep(1)
    

    最后,查看https://docs.python.org/2/library/multiprocessing.html 很有帮助。

    【讨论】:

      【解决方案2】:

      我无法解决多处理池的问题,但是我遇到了 loky 包,并且能够使用它通过以下几行运行我的代码:

      executor = loky.get_reusable_executor(timeout = 200, kill_workers = True)
      results = executor.map(discourse_process, entered_file)
      

      【讨论】:

        猜你喜欢
        • 2014-08-25
        • 2012-12-04
        • 2012-04-03
        • 1970-01-01
        • 2021-09-30
        • 2013-05-17
        • 1970-01-01
        • 2018-01-30
        • 2019-08-14
        相关资源
        最近更新 更多