【问题标题】:dask distributed memory errordask分布式内存错误
【发布时间】:2016-11-27 02:25:32
【问题描述】:

在分布式作业上运行 Dask 时,调度程序出现以下错误:

distributed.core - ERROR -
Traceback (most recent call last):
  File "/usr/local/lib/python3.4/dist-packages/distributed/core.py", line 269, in write
    frames = protocol.dumps(msg)
  File "/usr/local/lib/python3.4/dist-packages/distributed/protocol.py", line 81, in dumps
    frames = dumps_msgpack(small)
  File "/usr/local/lib/python3.4/dist-packages/distributed/protocol.py", line 153, in dumps_msgpack
    payload = msgpack.dumps(msg, use_bin_type=True)
  File "/usr/local/lib/python3.4/dist-packages/msgpack/__init__.py", line 47, in packb
    return Packer(**kwargs).pack(o)
  File "msgpack/_packer.pyx", line 231, in msgpack._packer.Packer.pack (msgpack/_packer.cpp:231)
  File "msgpack/_packer.pyx", line 239, in msgpack._packer.Packer.pack (msgpack/_packer.cpp:239)
MemoryError

这是调度程序或其中一个工作人员的内存不足吗?还是两者兼而有之??

【问题讨论】:

  • 鉴于错误来自调度程序,我的猜测是调度程序内存不足。这在正常操作下非常罕见。在您的问题中知道发生这种情况时您在做什么会很有趣。
  • 一种常见的情况是,您通过调用e.gather(...)collection.compute() 将大量数据从工作人员收集回客户端

标签: python dask


【解决方案1】:

此错误的最常见原因是试图收集太多数据,例如在以下使用 dask.dataframe 的示例中发生:

df = dd.read_csv('s3://bucket/lots-of-data-*.csv')
df.compute()

这会将所有数据加载到整个集群的 RAM 中(这很好),然后尝试通过调度程序将整个结果带回本地计算机(这可能无法处理您的 100 GB数据都在一个地方。)工作人员到客户端的通信通过调度程序,因此它是接收所有数据的第一台单机,并且第一台可能发生故障的机器。

如果是这种情况,那么您可能希望使用Executor.persist 方法来触发计算,但将其留在集群上。

df = dd.read_csv('s3://bucket/lots-of-data-*.csv')
df = e.persist(df)

一般来说,我们只将df.compute() 用于我们想在本地会话中查看的小结果。

【讨论】:

  • 在我的例子中,调度程序和客户端在同一台机器上运行。那么在使用persist 之后如何检索df 呢?会不会导致和以前一样的问题,因为df 需要流回调度程序/客户端?
  • 如果您在同一台机器上并且有一个适合内存的数据集,那么也许只使用 Pandas 代替?如果没有必要,我不会使用 dask.dataframe。或者考虑单机调度程序。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2017-01-18
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多