dask / dask/distributed

Scale by number of cores or amount of memory

Open
#2,208 9 comments 1 reaction 0 assignees View on GitHub
Dominant language
Python
Stars
1.7k
Forks
778
Avg merge
2h 50m
Merged PRs (30d)
3

Description

When creating a cluster object we currently scale by number of workers

cluster = KubeCluster()
cluster.scale(10)

Where `10` is the number of workers we want to have. However, it is common for users to think about clusters in terms of number of cores or amount of memory, rather than in terms of number of dask workers

cluster.scale(cores=100)
cluster.scale(memory='1 TB')

What is the best way to achieve this uniformly across the dask deployment projects? I currently see two approaches, though there are probably more that others might see.

1. Establish a convention where clusters define information about the workers they will produce, something like the following:

```python
>>> cluster.worker_info
{'cores': 4, 'memory': '16 GB'}
```

Then the core `Cluster.scale` method would translate this into number of workers and then call the subclass's `scale` method appropriately

2. Let the downstream classes handle this themselves, but ask them all to handle it uniformly. This places more burden onto downstream implementations, but also gives them more freedom to select worker types as they see fit based on their capabilities.

cc

- @guillaumeeb (who has shown interest in doing this work) @lesteve @jhamman from `dask-jobqueue`
- @jcrist from `dask-yarn`
- @jacobtomlinson from `dask-kubernetes`

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.