0
votes

I have a 2D grid on which there is a path. I want to calculate the distances of each point of the grid to each point on the path, then do some operations on those grid. I am using dask.dataframe and dask.array for this task.

The code is:

import dask.dataframe as dd
import dask.array as da 

x = np.linspace(-60, 60, 10000)
xv, yv = da.meshgrid(x, x, sparse='True')

path = da.from_array(np.random.rand(100, 2))

h = 100.0

# function to calculate distance to point
def dist_to_point(x, y, p):
    x_dist = x-p[0] 
    y_dist = y-p[1]
    dist = da.sqrt(x_dist**2+y_dist**2)
    d2 = da.sqrt(dist**2 + h**2)    
    return dd.from_dask_array(d2) 


distances = [dist_to_point(xv, yv, path[i, :]) for i in range(npath)]
distances_grid = dd.multi.concat(distances, axis=1, ignore_index=True)

So distances_grid should the concatenation of [grid distance to point 1, grid distance to point 2, ..., grid distance to point 100]

Now suppose I want to get the max across all dataframes I apply this

l_max = distances_grid.map_partitions(lambda x: x.groupby(level=0, axis=1).max())

The dask graph for this looks like this which to me does not look like proper parallelization of the tasks. Can anyone help point me to what I am doing wrong or how I can improve this? My final application will be on 100000x100000 grids hence the use of dask enter image description here

1

1 Answers

0
votes

So in case anyone runs into this I solved it by broadcasting the arrays and avoiding the for loop all together. The code I ended up using is

x = da.from_array(np.linspace(-60, 60, 10000), chunks=1000)
xv, yv = da.meshgrid(x, x, sparse='True')
path = da.from_array(np.random.rand(10, 2))

h = 100.0

ngrid = x.shape[0]

xd = x[:, np.newaxis] - path[:, 0]
yd = x[:, np.newaxis] - path[:, 1]
z = xd**2 + yd[:, np.newaxis]**2 + h**2

# euclidian distance at height = 100
z = xd**2 + yd[:, np.newaxis]**2 + h**2
distances_grid = z**0.5

l_max = distances_grid.max(axis=2)

This gave me a nicer graph which I am able to balance even more by changing the sizes of the chunks. enter image description here