【问题标题】:Dask Scheduler exits with output 'Killed' on 'ddf.persist()'Dask Scheduler 退出并在 'ddf.persist()' 上输出 'Killed'
【发布时间】:2019-11-19 23:11:56
【问题描述】:

我对 DASK 很陌生,这可能真的很明显.... 我正在尝试运行一个分布式 dask 设置,其中有 1 个用于调度程序的节点和足够的工作节点来容纳内存中的数据——在这种特殊情况下,我使用了 15 个工作人员。我能够很好地启动集群,我还可以加载一些数据并对其进行分析。

我已将数据复制到工作节点,但我的客户端计算机上没有可用的数据,因此我正在延迟加载数据,如下所示:

import dask
import dask.dataframe as dd
from dask import delayed

def load_data(path):
    return dd.read_csv(path)

然后我可以做一些简单的分析就好了:

taxi_df_2016 = delayed(load_data)('/tmp/2016/*.csv').compute()
taxi_df_2016['fare_amount'].mean().compute()

...会给我一个值

但是当我想将文件保存在内存中时,调度程序会在控制台上打印Killed 时死掉......

taxi_df_2016_pers = taxi_df_2016.compute().persist()

将在调度程序节点的控制台上显示这个:

distributed.scheduler - INFO - Register tcp://10.0.0.17:40385
distributed.scheduler - INFO - Register tcp://10.0.0.6:42847
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.0.0.17:40385
distributed.core - INFO - Starting established connection
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.0.0.6:42847
distributed.core - INFO - Starting established connection
distributed.scheduler - INFO - Register tcp://10.0.0.9:44627
distributed.scheduler - INFO - Register tcp://10.0.0.7:44419
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.0.0.9:44627
distributed.core - INFO - Starting established connection
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.0.0.7:44419
distributed.core - INFO - Starting established connection
distributed.scheduler - INFO - Register tcp://10.0.0.16:41907
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.0.0.16:41907
distributed.core - INFO - Starting established connection
distributed.scheduler - INFO - Register tcp://10.0.0.18:41879
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.0.0.18:41879
distributed.core - INFO - Starting established connection
distributed.scheduler - INFO - Register tcp://10.0.0.13:32993
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.0.0.13:32993
distributed.core - INFO - Starting established connection
distributed.scheduler - INFO - Register tcp://10.0.0.8:33265
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.0.0.8:33265
distributed.core - INFO - Starting established connection
distributed.scheduler - INFO - Register tcp://10.0.0.14:33851
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.0.0.14:33851
distributed.core - INFO - Starting established connection
distributed.scheduler - INFO - Register tcp://10.0.0.10:44653
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.0.0.10:44653
distributed.core - INFO - Starting established connection
distributed.scheduler - INFO - Register tcp://10.0.0.19:40201
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.0.0.19:40201
distributed.core - INFO - Starting established connection
distributed.scheduler - INFO - Register tcp://10.0.0.5:42207
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.0.0.5:42207
distributed.core - INFO - Starting established connection
distributed.scheduler - INFO - Register tcp://10.0.0.15:36087
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.0.0.15:36087
distributed.core - INFO - Starting established connection
distributed.scheduler - INFO - Register tcp://10.0.0.12:32827
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.0.0.12:32827
distributed.core - INFO - Starting established connection
distributed.scheduler - INFO - Register tcp://10.0.0.11:35405
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.0.0.11:35405
distributed.core - INFO - Starting established connection
distributed.scheduler - INFO - Clear task state
Killed

当我查看仪表板时,我看到所有 173 个分区都已成功加载,并且内存被使用,但在此之后的某个时间,调度程序死了。

关于如何调试这个有什么想法吗?

【问题讨论】:

  • 我想我知道我做错了什么:我有一个 compute() 太多了。 taxi_df_2016_pers = taxi_df_2016.persist() 工作得很好。想知道为什么额外的compute() 会使调度程序崩溃,以及在这种情况下如何获得更多的调试/日志信息。

标签: dask dask-distributed


【解决方案1】:

看起来您正在延迟一个 dask 数据帧函数,这似乎很奇怪。

https://docs.dask.org/en/latest/delayed-best-practices.html#don-t-call-dask-delayed-on-other-dask-collections

【讨论】:

  • 我延迟加载的原因是文件不存在于客户端机器上,而只存在于集群上。因此,如果我直接使用dd.read_csv('...2016/*.csv'),我会得到OSError: ...2016/*.csv resolved to no files
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2011-06-09
  • 1970-01-01
  • 1970-01-01
  • 2016-08-23
相关资源
最近更新 更多