【问题标题】:How do I ignore a worker whose tasks have failed and redistribute its tasks to other workers?如何忽略任务失败的工作人员并将其任务重新分配给其他工作人员?
【发布时间】: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

【问题讨论】:

标签: python distributed-computing dask dask-distributed


【解决方案1】:

作为一种解决方法,您可以保留一个不良工作人员的字典,每次确定它是不良工作时(可能在它引发一定数量的异常之后)都将主机名添加到其中。

然后当你想发布一些任务时,检查它是否在违规列表中。比如:

  if socket.gethostname() in badHosts:
    skip
  else:
    do_something()

如果您能提供更多关于如何管理所连接的池的详细信息,我或许可以就如何直接删除它们提供更多建议,而不必每次都检查。

【讨论】:

  • 在运行前我不知道谁是坏工人。我正在尝试在运行时自动处理异常,包括从池中删除任何最终的不良工作人员。
  • 我想我现在明白了。你不知道主机是坏的,直到它失败了,一旦它被识别出来,你不希望它获得更多的工作,对吧?你如何设置池?
  • 没错。抱歉编辑了这么多帖子。我正在尽我所能解释这个问题。我有一组机器。一台机器托管调度程序。其他每台机器都承载一个指向调度程序的单线程工作程序。除了我所做的模拟之外,我有多台机器托管工人应该没关系。重要的是我有多个工人。
猜你喜欢
  • 2013-06-17
  • 1970-01-01
  • 1970-01-01
  • 2019-12-25
  • 2022-08-22
  • 1970-01-01
  • 2018-06-04
  • 2014-06-28
  • 2012-06-19
相关资源
最近更新 更多