Dask - Searching for rows that match a value

Viewed 1310

Im trying to use Dask to read a folder of very large csv files (which all fit in memory, they're extremely large, but I have a lot of RAM) - my current solution looks like:

val = 'abc'

df = dd.read_csv('/home/ubuntu/files-*', parse_dates=['date'])
# 1 - df_pd = df.compute(get=dask.multiprocessing.get)
ddf_selected = df.map_partitions(lambda x: x[x['val_col'] == val])
# 2 - ddf_selected.compute(get=dask.multiprocessing.get)

Is 1 (and then using pandas) or 2 better? Just trying to get a sense of what to do?

1 Answers

You can also just do the following:

ddf_selected = ddf[ddf['val_col'] == val]

In terms of which is better it depends strongly on the operation. For large datasets that don't require in-memory shuffles dask.dataframe will likely perform better. For random access or full sorts pandas will likely perform better.

You may not want to use the multiprocessing scheduler. Generally for Pandas we recommend either the threaded or distributed schedulers.

Related