【问题标题】:Computation of two Dask arrays freezing at the end最后冻结的两个 Dask 数组的计算
【发布时间】:2018-12-22 21:39:09
【问题描述】:

我生成了两个长度为 450,000,000 的随机 dask 数组,我想将它们相除。当我去计算它们时,计算总是在最后冻结。

我有一个 8 核 32GB 实例来运行代码。

我已经尝试了下面的代码,并且我尝试的一些修改没有将数据保留在 x 或 y 中。

x = da.random.random(450000000, chunks=(10000,))
x = client.persist(x)
z1 = dd.from_array(x)

y = da.random.random(450000000, chunks=(10000,))
y = client.persist(y)
z2 = dd.from_array(y)

flux_ratio_sq = z1.div(z2)
flux_ratio_sq.compute() 

我得到的实际结果是,persist 将 x 和 y 保存在内存中(总共 8GB 内存),这是预期的,然后计算会增加更多内存。我遇到的一些错误如下。

很多这样的错误:

distributed.core - INFO - Event loop was unresponsive in Scheduler for 
3.74s.  This is often caused by long-running GIL-holding functions 
or moving large chunks of data. This can cause timeouts and instability.

tornado.application - ERROR - Exception in callback <bound method 
BokehTornado._keep_alive of <bokeh.server.tornado.BokehTornado 
object at 0x7fb48562a4a8>>

raise StreamClosedError(real_error=self.error)
tornado.iostream.StreamClosedError: Stream is closed

我希望最终结果出现在 dask Series 中,以便我可以将其与现有数据合并。

【问题讨论】:

  • 我无法重现您的错误。我正在使用 16GB 内存的本地机器上尝试,所以我使用 N=4.5e7 而不是 N=4.5e8,我发现:1) 即使您不坚持,它似乎也以相同的方式执行。 2) 使用 flux_ratio_sq = da.divide(x,y) 比使用系列快 2 倍。
  • 我在加载到数据框然后进行计算时一定遇到了一些问题。使用 da.divide(x,y) 解决了这个问题。如果您想重新发布您的评论作为答案,我很乐意“接受”您的回答。谢谢!

标签: dask dask-distributed


【解决方案1】:

我会尝试在这里扩展我的评论。拳头:给定比numpy 的性能优于pandasDataFrameSeries)最好使用numpy 进行计算,然后如果需要,将结果附加到DataFrame。与Dask 完全相同。在 documentation 之后的第二个,您应该只在需要多次调用同一个数据帧的情况下坚持。

所以对于你的具体问题,你可以做的是

import dask.array as da
N = int(4.5e7)

x = da.random.random(N, chunks=(10000,))
y = da.random.random(N, chunks=(10000,))
flux_ratio_sq = da.divide(x, y).compute()

附录:使用dask.dataframe,您可以使用to_parquet() 而不是compute(),并将您的结果存储到文件中。在像这样一个令人尴尬的并行问题中,对 RAM 的影响小于使用compute()。想知道在dask.array 的情况下是否可以应用类似的东西会很有趣

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2022-10-12
    • 2020-12-04
    • 2022-08-03
    • 2022-08-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多