【发布时间】:2019-03-27 19:12:24
【问题描述】:
我正在使用 client.map 在 N 个单线程工作人员(在 N 台机器上)池上运行一个函数,其中一个工作人员失败了。我想知道是否有一种方法可以自动处理工作人员引发的异常,将其失败的任务重新分配给其他工作人员,并将其从池中忽略或排除?
我尝试使用下面显示的方法模拟问题。为了导致一名工人失败,我在my_function 中引发了一个 OSError,它被提交给client.map,如下所示:futures = client.map(my_function, range(100))。在我的示例中,“Computer123”上的工作人员将失败。为了处理my_function 抛出的异常,我在exception_handler 中使用了sys.exit。因此,当工作人员的任务失败时,会调用 sys.exit。结果是坏的worker的distributed.nanny捕获了故障并重新启动了worker,而客户端重新分配了它失败的任务。但是一旦坏工人再次备份,它会再次接收任务,因为它仍在池中。它再次失败并重复该过程。随着它继续失败,最终其他工人完成了所有任务。如果我可以自动处理来自“Computer123”等不良工作人员的异常并将其从池中删除,那将是理想的。也许我只需要将它从池中移除?
@exception_handler
def my_function(x):
import socket
import time
time.sleep(5)
if socket.gethostname() == 'Computer123':
raise(OSError)
else:
return x**2
def exception_handler(orig_func):
def wrapper(*args,**kwargs):
try:
return orig_func(*args,**kwargs)
except:
import sys
sys.exit(1)
return wrapper
【问题讨论】:
-
我会在 How to find why a task fails in dask distributed? 上发表评论,但我没有这样做的声誉。
标签: python distributed-computing dask dask-distributed