How to apply a map reduce approach in python

Viewed 124

I try to perform a Monte Carlo computation in parallel using python. The problem is an extremely parallel: I need to compute a function N times and add the output together, each computation is independent and the addition is a simple addition between tables.

So far I have tried two approaches:

  1. Using multiprocessing.map() and then python reduce. The problem is that I run out of memory because map is storing all the data even if I do not need to. The code looks like this:

    from multiprocessing import Pool
    import tqdm
    import numpy as np
    
    n_cpu = 8
    pool = Pool(n_cpu)
    out1 = list(tqdm.tqdm(pool.imap_unordered(f, args, chunksize = 1000)))
    out = reduce(np.add, out1)
    

    This way I obtain a poor scaling with n_cpu and the code crashes for a memory error if the input size len(args) is too large.

  2. I tried to solve using pyspark and the following code:

    import pyspark, findspark
    import numpy as np
    
    findspark.init()
    number_cores = 8
    memory_gb = 8
    conf = (
        pyspark.SparkConf()
            .setMaster('local[{}]'.format(number_cores))
            .set('spark.driver.memory', '{}g'.format(memory_gb))
    )
    sc = pyspark.SparkContext(conf=conf)
    out = sc.parallelize(range(N_samples)).repartition(number_cores).map(function).reduce(lambda a, b: np.add(a, b))
    

The repartition is done explicitly for clarity and is equal to the number of cores because I thought this is the best way to do it since the function to compute is computationally heavy. The problem is that I obtain similar performance to the multiprocessing method.

My question is:
Is there a method to make the code scale better with the number of cores? Is there a way to use multiprocessing imap_unordered() and reduce it before the computation is done?

1 Answers

Correct me if i'm wrong, but to generalize your problem:

You have 1 machine with multiple cores and you try to run some algorithm and want to get best performance.

If the statement above is true - then Spark, most probably, is not what you need. Spark is primarily for distributed computing - i.e if you have multiple machines - then you can distribute your work and perform computations faster.

If you have only one machine - you can hardly get any better performance than any multithreaded approach, since Spark does a lot of stuff which is not required for single machine computations.

Related