Using Dask at NAS

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 Operations

This section provides three different examples of parallelization in Dask.

Example 1

To create a one-dimensional array with entries equal to unity:

import dask.array as da
x = da.ones(15, chunks=(5,)) 

To sum the array:

x.sum

There are three groups, each with five chunks of data (3x5) to be summed in parallel.

Example 2

To sum a two-dimensional array:

import dask.array as da
x = da.ones((15,15), chunks=(5,5)) 
x.sum(axis=0)

     

The 3x3 groups are summed in parallel.

Example 3

To add an array to its transpose:

import dask.array as da
x = da.ones((15,15), chunks=(5,5)) 
x + x.T

This is done by obtaining the transposes in parallel before adding the original matrix.

Parallelizing Existing Systems

Dask can scale existing codebases with minor changes. For example, consider the following code:

results = []

for x in A:
   for y in B:
       if x < y:
           results.append(f(x,y))
       else :
           results.append(g(x,y))

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:

import dask
results = []

for x in A:
   for y in B:
       if x < y:
           results.append(dask.delayed(f)(x,y))
       else :
           results.append(dask.delayed(g)(x,y))

Results = dask.compute(results)

Running Across Multiple Nodes

To run Dask across multiple nodes, load the mpi-hpe/mpt module and use the dask_mpi conda environment:

PBS r429i5n7~> module use /swbuild/analytix/tools/modulefiles
PBS r429i5n7~> module load miniconda3/v4
PBS r429i5n7~> module load mpi-hpe/mpt
PBS r429i5n7~> source activate dask_mpi

At the command prompt on the master node, run dask-mpi:

PBS r429i5n7~> mpirun -np 20 dask-mpi --worker-class distributed.Worker --scheduler-file 
~/dask_sch/sched.json & your_python_script

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:

import dask
from dask.distributed import Client

#initialize dask client
path_to_sched = '/dask_sch/sched.json'
client = Client(scheduler_file=path_to_sched)

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:

##!/bin/bash
#PBS -q devel -l select=25:ncpus=4:model=has
#PBS -l walltime=01:00:00
#PBS -j oe
#PBS -N air_dask_test

#set environment
module use /swbuild/analytix/tools/modulefiles
module load miniconda3/v4
module load mpi-hpe/mpt
source activate dask_mpi

#clear and create directory for dask scheduler
if [ -d dask_sch ]
then 
    echo "removing old scheduler directory"
    rm -r dask_sch
    echo "creating new scheduler directory"
    mkdir dask_sch
else
    echo "creating scheduler directory"
    mkdir dask_sch
fi
 
#run program
mpirun -np 20 dask-mpi --worker-class distributed.Worker --scheduler-file 
$(pwd)/dask_sch/sched.json & your_python_script

Contact Us

General User Assistance

Security

  • Report security issues 24x7x365
  • Toll-free: 1-877-NASA-SEC (1-877-627-2732)
  • E-mail: soc@nasa.gov

User Documentation

High-End Computing Capability (HECC) Portfolio Office

NASA High-End Computing Program

Tell Us About It

Our goal is furnish all the information you need to efficiently and effectively use the HECC resources needed for your NASA computational projects.

We welcome your input on features and topics that you would like to see included on this website.

Please send us email with your wish list and other feedback.