【问题标题】:Exception thrown in multiprocessing Pool not detected未检测到多处理池中引发的异常
【发布时间】:2011-07-18 02:46:12
【问题描述】:

似乎当从 multiprocessing.Pool 进程引发异常时,没有堆栈跟踪或任何其他表明它失败的迹象。示例:

from multiprocessing import Pool 

def go():
    print(1)
    raise Exception()
    print(2)

p = Pool()
p.apply_async(go)
p.close()
p.join()

打印 1 并静默停止。有趣的是,引发 BaseException 反而有效。有没有办法让所有异常的行为都与 BaseException 相同?

【问题讨论】:

  • 我遇到了同样的问题。原因如下:工作进程捕捉到异常,并将失败代码和异常放在结果队列中。回到主进程,Pool 的结果处理线程获取失败代码并忽略它。某种猴子补丁调试模式可能是可能的。另一种方法是确保您的工作函数捕获任何异常,返回它并为您的处理程序打印一个错误代码。
  • 这里已经回答了这个问题:stackoverflow.com/a/26096355/512111

标签: python exception multiprocessing


【解决方案1】:

也许我遗漏了一些东西,但这不是 Result 对象的 get 方法返回的内容吗?见Process Pools

类 multiprocessing.pool.AsyncResult

Pool.apply_async() 和 Pool.map_async().get([timeout]) 返回结果的类
到达时返回结果。如果 timeout 不是 None 并且结果没有在 timeout 秒然后 multiprocessing.TimeoutError 被提出。如果遥控器 call 引发异常,然后 get() 将重新引发该异常。

所以,稍微修改一下你的例子,就可以了

from multiprocessing import Pool

def go():
    print(1)
    raise Exception("foobar")
    print(2)

p = Pool()
x = p.apply_async(go)
x.get()
p.close()
p.join()

结果是什么

1
Traceback (most recent call last):
  File "rob.py", line 10, in <module>
    x.get()
  File "/usr/lib/python2.6/multiprocessing/pool.py", line 422, in get
    raise self._value
Exception: foobar

这并不完全令人满意,因为它不打印回溯,但总比没有好。

更新:此错误已在 Python 3.4 中得到修复,由 Richard Oudkerk 提供。请参阅问题get method of multiprocessing.pool.Async should return full traceback

【讨论】:

  • 如果您弄清楚为什么它不返回回溯,请告诉我。既然它能够返回错误值,它也应该能够返回回溯。我可能会在一些合适的论坛上提问——也许是一些 Python 开发列表。顺便说一句,正如您可能已经猜到的那样,我在试图找出同样的事情时遇到了您的问题。 :-)
  • 注意:要为一堆同时运行的任务执行此操作,您应该将所有结果保存在一个列表中,然后使用 get() 遍历每个结果,如果不这样做,可能会被 try/catch 包围'不想在第一个错误上废话。
  • @dfrankow 这是一个很好的建议。您是否愿意在新答案中提出可能的实现?我打赌它会非常有用。 ;)
  • 遗憾的是,一年多之后,我完全忘记了这一切。
  • 答案中的代码将等待x.get(),这破坏了异步应用任务的意义。 @dfrankow 关于将结果保存到列表然后在最后getting 他们的评论是一个更好的解决方案。
【解决方案2】:

我有一个合理的解决方案,至少用于调试目的。我目前没有可以在主进程中引发异常的解决方案。我的第一个想法是使用装饰器,但你只能腌制functions defined at the top level of a module,所以就这样了。

取而代之的是一个简单的包装类和一个将其用于apply_async(以及因此apply)的Pool 子类。我会把map_async留给读者作为练习。

import traceback
from multiprocessing.pool import Pool
import multiprocessing

# Shortcut to multiprocessing's logger
def error(msg, *args):
    return multiprocessing.get_logger().error(msg, *args)

