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