5.5 Distributed versus GPU-accelerated frameworks
What Dask adds, how it compares with Spark, and Dask-CUDA for multi-GPU.
Key points
A distributed framework spreads work over many machines or processes. Dask scales NumPy, pandas and scikit-learn style code.
What NVIDIA says (2)
“Dask is an open-source library designed to provide parallelism to the existing Python stack.”
“It provides integrations with Python libraries like NumPy Arrays, Pandas DataFrames, and scikit-learn to enable parallel execution across multiple cores, processors, and computers without having to learn new libraries or languages.”
Spark has its own API and engine. RAPIDS keeps pandas-like and scikit-learn-like APIs and uses Dask to scale out. API means application programming interface.
What NVIDIA says (1)
“However, unlike Apache Spark, it does not introduce a new API but provides a familiar programming interface of tools found in the PyData ecosystem, like pandas, scikit-learn, NetworkX, and more.”
Dask cuDF is the GPU backend for Dask DataFrames. Running across GPUs or nodes needs a cluster of workers, which Dask-CUDA sets up.
What NVIDIA says (2)
“Note Neither Dask cuDF nor Dask DataFrame provide support for multi-GPU or multi-node execution on their own.”
“You must also deploy a dask.distributed cluster to leverage multiple GPUs.”
Key terms
- Parallel processing: Splitting a job into many pieces that run at the same time.
- Dask: A Python library that runs familiar tools such as pandas in parallel across cores, GPUs and machines.
- Dask cuDF: Dask DataFrames backed by cuDF, for data split into partitions across GPUs.
Sample question
What is the main idea behind Dask?
Show the answer
Answer: Parallelism for the existing Python stack across cores, processors and computers, without learning new libraries
A distributed framework spreads work over many machines or processes. Dask scales NumPy, pandas and scikit-learn style code.
What NVIDIA says (2)
“Dask is an open-source library designed to provide parallelism to the existing Python stack.”
“It provides integrations with Python libraries like NumPy Arrays, Pandas DataFrames, and scikit-learn to enable parallel execution across multiple cores, processors, and computers without having to learn new libraries or languages.”
Practice 5.5 (3 questions) Full Foundations of Accelerated Data Science guide
← 5.4 The end-to-end data science workflow · 5.6 Parameters, tuning and fitting →