multiprocessing: Mapping a function using a list of Processes instead of a Pool

Viewed 93

I have a multiprocessing task that, in its simplest form, looks like the following:

def fun(x):
     y = setup()
     return y.f(x) 

pool = mp.Pool(4)
pool.map(fun, my_list)

However, the setup() is expensive and so I only want to do it once in each process, as opposed to doing it once per item in my_list.

I also do not want to pickle y and send it into each process, in this instance I require that the setup occurs within each process separately.

Hence, I could do something like this to set up each process:

class MyProcess(mp.Process):
    def __init__(self):
        self.y = setup()

    def fun(x):
        return self.y.f(x)

workers = [MyProcess() for _ in range(4)]

Is there any way I can now use workers as if it was a Pool? i.e. mapping some worker's worker.fun to each item in my_list? Ideally, I would want something like this:

for result in workers.imap_unordered(MyProcess.fun, my_list):
    # do something

I suspect a solution using a Queue would also work, but I'm not entirely sure how I can implement this.

3 Answers

A Pool already supports customising each process at startup. Define an initialiser of the Pool processes that creates y and makes it accessible:

def init_process():
    global y     # make y accessible to everything
    y = setup()  # ... and initialise it

def fun(x):
     # use already initialised y
     return y.f(x) 

pool = mp.Pool(4, initializer=init_process)
pool.map(fun, my_list)

To initialize processes in the pool when they are created you can use initializer and initargs parameters of Pool.

An example to illustrate the approach:

import multiprocessing as mp

init_obj = {}


def setup(a):
    global init_obj
    init_obj = {"one": a}


def fun(x):
    y = init_obj
    print(y)


pool = mp.Pool(None, initializer=setup, initargs=(1,))
pool.map(fun, [0, 1, 2])

I feel that you could still use a pool:

def fun(x)
    y = setup()
    return [y.f(item) for item in x]

processes = 4
newlist = []
for i in range(processes):
    mylist[i * len(mylist)//4: (i + 1) * len(mylist)//4]
    newlist.append(mylist)

pool = mp.Pool(4)
pool.map(fun, new_list)

There's an overhead for splitting the list, but this is the simplest solution I can think of to reduce the number of times setup is called.

NOTE: x in this version of the code is a list of values, instead of a single value.

Related