class LogExceptions(object):
    def __init__(self, callable):
        self.__callable = callable

    def __call__(self, *args, **kwargs):
        try:
            result = self.__callable(*args, **kwargs)

        except Exception as e:
            # Here we add some debugging help. If multiprocessing's
            # debugging is on, it will arrange to log the traceback
            error(traceback.format_exc())
            # Re-raise the original exception so the Pool worker can
            # clean up
            raise

        # It was fine, give a normal answer
        return result

class LoggingPool(Pool):
    def apply_async(self, func, args=(), kwds={}, callback=None):
        return Pool.apply_async(self, LogExceptions(func), args, kwds, callback)

def go():
    print(1)
    raise Exception()
    print(2)

multiprocessing.log_to_stderr()
p = LoggingPool(processes=1)

p.apply_async(go)
p.close()
p.join()

这给了我:

1
[ERROR/PoolWorker-1] Traceback (most recent call last):
  File "mpdebug.py", line 24, in __call__
    result = self.__callable(*args, **kwargs)
  File "mpdebug.py", line 44, in go
    raise Exception()
Exception

【讨论】:

  • 太糟糕了,没有更简单的解决方案(或者我的错误),但这将完成工作 - 谢谢!
  • 我已经意识到可以使用装饰器,如果你使用@functools.wraps(func) 来装饰你的包装器。这使您的装饰函数看起来像定义在模块顶层的函数。
  • this answer中的解决方案比较简单并且支持在主进程中重提错误!
  • @j08lue - 这个答案很好,但有 3 个缺点:1)额外的依赖 2)必须用 try/except 包装你的工作函数,并且返回包装器对象的逻辑 3)必须嗅探返回类型并重新加注。从好的方面来说,在你的主线程中获得实际的回溯会更好,我同意。
  • @RupertNash 我的意思实际上更像是this new answer 中的用法。这解决了缺点 3。
【解决方案3】:

撰写本文时得票最多的解决方案有问题:

from multiprocessing import Pool

def go():
    print(1)
    raise Exception("foobar")
    print(2)

p = Pool()
x = p.apply_async(go)
x.get()  ## waiting here for go() to complete...
p.close()
p.join()

正如@dfrankow 所说,它将等待x.get(),这会破坏异步运行任务的意义。因此,为了提高效率(特别是如果您的工作函数 go 需要很长时间),我会将其更改为:

from multiprocessing import Pool

def go(x):
    print(1)
    # task_that_takes_a_long_time()
    raise Exception("Can't go anywhere.")
    print(2)
    return x**2

p = Pool()
results = []
for x in range(1000):
    results.append( p.apply_async(go, [x]) )

p.close()

for r in results:
     r.get()

优点:worker函数是异步运行的,所以如果你在多个内核上运行许多任务,它会比原来的解决方案效率高很多。

缺点:如果worker函数中出现异常,只有在池完成所有任务后才会引发。这可能是也可能不是理想的行为。 根据@colinfang 的评论编辑,修复了这个问题。

【讨论】:

  • 努力。但是,由于您的示例是基于有多个结果的假设,也许可以稍微扩展一下,以便实际上有多个结果?另外,您写道:“特别是如果您是工作人员”。那应该是“你的”。
  • 你是对的,谢谢。我稍微扩展了这个例子。
  • 酷。此外,您可能想要尝试/排除,具体取决于您希望如何容忍 fetch 中的错误。
  • @gozzilli 你能把for r in ... r.get() 放在p.close()p.join() 之间,这样一遇到异常就退出
  • @colinfang 我相信会return null,因为计算还没有发生——除非你join(),否则它不会等待它。
【解决方案4】:

我已经用这个装饰器成功记录了异常:

import traceback, functools, multiprocessing

def trace_unhandled_exceptions(func):
    @functools.wraps(func)
    def wrapped_func(*args, **kwargs):
        try:
            func(*args, **kwargs)
        except:
            print 'Exception in '+func.__name__
            traceback.print_exc()
    return wrapped_func

问题中的代码是

@trace_unhandled_exceptions
def go():
    print(1)
    raise Exception()
    print(2)

p = multiprocessing.Pool(1)

p.apply_async(go)
p.close()
p.join()

只需装饰您传递给进程池的函数。这个工作的关键是@functools.wraps(func),否则多处理会抛出PicklingError

