How do you use dask to groupby column and apply with custom function without running out of memory?

Viewed 386

I know I can do this easily on spark but have been trying with dask and keep getting out of memory errors, maybe I'm not doing it correctly.

Here is the situation: I have a very large dataframe, let's call it df. df has 2 columns: key and amount

   key  amount 
    1 4.2
    1 4.3
    1 4.2
    2 4.1 
    2 4.1
    2 4.3

I have a custom function:

def process(df):
    // a lot of processing 
    // cant be written in dask
    return df 

As you can see process operates on the entire dataframe and returns a new dataframe that has been processed. The stuff inside process cannot be translated into dask.

I need to partition by key (so I get 1 dataframe per key), run process on each key's dataframe and then combine them back into a single data frame and finally write it as csv.

I have tried df.groupby("key").apply(process).to_csv(...) but run out of memory. I have also tried df.map_partitions(...).compute() but also run out of memory

I read the docs and even tried df.map_partitions(lambda x: x.groupby(..).apply(proc)) but realized this does not work as the groups are not isolated by partition.

I know this is an easy thing to do what am I doing wrong?

0 Answers
Related