Sharing large read-only array and using the apply_async method of Python's multiprocessing.pool module

Viewed 196

I developed a little program that computes several things on a large number of independent combination of parameters. To that extent, I've parallelized my code with the pool.apply_async() method of Python's multiprocessing module. I'm using apply_async because the parallelized function takes several arguments.

At some point in the parallelized function, I need to read and retrieve the values of certain entries of a large hypercube. I've checked the memory activity on my computer and it seems like the hypercube - passed as an argument in the function to parallelized - is copied on ALL the child running processes. I'd like to know if it is possible to share my hypercube to the function parallelized with apply_async without it being copied on every child process. I only need to read some values of the hypercube at each iteration but there is a lot of iteration so I need the memory access to be quick.

I've been searching answers for my problem everywhere, a lot of people seem to have had similar issues, most of them use multiprocessing.Array but I just can't make it work with the apply_async method, the hypercube is still copied everywhere...

I'm running my code on macOS Catalina so I guess I am not subject to the problems caused with the fork() on Windows.

Here's a little example that mimics what I'm trying to do :

import numpy as np
import ctypes
import time
import multiprocessing as mp

N = 10000 # size of array

M = np.arange(0, N*N) # large array (equivalent of my hypercube)

# Sharing memory and storing M
M_shared_base = mp.Array(ctypes.c_double, M.size)
M_shared = np.frombuffer(M_shared_base.get_obj())
M_shared[:] = M[:]
del M

params = np.array([-2,0.5,8,4,45]) # a bunch of parameters

def func(idx, parameters, def_param = M_shared):
    ''' Function to parallelize using the shared memory for M '''
    # Computes independant stuff on parameters for every iteration
    # ...
    # ...
    time.sleep(10)

def func2(idx, parameters, matrix):
    ''' Function to parallelize NOT using the shared memory for M '''
    # Computes independant stuff on parameters for every iteration
    # ...
    # ...
    time.sleep(10)

if __name__ == '__main__':
    
    print("\nStarting parallelization on 'func'...")
    pool = mp.Pool(4)
    output = [pool.apply_async(func, args=(i, params)) for i in range(4)]
    results = [r.get() for r in output]
    pool.close()
    pool.join()
    
    print("\nStarting parallelization on 'func2'...")
    pool = mp.Pool(4)
    output = [pool.apply_async(func2, args=(i, params, M_shared)) for i in range(4)]
    results = [r.get() for r in output]
    pool.close()
    pool.join()

When I execute this code and look at the memory, both func and func2 are running 4 processes, each using around 780 MB. What I want is the main process to use around 780 MB and the 4 child processes only a few MB as the matrix is shared by the parent process.

Thank you very much to those who'll answer, I hope I'm not bothering with something that already has been answered somewhere else.

A very pleasant evening to all ! :)

0 Answers
Related