0
votes

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.

Please show us what error you are getting - mdurant
@mdurant I am not getting any error. The script just terminates as soon as the memory is full. - McToel
You mean, df.npartitions == 1? Note that you probably need the divisions to have the "natural" ordering of the DB, probably the same as sorted(divisions). - mdurant