【问题标题】:are there any limits on number of the dask workers/cores/threads?对 dask 工作人员/核心/线程的数量有任何限制吗?
【发布时间】:2021-08-13 16:34:59
【问题描述】:

当我使用超过 25 个工作人员时,我发现我的数据分析性能有所下降,每个工作人员有 192 个线程。调度程序有任何限制吗?通信(使用ib)或cpu或ram)没有负载足迹。 例如,最初我在 lustrefs 上有 170K hdf 文件:

ddf=dd.read_hdf(hdf5files,key="G18",mode="r")
ddf.repartition(npartitions=4096).to_parquet(splitspath+"gdr3-input-cache")

代码在 64 名工作人员上的运行速度比 25 名工作人员慢。看起来初始任务设计阶段的调度程序非常过载。

编辑: dask-2021.06.0 分布式-2021.06.0

【问题讨论】:

  • 这是最新的dask 版本吗? (最近升级了处理大型任务图的方式)
  • 不是最近,但是 dask 2021.06.0(我想这个已经有那个更新了)你有 git-commit id 吗?
  • 不确定 id,但在此处查看更改日志docs.dask.org/en/latest/changelog.html 自 2021.6.0 以来有很多更改...它可能无法解决问题,但可能值得检查一下最新版本...
  • 感谢提示,等待下一个版本,他们修复了一个非常重要的 showstopper 错误。

标签: dask dask-distributed


【解决方案1】:

有许多潜在的瓶颈。这里有一些提示。

是的,调度程序是所有任务都必须通过的单个进程,并且它引入了每个任务的开销(

同样,如果您有很多工作人员,那么您将有大量的网络流量用于任务分配和工作人员之间的任何数据混洗。更多的工人,更多的流量。

第三,python 在运行代码时使用全局锁 GIL。即使您的任务是 GIL 友好的(例如,数组/数据帧操作),线程有时仍可能需要 GIL,这可能会导致争用和性能下降。

最后,您说您正在使用 lustre,因此您有许多任务同时处理网络存储,这对于元数据访问和数据流量都有其自身的限制。

【讨论】:

  • 做几个测试看起来像一个主要瓶颈是单线程调度程序(如果它是真的)。如果我将分区从 10k 减少到 512,则速度会更快,但还有另一个缺点:稍后内存压力会更高。 fs 不是问题,3k 核的持久性需要几秒钟。
  • 我们正在积极努力减少由于调度程序导致的每个任务的开销。
  • 好消息,谢谢,我希望在某个时候我们能够在 10K 节点上使用我们的工作负载 :)
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2021-04-25
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-10-27
相关资源
最近更新 更多