dtr*_*iak 5 python numpy dot-product dask
我正在努力比较不同数据大小的Dask和Numpy的计算速度.我知道Dask可以并行执行数据计算,并将数据拆分成块,这样数据大小就可以大于RAM.当使用下面的Dask代码时,我得到一个内存错误(显示在底部),方形数组大小为42000.
import dask as da
import time
size = 42000
y = da.random.random(size = (size,size), chunks = (size/8,size/8))
start = time.time()
y = y.dot(y*2) #arbitrary dot product calculation
y.compute()
end = time.time()
print(str(end-start) + " seconds")
Run Code Online (Sandbox Code Playgroud)
但是,在使用Numpy运行类似代码时,我不会收到任何错误.
import numpy as np
import time
size = 42000
x = np.random.random(size = (size,size))
start = time.time()
x = x.dot(x*2) #arbitrary dot product calculation
end = time.time()
print(str(end-start) + " seconds")
Run Code Online (Sandbox Code Playgroud)
因此,我不明白为什么当Numpy不特别因为Dask应该能够对数据进行分区时Dask会抛出内存错误.这有什么解释/解决方案吗?
编辑:我只有点产品的这个问题.我已经用平均值测试没有任何问题.
MemoryError Traceback (most recent call last)
<ipython-input-3-a3af599b673a> in <module>()
3 start = time.time()
4 y = y.dot(y*2)
----> 5 y.compute()
6 end = time.time()
7 print(str(end-start) + " seconds")
~\AppData\Local\Continuum\anaconda3\lib\site-packages\dask\base.py in compute(self, **kwargs)
152 dask.base.compute
153 """
--> 154 (result,) = compute(self, traverse=False, **kwargs)
155 return result
156
~\AppData\Local\Continuum\anaconda3\lib\site-packages\dask\base.py in compute(*args, **kwargs)
405 keys = [x.__dask_keys__() for x in collections]
406 postcomputes = [x.__dask_postcompute__() for x in collections]
--> 407 results = get(dsk, keys, **kwargs)
408 return repack([f(r, *a) for r, (f, a) in zip(results, postcomputes)])
409
~\AppData\Local\Continuum\anaconda3\lib\site-packages\dask\threaded.py in get(dsk, result, cache, num_workers, **kwargs)
73 results = get_async(pool.apply_async, len(pool._pool), dsk, result,
74 cache=cache, get_id=_thread_get_id,
---> 75 pack_exception=pack_exception, **kwargs)
76
77 # Cleanup pools associated to dead threads
~\AppData\Local\Continuum\anaconda3\lib\site-packages\dask\local.py in get_async(apply_async, num_workers, dsk, result, cache, get_id, rerun_exceptions_locally, pack_exception, raise_exception, callbacks, dumps, loads, **kwargs)
519 _execute_task(task, data) # Re-execute locally
520 else:
--> 521 raise_exception(exc, tb)
522 res, worker_id = loads(res_info)
523 state['cache'][key] = res
~\AppData\Local\Continuum\anaconda3\lib\site-packages\dask\compatibility.py in reraise(exc, tb)
65 if exc.__traceback__ is not tb:
66 raise exc.with_traceback(tb)
---> 67 raise exc
68
69 else:
~\AppData\Local\Continuum\anaconda3\lib\site-packages\dask\local.py in execute_task(key, task_info, dumps, loads, get_id, pack_exception)
288 try:
289 task, data = loads(task_info)
--> 290 result = _execute_task(task, data)
291 id = get_id()
292 result = dumps((result, id))
~\AppData\Local\Continuum\anaconda3\lib\site-packages\dask\local.py in _execute_task(arg, cache, dsk)
269 func, args = arg[0], arg[1:]
270 args2 = [_execute_task(a, cache) for a in args]
--> 271 return func(*args2)
272 elif not ishashable(arg):
273 return arg
~\AppData\Local\Continuum\anaconda3\lib\site-packages\dask\compatibility.py in apply(func, args, kwargs)
46 def apply(func, args, kwargs=None):
47 if kwargs:
---> 48 return func(*args, **kwargs)
49 else:
50 return func(*args)
~\AppData\Local\Continuum\anaconda3\lib\site-packages\numpy\core\fromnumeric.py in sum(a, axis, dtype, out, keepdims)
1880 return sum(axis=axis, dtype=dtype, out=out, **kwargs)
1881 return _methods._sum(a, axis=axis, dtype=dtype,
-> 1882 out=out, **kwargs)
1883
1884
~\AppData\Local\Continuum\anaconda3\lib\site-packages\numpy\core\_methods.py in _sum(a, axis, dtype, out, keepdims)
30
31 def _sum(a, axis=None, dtype=None, out=None, keepdims=False):
---> 32 return umr_sum(a, axis, dtype, out, keepdims)
33
34 def _prod(a, axis=None, dtype=None, out=None, keepdims=False):
MemoryError:
Run Code Online (Sandbox Code Playgroud)
在最后阶段,当 Dask 将事物拼接在一起时,它可能需要大约 2 倍的内存用于输出。
一般来说,如果您的计算适合内存,您可能不应该使用 Dask。具有现代 BLAS 实现(OpenBLAS、MKL...)的 NumPy 可能会表现得更好。