Trying to group out data and write them out to files

Viewed 42

I was wondering if anyone knew the proper way to write out a group of files based on the value of a column in Dask. In other words, if I want to group a bunch of columns based on a value in a column and write those out to CSVs. I've been trying to use the groupby-apply paradigm with Dask, but the problem is that it does not return a dask.dataframe object, so the function I apply it with uses the Pandas API.

Is there a better way to approach what I'm trying to do? A scalable solution would be much appreciated because some of the data that I'm dealing with is very large.

Thanks!

1 Answers

If you were saving to parquet, then partition_on kwarg would be useful. If you are saving to csv, then it's possible to do something similar with (rough pseudocode):


def save_partition(df, partition_info=None):
    for group_label, group_df in df.groupby('some_col'):
        csv_name = f"{group_label}_partition_{partition_info['number']}.csv"
        group_df.to_csv(csv_name)

delayed_save = ddf.map_partitions(save_partition)

The delayed_save can then be computed when convenient.

Related