dask / dask/distributed

SpecCluster correct state can cause inconsistencies

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

Description

There are a bunch of problems in `SpecCluster._correct_state_internal` such that I believe it should be rewritten

* If one worker fails during startup, all workers are rejected. This can cause the cluster to spin up too many workers
* Similarly, while closing, if one worker fails but others properly shut down, they are not removed from the internal state
* `correct_state_internal` is only called once, i.e. any kind of exception would abort the entire up/downscaling without further attempt to correct the state.
* If the cluster is closing while correct_state is running, nothing is actually cancelled.
* `self._correct_state_waiting` is actually never cancelled.
* `SpecCluster.scale` schedules a callback to self `_correct_state`. This can cause all sorts of race conditions, e.g. by creating more futures even if the cluster is already closing

This list probably continues long and most issues _could_ be addressed individually but I believe we're better off rewriting this section.

cc @graingert

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.