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 ! :)