Parallel processing of each row in Pandas iteration

Viewed 1572

I have df_fruits, which is a dataframe of fruits.

index      name
1          apple
2          banana
3          strawberry

and, its market prices are in mysql database like below,

category      market      price
apple         A           1.0
apple         B           1.5
banana        A           1.2
banana        A           3.0
apple         C           1.8
strawberry    B           2.7        
...

During the iteration in df_fruits, I'd like to do some processes.

The code below is a non-parallel version.

def process(fruit):
   # make DB connection
   # fetch the prices of fruit from database
   # do some processing with fetched data, which takes a long time
   # insert the result into DB
   # close DB connection

for idx, f in df_fruits.iterrows():
    process(f)

What I want to do is to do process on each row in df_fruits in parallel, since df_fruits has plenty of rows and the table size of the market prices is quite large (fetching data takes a long time).

As you can see, the order of execution between rows does not matter and there's no sharing data.

Within iteration in df_fruits, I'm confused about where to locate `pool.map(). Do I need to split the rows before parallel execution and distribute chunks to each process? (If so, a process which finished its job earlier than other process would be idle?)

I've researched of pandarallel but I can't use it (my os is windows).

Any help would be appreciated.

3 Answers

There is no need to use pandas at all. You can simply use Pool from the multiprocessing package. Pool.map() takes two inputs: one function and one list of values.

So you can do:

from multiprocessing import Pool

n = 5  # Any number of threads
with Pool(n) as p:
    p.map(process, df_fruits['name'].values)

This will go through all fruits in the df_fruits dataframe one-by-one. Note that there is no result returned here since the process function is designed to write the result back to the database.


If you have multiple columns that you want to consider in each row, you can change df_fruits['name'].values to:

df_fruits[cols].to_dict('records')

this will give a dictionary as the input to preprocess, e.g.:

{'name': 'apple', 'index': 1, ...}

Yeah its possible, although not really provided in the pandas library straight out of the box.

Maybe you can attempt something like this:

def do_parallel_stuff_on_dataframe(df, fn_to_execute, num_cores):
    # create a pool for multiprocessing
    pool = Pool(num_cores)

    # split your dataframe to execute on these pools
    splitted_df = np.array_split(df, num_cores)

    # execute in parallel:
    split_df_results = pool.map(fn_to_execute, splitted_df)

    #combine your results
    df = pd.concat(split_df_results)

    pool.close()
    pool.join()
    return df

You might to be able do something like:

with Pool() as pool:
    # create an iterator that just gives you the fruit and not the idex
    rows = (f for _, f in df_fruits.iterrows())
    pool.imap(process, rows)

You may want to use one of the other pool primitives other than map if you don't care for the results, or are willing to get the results in any order or don't care about the results.

Related