5.5 Distributed versus GPU-accelerated frameworks

NCA-ADS · Foundations of Accelerated Data Science (12% of the exam) · Official objective: “Distributed vs. GPU-accelerated computing frameworks”

What Dask adds, how it compares with Spark, and Dask-CUDA for multi-GPU.

Key points

  1. 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.”

    — NVIDIA Glossary: Dask

    “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.”

    — NVIDIA Glossary: Dask

  2. 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.”

    — Beginner's Guide to GPU-Accelerated DataFrames for pandas Users

  3. 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.”

    — Dask cuDF documentation

    “You must also deploy a dask.distributed cluster to leverage multiple GPUs.”

    — Dask cuDF documentation

Key terms

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.”

— NVIDIA Glossary: Dask

“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.”

— NVIDIA Glossary: Dask

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 →