I have a dask array where I want to define each chunk on a different node within a cluster. Say I have something like
Ordinarily on a local machine I would have an array that looks something like
arr = np.random.rand(16, 16, 16)
Afterwards I would need to transform this array by multiplying it by a matrix M
M = np.random.rand(16**3, 16**3)
Basically I would have a function that looks like
def transform_arr(M, a):
shape = a.shape
a = (M.dot(a.ravel())).reshape(shape)
So far so good. Now let's scale this up to the sizes I actually want to work with. Lets have
import dask.array as da
da_arr = da.random.random((512, 512, 512), chunks=(256, 256, 256))
where I want the data in each chunk to be in allocated on a unique node in a cluster.
My first question
- How do I define the chunk sizes for my dask matrix I need to transform my array
da_M = da.random.random((512**3, 512**3), chunks=(?, ?))
My second question
- If I define a corresponding
def da_transform_arr(da_M, da_a)function, when I "ravel" my dask arrayda_arr, what happens to the chunks of data that are split across nodes in the cluster? Do they get reassigned to other nodes?
def da_transform_arr(da_M, da_a):
shape = da_a.shape
da_a = (da_M.dot(da_a.ravel())).reshape() # how much data is traveling across nodes?
I am just afraid that a lot of time will be spent on data traveling between nodes. I would rather keep the communication between nodes at a minimum.
I am also left with my transformation matrix M as is. I am not sure how to define a multi-dimensional version of M, or if it is worthwhile to even do so.