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?