I have a dask process that runs a function on each dataframe partition. I let to_parquet do the
compute() that runs the functions.
But I also need to know the number of records in the parquet table. For that, I use ddf.map_partitions(len). Problem is that when I count the number of records, a compute() is done again on the dataframe, and that makes the map_partitions functions run again.
What should be the approach to run map_partitions, save the result in parquet, and count the number of records?
def some_func(df):
df['abc'] = df['def'] * 10
return df
client = Client('127.0.0.1:8786')
ddf.map_partitions(some_func) # some_func executes twice for each partition
ddf.to_parquet('/some/folder/data', engine='pyarrow')
total = ddf.map_partitions(len).compute().sum()