上面的代码给出了

1
Exception in go
Traceback (most recent call last):
  File "<stdin>", line 5, in wrapped_func
  File "<stdin>", line 4, in go
Exception

【讨论】:

  • 如果并行运行的函数——在这种情况下是 go()——返回一个值,这将不起作用。装饰器不会传递返回值。除此之外,我喜欢这个解决方案。
  • 为了传递返回值,只需像这样修改 wrapper_func:` def Wrapper_func(*args, **kwargs): result = None try: result = func(*args, **kwargs) except: print ('Exception in '+func.__name__) traceback.print_exc() return result ` 像魅力一样工作 ;)
【解决方案5】:

由于multiprocessing.Pool 已经有了不错的答案,我将使用不同的方法提供一个解决方案以确保完整性。

对于python &gt;= 3.2,以下解决方案似乎是最简单的:

from concurrent.futures import ProcessPoolExecutor, wait

def go():
    print(1)
    raise Exception()
    print(2)


futures = []
with ProcessPoolExecutor() as p:
    for i in range(10):
        futures.append(p.submit(go))

results = [f.result() for f in futures]

优点:

  • 很少的代码
  • 在主进程中引发异常
  • 提供堆栈跟踪
  • 没有外部依赖

有关 API 的更多信息,请查看this

此外,如果您要提交大量任务,并且希望您的主进程在您的一项任务失败后立即失败,您可以使用以下 sn-p:

from concurrent.futures import ProcessPoolExecutor, wait, FIRST_EXCEPTION, as_completed
import time


def go():
    print(1)
    time.sleep(0.3)
    raise Exception()
    print(2)


futures = []
with ProcessPoolExecutor(1) as p:
    for i in range(10):
        futures.append(p.submit(go))

    for f in as_completed(futures):
        if f.exception() is not None:
            for f in futures:
                f.cancel()
            break

[f.result() for f in futures]

只有在所有任务都执行完毕后,所有其他答案才会失败。

【讨论】:

    【解决方案6】:

    既然你用过apply_sync,我猜这个用例是想做一些同步任务。使用回调进行处理是另一种选择。请注意此选项仅适用于python3.2及以上版本,不适用于python2.7。

    from multiprocessing import Pool
    
    def callback(result):
        print('success', result)
    
    def callback_error(result):
        print('error', result)
    
    def go():
        print(1)
        raise Exception()
        print(2)
    
    p = Pool()
    p.apply_async(go, callback=callback, error_callback=callback_error)
    
    # You can do another things
    
    p.close()
    p.join()
    

    【讨论】:

    【解决方案7】:
    import logging
    from multiprocessing import Pool
    
    def proc_wrapper(func, *args, **kwargs):
        """Print exception because multiprocessing lib doesn't return them right."""
        try:
            return func(*args, **kwargs)
        except Exception as e:
            logging.exception(e)
            raise
    
    def go(x):
        print x
        raise Exception("foobar")
    
    p = Pool()
    p.apply_async(proc_wrapper, (go, 5))
    p.join()
    p.close()
    

    【讨论】:

      【解决方案8】:

      我创建了一个模块RemoteException.py,它显示了进程中异常的完整回溯。蟒蛇2。 Download it 并将其添加到您的代码中:

      import RemoteException
      
      @RemoteException.showError
      def go():
          raise Exception('Error!')
      
      if __name__ == '__main__':
          import multiprocessing
          p = multiprocessing.Pool(processes = 1)
          r = p.apply(go) # full traceback is shown here
      

      【讨论】:

        【解决方案9】:

        我会尝试使用 pdb:

        import pdb
        import sys
        def handler(type, value, tb):
          pdb.pm()
        sys.excepthook = handler
        

        【讨论】:

        • 在那种情况下它永远不会到达处理程序,奇怪
        猜你喜欢
        • 2013-01-29
        • 2018-03-15
        • 2017-02-23
        • 2013-02-10
        • 2014-06-06
        • 2014-08-03
        • 2013-10-09
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多