|
Article ID: 648
Last updated: 27 Nov, 2023
Dask is a Pythonic distributed programming framework that enables development of scalable Python code for big data. You can think of Dask as a Python library for parallel computing, and also as a tool to scale computational libraries such as Numpy, Pandas, and Scikit-Learn. To use Dask, load and activate one of our conda environments for machine learning that has the Dask software package installed. TIP: For multi-node applications, we recommend using the dask_mpi conda environment.
This article provides a very basic overview of how Dask parallelizes common Python operations. For detailed information, see the Dask website. Parallelizing Python OperationsThis section provides three different examples of parallelization in Dask. Example 1To create a one-dimensional array with entries equal to unity:
To sum the array:
There are three groups, each with five chunks of data (3x5) to be summed in parallel. Example 2To sum a two-dimensional array:
The 3x3 groups are summed in parallel. Example 3To add an array to its transpose:
This is done by obtaining the transposes in parallel before adding the original matrix. Parallelizing Existing SystemsDask can scale existing codebases with minor changes. For example, consider the following code:
Suppose there is no parallelization for the work at f and g, because it does not look like a big array or dataframe. You can use the dask.delayed function to parallelize the existing codebase:
Running Across Multiple NodesTo run Dask across multiple nodes, load the mpi-hpe/mpt module and use the dask_mpi conda environment:
At the command prompt on the master node, run dask-mpi:
where ~/dask_sch is the path to the Dask scheduler file, sched.json. In your Python set the Dask Client scheduler file path to the path of the Dask scheduler file (sched.json) that was created when dask-mpi was run at the command prompt:
Dask will now run across the nodes, and the scheduler information can be found in the ~/dask_sch/sched.json directory. You can use this sample PBS script:
|

