Merging dask dataframes results in disk error due to temporary .partd files in /tmp

Viewed 178

I have been trying to merge 32 files with a common id column in dask. In total, the files are 82.4GB big. I am iteratively going through the files and doing an outer merge on some and then a left merge on the others. I did not set a distributed client as I wasn't sure it was needed. I also tried setting the id column as index and using it to merge but it wasn't working. I also tried pandas but got a memory error even though I thought this would fit into the RAM (524GB).

The problem is that dask creates so many big temporary .partd files that it fills up the 950GB available. It didn't look like it was making full use of the RAM. I suspect the workflow creates many temporary files and they don't get cleared during the process even if they are unused. They are cleared once the scripts fails with an error due to disk space.

Any suggestion to make this better would be helpful! Thank you very much!

Code looks like this:

#!/usr/bin/env python

import pandas as pd
import dask.dataframe as dd


if __name__ == '__main__':

    list_df=pd.read_csv("path_2_files.csv")
    for index, row in list_df.iterrows():
        print("Starting "+str(index))
        if (index==0):
            ddf = dd.read_csv(row['path']+row['filename'] , sep="\t",\
                              dtype={'id': 'object', \
                                     'data1': 'float64', \
                                     'data2': 'float64', \
                                     'data3': 'float64' \
                                     })
            #Check for duplicates and drop
            ddf=ddf.drop_duplicates()
            print(len(ddf))
        else:
            ddf2 = dd.read_csv(row['path'] + row['filename'], \
                                dtype={'id': 'object', \
                                       'data1': 'float64', \
                                       'data2': 'float64', \
                                       'data3': 'float64' \
                                       })
            #Check for duplicates and drop
            ddf2=ddf2.drop_duplicates()
            
            if (row['type']=='test1'):
                ddf = dd.merge(ddf, ddf2, on='id', how="outer")
                print(len(ddf))
                del ddf2
            elif (row['type']=='test2' or row['type']=='test3'):
                ddf = dd.merge(ddf, ddf2, on='id', how="left")
                print(len(ddf))
                del ddf2

    ddf.to_parquet("s3://some/path/merged_file.parquet", object_encoding='utf8')
0 Answers
Related