I have to do an operation in which a function receives a big pandas dataframe, performs some calculations and returns a result (each iteration takes ~1 min). I have to repeat this operation several times by "passing" this pandas dataframe to the function (~5000 times).
I have decided to use the Python library multiprocessing to perform all these operations, as they are independent between them. However, when using the Pool method (with imap), all the values that are passed as input of the function that is going to be parallelized are inferred at once (first get a list with all the inputs and the parallelize). In my case, this results in a very big list of pandas dataframes that gives memory issues (it is a list with the same pandas dataframe repeated several times).
Is there a way to create an iterator and generate new inputs (the pandas dataframe) as the processes that are currently running finish?
Next it is shown a small dummy example in which all the inputs are calculated at once and then the function do_something() is parallelized (this is not the behaviour that I want):
import multiprocessing
import time
def do_something(a):
print(f"Sleep {a}")
time.sleep(1)
return a**2, a + a
def test_multiprocessing(inputs):
class data_generator:
def __init__(self, inputs):
self.inputs = inputs
self.i = 0
def __next__(self):
if self.i >= len(self):
raise StopIteration
return_value = self.inputs[self.i]
self.i += 1
print(f"Input id: {self.i}")
return return_value
def __iter__(self):
return self
def __len__(self):
return len(self.inputs)
# Parallelize
with multiprocessing.Pool(processes=4) as executor:
results = executor.imap(do_something, data_generator(inputs))
executor.close()
executor.join()
for result in results:
print(result)
print("Done")
if __name__ == "__main__":
inputs = [(i) for i in range(1, 30)]
test_multiprocessing(inputs)
Is there any other solution to this problem?