I have a multiprocessed code that performs the following:
- Parse a series of huge files line by line
- Add lines to a list as long as a condition is met
- When condition is
False,apply_async()a function to the list, in aPool()
def main_function(func_args):
[...]
# some heavy calculations
# queue is among the func_args
queue.put(result)
if __name__ == "__main__":
manager = mp.Manager()
lock = manager.Lock()
queue = manager.Queue()
pool = manager.Pool(processes=20, maxtasksperchild=20)
for line in open(infile, "r"):
if condition == True:
# fill list
elif condition == False:
pool.apply_async(target=main_function, args=func_args)
else:
# run last list before closing file
pool.apply_async(target=main_function, args=func_args)
pool.close()
pool.join()
Now, the code works per se, but I noticed a couple of things:
- Until all the 20 first processes are over, no new 20 are spawned. Am I doing the
apply_asynccorrectly? - The memory usage is pretty high even with a small subset of data: does each process generate a "personal" copy of the arguments that goes into the memory, unless I share it with
Manager()?
What I'm after, is:
- I would like each process to be run as soon as possible, and up to 20 at once (the 20 is passed by command line from the user via
argparse)
Any advice from more expert people? :)