Run code in multiprocessing pool processes

Viewed 178

How can I handle multiprocessing.Pool() error in initialization?

In short, I use
pool = multiprocessing.Pool(os.cpu_count(), initializer=init)

def init(): 
   # stop everything if that fails with exception
   # init..

Related question How to handle initializer error in multiprocessing.Pool?

But it does not show how to exit the main process cleanly.

Is it possible to get each Process instance fromthe pool, and run a function in each of them? Note this is not the same as pool.map, because I don't know in which process my function will run.

1 Answers

How about using a Queue to listen for initialisation failures?

For example:

from multiprocessing import Pool, Manager
from random import random

# an initialisation function that fails 10% of the time
def init(initdone):
    try:
        # fail sometimes!
        assert random() < 0.9
    except Exception as err:
        # record the failure
        initdone.put(err)
    finally:
        # record that initialisation was successful
        initdone.put(None)

# we need a manager to maintain the queue
with Manager() as manager:
    # somewhere to store initialisation state
    initresult = manager.Queue()

    # number of workers we want to wait for
    nprocs = 5

    with Pool(nprocs, initializer=init, initargs=(initresult,)) as pool:
        # wait for initializers to run
        for i in range(nprocs):
            res = initresult.get()
            # reraise if it failed (or whatever logic is appropriate)
            if res is not None:
                raise res

        # do something now we've got this pool set up
        print(sum(pool.map(int, range(20))))

I should note that multiprocessing can restart failed workers, so you might get more than nprocs entries entered into the queue. In the OP's case this is probably OK as they said they wanted to abort in this case, but this might not be true in other situations.

I can't think of any ways in which you'd get less than nprocs entries in the queue, please comment if you see any! Note that I'm deliberately not catching BaseException in init as that's handled by code in multiprocessing (along with a few others) and it seems sensible to let it do the right thing.

Related