Dask reshape and rechunk, and matrix vector multiplication on a distributed cluster

Viewed 63

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 array da_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.

0 Answers
Related