【发布时间】: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