worker config set by config.set is not read by worker
- 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
Assessment
This issue has not been assessed yet.