I am using Dask to read a dataset of 2.8GB. This dataset is split into 2 csv files of 1.4GB each. I want to simply read this dataset and save it into parquet format, the code is the following and it is pretty simple:
import dask
import dask.dataframe as dd
import RootPath
import pandas as pd
import numpy as np
import gc
### Create intermediate parquet full dataset
# Read data
df = dd.read_csv(original_dataset_path,
sep='\x01',
names=all_features_dtype.keys(),
dtype=all_features_dtype,
)
# Write to parquet
df.to_parquet(dataset_path, write_index=False, compression="snappy", engine="pyarrow", overwrite="True")
I tried many partition size/number but I keep getting an OOM error during the computation of this simple code. I am using a 8GB RAM Laptop. While the code is running, it seems that it keeps reading from disk into RAM memory, never writing anything to disk until I get a OOM error.
How could I debug my code and find which is the problem? I tried to change the partition size ecc... as suggested in other questions, but it doesn't matter at all, I keep getting OOM.
IF further informations are needed, I will provide them as requested.
NB: I am aware that in this case a simple Pandas approach is enough, however I am just messing around with Dask in order to scale it to bigger dataset (the final dataset is about 350GB). I am doing this to test my code before running it on an AWS instance to save some money.
EDIT: I moved from the single scheduler to the distributed one. I basically modified the code by just adding c = Client()
import dask
import dask.dataframe as dd
import RootPath
import pandas as pd
import numpy as np
import gc
from dask.distributed import Client,wait,LocalCluster
### Create intermediate parquet full dataset
# Read data
df = dd.read_csv(original_dataset_path,
sep='\x01',
names=all_features_dtype.keys(),
dtype=all_features_dtype,
)
c = Client()
# Write to parquet
df.to_parquet(dataset_path, write_index=False, compression="snappy", engine="pyarrow", overwrite="True")
And now it works fine: it reads the data into the RAM memory and writes back to disk in a parallel manner without getting out of memory. In seems to read each chunk, load it into memory, write it back to disk in parquet format, release the memory.
Compared to the single scheduler version, it seems that the single one waits to load the whole dataset to memory before starting to write to disk. Why is this happening?