I have to load data from a database into a dask dataframe. I am using dd.read_sql_table for that. I need to pass a meta in order to make some int with nan columns use the pd.Int64 datatype. As the index column is string based (e.g. "-7874571842864554321-1403311221-11"), I need to set divisions. But I am failing to do so. What I tried so far:
divisions = [
'9000000000000000000-0000000000-0',
'8000000000000000000-0000000000-0',
'7000000000000000000-0000000000-0',
'6000000000000000000-0000000000-0',
'5000000000000000000-0000000000-0',
'4000000000000000000-0000000000-0',
'3000000000000000000-0000000000-0',
'2000000000000000000-0000000000-0',
'1000000000000000000-0000000000-0',
'0000000000000000000-0000000000-0',
'-1000000000000000000-0000000000-0',
'-2000000000000000000-0000000000-0',
'-3000000000000000000-0000000000-0',
'-4000000000000000000-0000000000-0',
'-5000000000000000000-0000000000-0',
'-6000000000000000000-0000000000-0',
'-7000000000000000000-0000000000-0',
'-8000000000000000000-0000000000-0',
'-9000000000000000000-0000000000-0'
]
with ProgressBar():
df = dd.read_sql_table(tablename, DB_CONNECT_STRING, index_col='id', meta=meta,
divisions=divisions)
df.to_parquet(LOCAL_BUFFER_PATH, engine='pyarrow')
It would be fine if the data was just divided into n divisions, as the index does not say anything meaningful anyways. But what happens with this code right now is that it just seem to use one partition which overflows my memory. What I would expect is that 19 divisions get downloaded and saved separately. However, it just downloads data from my db until the memory is full without saving any.