dask / dask/distributed

worker config set by config.set is not read by worker

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

Description

The configuration directly within Python is explained in the documentation here :
[Configuration - Directly within Python](https://docs.dask.org/en/latest/configuration.html#directly-within-python)

When using `dask.config.set`, I expect the worker to use those values. Instead, the worker reads the default values and does not use the values set using `dask.config.set`.

I modified distributed\worker.py as below to print the values received by the worker.

```python
if "memory_spill_fraction" in kwargs:
self.memory_spill_fraction = kwargs.pop("memory_spill_fraction")
print("self.memory_spill_fraction from kwargs = {}".format(self.memory_spill_fraction))
else:
self.memory_spill_fraction = dask.config.get(
"distributed.worker.memory.spill"
)
print("self.memory_spill_fraction from dask.config = {}".format(self.memory_spill_fraction))
```

```python
import dask
import dask.dataframe as dd

from dask.distributed import Client, LocalCluster

import pandas as pd

cluster = LocalCluster()
client = Client(cluster)

new = {"distributed.worker.memory.target": 0.1,
"distributed.worker.memory.spill": 0.2,
"distributed.worker.memory.pause": 0.3}

with dask.config.set(new):
print(dask.config.get("distributed.worker.memory"))
timestamp = pd.date_range('2018-01-01', periods=4, freq='S')
col1 = pd.Series(["1", "3", "5", "7"], dtype="string")
df = pd.DataFrame({"timestamp": timestamp,"col1": col1}).set_index('timestamp')
ddf = dd.from_pandas(df, npartitions=1)
ddf.compute()
ddf.head(2)

```
Outputs

```python
self.memory_spill_fraction from dask.config = 0.7
self.memory_spill_fraction from dask.config = 0.7
self.memory_spill_fraction from dask.config = 0.7
self.memory_spill_fraction from dask.config = 0.7
{'target': 0.1, 'spill': 0.2, 'pause': 0.3, 'terminate': 0.4}
```
Notice the 0.7 value which is the default.

Passing the configuration by kwargs works.

```python
import dask
import dask.dataframe as dd

from dask.distributed import Client, LocalCluster

import pandas as pd

cluster = LocalCluster(
memory_target_fraction=0.1,
memory_spill_fraction=0.2,
memory_pause_fraction=0.3)
client = Client(cluster)

timestamp = pd.date_range('2018-01-01', periods=4, freq='S')
col1 = pd.Series(["1", "3", "5", "7"], dtype="string")
df = pd.DataFrame({"timestamp": timestamp,"col1": col1}).set_index('timestamp')
ddf = dd.from_pandas(df, npartitions=1)
ddf.compute()
ddf.head(2)

```

Outputs

```python
self.memory_spill_fraction from kwargs = 0.2
self.memory_spill_fraction from kwargs = 0.2
self.memory_spill_fraction from kwargs = 0.2
self.memory_spill_fraction from kwargs = 0.2
```

**Environment**:

- Dask version: 2.18.1
- distributed version : 2.18.0
- Python version: 3.8.3
- Operating System: Windows
- Install method : pip

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.