【问题标题】:Kubernetes and Dask and SchedulerKubernetes 和 Dask 和调度程序
【发布时间】:2020-09-05 19:11:25
【问题描述】:

我的代码看起来像这样

def myfunc(param):
    # expensive stuff that takes 2-3h

mylist = [...]
client = Client(...)
mgr = DeploymentMgr()
# ... setup stateful set ...
futures = client.map(myfunc, mylist, ..., resources={mgr.hash.upper(): 1})
client.gather(futures)

我在 Kubernetes 集群上运行 dask。在程序开始时,我创建了一个有状态集。这是通过kubernetes.client.AppsV1Api() 完成的。然后我最多等待 30 分钟,直到我请求的所有工作人员都可用。对于此示例,假设我请求 10 个工作人员,但 30 分钟后,只有 7 个工作人员可用。最后,我调用client.map() 并将一个函数和一个列表传递给它。这个列表有 10 个元素。但是,dask 只会使用 7 个 worker 来处理这个列表!即使几分钟后剩余的 3 个工作人员可用,即使第一个元素的处理都没有完成,dask 也不会为它们分配任何列表元素。

如何改变 dask 的行为?有没有办法告诉 dask(或 dask 的调度程序)定期检查新到达的工作人员并更“正确地”分配工作?或者我可以手动影响这些列表元素的分布吗?

谢谢。

【问题讨论】:

    标签: python kubernetes dask


    【解决方案1】:

    一旦 Dask 更好地了解任务需要多长时间,它就会平衡负载。您可以使用配置值估算任务长度

    distributed:
      scheduler:
        default-task-durations:
          myfunc: 1hr
    

    或者,一旦 Dask 完成其中一项任务,它将知道如何在未来围绕该任务做出决策。

    我相信这在 GitHub 问题跟踪器上也出现过几次。您可能需要搜索https://github.com/dask/distributed/issues 以获取更多信息。

    【讨论】:

    • 感谢您的帮助。我尝试设置默认任务持续时间,但我还没有测试过。在我这样做之前,我用几个 client.submit() 调用替换了单个 client.map() 调用。然后我让调度程序给我所有“正在运行”的工作人员的所有 IP 地址,然后在 client.submit() 的 workers 参数中使用这些 IP 地址。然而,这并没有解决问题。所以我希望 default-task-duration 能解决它。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-10-27
    • 2019-04-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多