# What is the most efficient Python package to run a parallel job over multiple nodes?

**URL:** <https://ask.cyberinfrastructure.org/t/what-is-the-most-efficient-python-package-to-run-a-parallel-job-over-multiple-nodes/171>\
**Category:** Q&A\
**Tags:** parallelization, programming-for-hpc, python, researcher\
**Created:** [April 13, 2018, 2:58pm UTC](https://ask.cyberinfrastructure.org/t/what-is-the-most-efficient-python-package-to-run-a-parallel-job-over-multiple-nodes/171 "2018-04-13T14:58:46Z")\
**Posts on this page:** 4\
**Page:** 1

<div class="post-metadata">

**Author:** ![ktrn](https://ask.cyberinfrastructure.org/user_avatar/ask.cyberinfrastructure.org/ktrn/32/554_2.png) [@ktrn](https://ask.cyberinfrastructure.org/u/ktrn)\
**Post date:** [April 13, 2018, 2:58pm UTC](https://ask.cyberinfrastructure.org/t/what-is-the-most-efficient-python-package-to-run-a-parallel-job-over-multiple-nodes/171/1 "2018-04-13T14:58:46Z")

</div>

I would like to parallelize my Python script over multiple nodes on the cluster. Which Python package is most efficient for this purpose?

**CURATOR:** Katia

---

<div class="post-metadata">

**Author:** ![aculich](https://ask.cyberinfrastructure.org/user_avatar/ask.cyberinfrastructure.org/aculich/32/79_2.png) [@aculich](https://ask.cyberinfrastructure.org/u/aculich)\
**Post date:** [August 20, 2018, 9:44pm UTC](https://ask.cyberinfrastructure.org/t/what-is-the-most-efficient-python-package-to-run-a-parallel-job-over-multiple-nodes/171/4 "2018-08-20T21:44:38Z")

</div>

For the python ecosystem, consider using [Dask](http://dask.pydata.org/en/latest/) which provides advanced parallelism for analytics. [Why use Dask](http://dask.pydata.org/en/latest/why.html) versus (or along with) other options? Dask integrates with Numpy, Pandas, and Scikit-Learn, and it also:

- [scales up to clusters](http://dask.pydata.org/en/latest/why.html#scales-out-to-clusters) with multiple nodes
- [deployable on job queuing systems](https://dask-jobqueue.readthedocs.io/en/latest/) like PBS, Slurm, MOAB, SGE, and LSF
- also [scales down to parallel usage of a single-node such as a server or laptop](http://dask.pydata.org/en/latest/why.html#scales-down-to-single-computers)— modern laptops often have a multi-core CPU, 16-32GB of RAM, and flash-based hard drives that can stream through data several times faster than HDDs or SSDs of even a year or two ago.
- supports a [map-shuffle-reduce pattern](http://dask.pydata.org/en/latest/why.html#supports-complex-applications) popularized by Hadoop and is a smaller, lightweight [alternative to Spark](http://dask.pydata.org/en/latest/spark.html).
- works with [MPI via mpi4py library](http://dask.pydata.org/en/latest/setup/hpc.html#using-mpi) (see @jpessin1’s [suggestion above](https://ask.cyberinfrastructure.org/t/what-is-the-most-efficient-python-package-to-run-a-parallel-job-over-multiple-nodes/171/2)) and compatible with [infiniband or other high speed networks](http://dask.pydata.org/en/latest/setup/hpc.html#high-performance-network).

See example [Dask Jobqueue for PBS cluster](https://dask-jobqueue.readthedocs.io/en/latest/#example):

```auto
from dask_jobqueue import PBSCluster
cluster = PBSCluster()
cluster.scale(10) # Ask for ten workers

from dask.distributed import Client
client = Client(cluster) # Connect this local process to remote workers

# wait for jobs to arrive, depending on the queue, this may take some time

import dask.array as da
x = ... # Dask commands now use these distributed resources

```

---

<div class="post-metadata">

**Author:** ![jpessin1](https://ask.cyberinfrastructure.org/user_avatar/ask.cyberinfrastructure.org/jpessin1/32/339_2.png) [@jpessin1](https://ask.cyberinfrastructure.org/u/jpessin1)\
**Post date:** [May 25, 2018, 10:14pm UTC](https://ask.cyberinfrastructure.org/t/what-is-the-most-efficient-python-package-to-run-a-parallel-job-over-multiple-nodes/171/2 "2018-05-25T22:14:16Z")

</div>

If you are looking for parallel processing the traditional (and still very valid) approach is use an MPI library,  
mpi4py is an example of a python based wrapper, and [https://mpi4py.readthedocs.io/en/stable/intro.html](https://mpi4py.readthedocs.io/en/stable/intro.html)  
includes a good overview of the concepts and related methods. (Not an endorsement, just not reinventing the wheel here)

Some other things to consider:  
Would it be less work to make the job fit on a single node? With tools like concurrent.futures (or the underlying multiprocessing & threading modules) or mixed tools like numpy/scipy/pandas with Cython?

_Would a faster python implementation (like pypy) provide enough speed?_

Not that they go away when you move to multi-node, but they are often, though not always sufficient and less demanding of the user/developers time.

---

<div class="post-metadata">

**Author:** ![jpessin1](https://ask.cyberinfrastructure.org/user_avatar/ask.cyberinfrastructure.org/jpessin1/32/339_2.png) [@jpessin1](https://ask.cyberinfrastructure.org/u/jpessin1)\
**Post date:** [August 18, 2018, 8:35pm UTC](https://ask.cyberinfrastructure.org/t/what-is-the-most-efficient-python-package-to-run-a-parallel-job-over-multiple-nodes/171/3 "2018-08-18T20:35:53Z")

</div>

@ktrn As an aside if you are trying to bring more hardware resources to bear but the jobs don’t require true parallelization, message queuing is an alternate/async approach, (queue is in the standard lib), there are several message queuing (MQ) systems out there for or in python off the top of my head zMQ (null-MQ), and RabbitMQ

[Next page](https://ask.cyberinfrastructure.org/t/what-is-the-most-efficient-python-package-to-run-a-parallel-job-over-multiple-nodes/171.md?page=2)
