How to get the amount of "work" left to be done by a Python multiprocessing Pool?

multiprocessing, parallel-processing, pool, process, python

Solution

Use a `Manager` queue. This is a queue that is shared between worker processes. If you use a normal queue it will get pickled and unpickled by each worker and hence copied, so that the queue can't be updated by each worker.

You then have your workers add stuff to the queue and monitor the queue's state while the workers are working. You need to do this using `map_async` as this lets you see when the entire result is ready, allowing you to break the monitoring loop.

Example:

import time
from multiprocessing import Pool, Manager


def play_function(args):
    """Mock function, that takes a single argument consisting
    of (input, queue). Alternately, you could use another function
    as a wrapper.
    """
    i, q = args
    time.sleep(0.1)  # mock work
    q.put(i)
    return i

p = Pool()
m = Manager()
q = m.Queue()

inputs = range(20)
args = [(i, q) for i in inputs]
result = p.map_async(play_function, args)

# monitor loop
while True:
    if result.ready():
        break
    else:
        size = q.qsize()
        print(size)
        time.sleep(0.1)

outputs = result.get()

Problem

So far whenever I needed to use `multiprocessing` I have done so by manually creating a "process pool" and sharing a working Queue with all subprocesses. For example: ``` from multiprocessing import Process, Queue class MyClass: def __init__(self, num_processes): self._log = logging.getLogger() self.process_list = [] self.work_queue = Queue() for i in range(num_processes): p_name = 'CPU_%02d' % (i+1) self._log.info('Initializing process %s', p_name) p = Process(target = do_stuff, args = (self.work_queue, 'arg1'), name = p_name) ``` This way I could add stuff to the queue, which would be consumed by the subprocesses. I could then monitor how far the processing was by checking the `Queue.qsize()`: ``` while True: qsize = self.work_queue.qsize() if qsize == 0: self._log.info('Processing finished') break else: self._log.info('%d simulations still need to be calculated', qsize) ``` Now I figure that `multiprocessing.Pool` could simplify a lot this code. What I couldn't find out is how can I monitor the amount of "work" still left to be done. Take the following example: ``` from multiprocessing import Pool class MyClass: def __init__(self, num_processes): self.process_pool = Pool(num_processes) # ... result_list = [] for i in range(1000): result = self.process_pool.apply_async(do_stuff, ('arg1',)) result_list.append(result) # ---> here: how do I monitor the Pool's processing progress? # ...? ``` Any ideas?

Original source