Project reference ↗

Dask Distributed lets Python programs send pieces of a computation to workers on multiple machines. A scheduler coordinates that work and its dependencies, while the client submits tasks and receives results. It is useful when a computation needs more parallel processing or memory than one process can provide. This entry focuses on the distributed scheduler-and-worker setup; deploying more workers does not remove the need for compatible Python environments and accessible input data.

Deployment and operating notes

Dask Distributed is the scheduler and worker system used for distributed Python computation. Dask’s current Kubernetes documentation describes several deployment options, including Dask Kubernetes and Helm-based clusters. The historical chart’s public load balancers and bundled notebook are not required architecture and should not be carried forward without an access-control decision.

For an existing cluster, pin compatible Python, Dask, Distributed and application-library environments across client, scheduler and workers. Inventory data-access credentials, worker memory limits, spill storage and dashboard exposure. Choose whether a shared long-running cluster or per-job clusters better match isolation and cost requirements. Test representative computations and failure recovery, including a worker disappearing and data being recomputed or reloaded. A notebook session or scheduler can hold transient state that is not recovered by restoring a Deployment manifest. Preserve input data and reproducible environment definitions, and keep the old endpoint available while clients are moved deliberately. Performance claims below describe historical examples, not measured results for current workloads.

Historical upstream link check · 2026-10-09

The recorded upstream address responded successfully (HTTP 200) on 2026-10-09. GitHub confirms that helm/charts is archived: this is a historical chart distribution, not evidence that the application itself is retired. Link availability does not certify the historical installation instructions or current security support.

Source for this check ↗

Website availability is separate from project, chart and image support. Use the current guidance and primary sources on this page to assess the distribution.

The original record

Historical Kubedex content

Preserved for context. Commands, versions, prices and results below reflect the original research.

Dask.distributed. Dask.distributed is a lightweight library for distributed computing in Python. It extends both the concurrent.futures and dask APIs to moderate sized clusters.

 

Chart Details

This chart will do the following:

  • 1 x Dask scheduler with port 8786 (scheduler) and 80 (Web UI) exposed on an external LoadBalancer
  • 3 x Dask workers that connect to the scheduler
  • 1 x Jupyter notebook with port 80 exposed on an external LoadBalancer
  • All using Kubernetes Deployments

Motivation

 

Distributed serves to complement the existing PyData analysis stack. In particular it meets the following needs:

  • Low latency: Each task suffers about 1ms of overhead. A small computation and network roundtrip can complete in less than 10ms.
  • Peer-to-peer data sharing: Workers communicate with each other to share data. This removes central bottlenecks for data transfer.
  • Complex Scheduling: Supports complex workflows (not just map/filter/reduce) which are necessary for sophisticated algorithms used in nd-arrays, machine learning, image processing, and statistics.
  • Pure Python: Built in Python using well-known technologies. This eases installation, improves efficiency (for Python users), and simplifies debugging.
  • Data Locality: Scheduling algorithms cleverly execute computations where data lives. This minimizes network traffic and improves efficiency.
  • Familiar APIs: Compatible with the concurrent.futures API in the Python standard library. Compatible with dask API for parallel algorithms
    Easy Setup: As a Pure Python package distributed is pip installable and easy to set up on your own cluster.

 

Architecture

Dask.distributed is a centrally managed, distributed, dynamic task scheduler. The central dask-scheduler process coordinates the actions of several dask-worker processes spread across multiple machines and the concurrent requests of several clients.

The scheduler is asynchronous and event driven, simultaneously responding to requests for computation from multiple clients and tracking the progress of multiple workers. The event-driven and asynchronous nature makes it flexible to concurrently handle a variety of workloads coming from multiple users at the same time while also handling a fluid worker population with failures and additions. Workers communicate amongst each other for bulk data transfer over TCP.

Internally the scheduler tracks all work as a constantly changing directed acyclic graph of tasks. A task is a Python function operating on Python objects, which can be the results of other tasks. This graph of tasks grows as users submit more computations, fills out as workers complete tasks, and shrinks as users leave or become disinterested in previous results.

Users interact by connecting a local Python session to the scheduler and submitting work, either by individual calls to the simple interface client.submit(function, *args, **kwargs) or by using the large data collections and parallel algorithms of the parent dask library. The collections in the dask library like dask.array and dask.dataframe provide easy access to sophisticated algorithms and familiar APIs like NumPy and Pandas, while the simple client.submit interface provides users with custom control when they want to break out of canned “big data” abstractions and submit fully custom workloads.

Futures

Dask supports a real-time task framework that extends Python’s concurrent.futures interface. This interface is good for arbitrary task scheduling, like dask.delayed, but is immediate rather than lazy, which provides some more flexibility in situations where the computations may evolve over time.

These features depend on the second generation task scheduler found in dask.distributed (which, despite its name, runs very well on a single machine).

 

Single Machine: Dask.distributed

The dask.distributed scheduler works well on a single machine. It is sometimes preferred over the default scheduler for the following reasons:

It provides access to asynchronous API, notably Futures
It provides a diagnostic dashboard that can provide valuable insight on performance and progress
It handles data locality with more sophistication, and so can be more efficient than the multiprocessing scheduler on workloads that require multiple processes.
You can create a dask.distributed scheduler by importing and creating a Client with no arguments. This overrides whatever default was previously set.

Sources & further reading

  1. Dask Kubernetes deployment options
  2. Dask Distributed documentation
  3. Recovered historical source (Common Crawl index)

Spotted something that needs another look?

Help improve this page →