`SpecCluster` fails to remove grouped workers when they die, breaking adaptive scaling
I maintainer di solito rispondono entro 1 giorno
Nessuno ha ancora preso questa issue.
Valutazione
- Difficoltà
- 4/5
- Tempo stimato
- 3-5 giorni
- Idoneità per principianti
- 45/100
- Tipo di issue
- Bug
- Chiarezza
- Abbastanza chiara
- Stato di attività
- Ferma
- Stack tecnologico
- python
- Ambito
- distributed-systems
Direzione di ricerca
Inizia dal punto di ingresso _update_worker_status e traccia il modo in cui i nomi dei worker raggruppati provenienti da scheduler_info vengono associati a self.workers e self.worker_spec. Controlla la copertura dei test per adaptive scaling e SpecCluster, quindi verifica che i job raggruppati terminati vengano rimossi e che la capacità richiesta possa scalare nuovamente verso l’alto.
Scritto dal modello di indicizzazione a partire dal testo della issue.
Descrizione
When using SpecCluster with grouped workers (e.g., JobQueueCluster from dask-jobqueue with processes > 1, c.f. dask/dask-jobqueue#498), dead workers are not properly removed from self.workers or self.worker_spec. This causes the adaptive system to incorrectly believe workers are still "requested" when they have actually died, preventing scale-up.
Example
Apologies it's not a full reproducer
# Using SLURMCluster (subclass of JobQueueCluster -> SpecCluster)
from dask_jobqueue import SLURMCluster
cluster = SLURMCluster(
cores=8,
processes=8, # Creates grouped workers
memory="32GB",
walltime="01:00:00"
)
cluster.adapt(minimum=0, maximum=10)
# When jobs die (timeout, killed, etc.):
# - cluster.observed drops to 0 (scheduler sees no workers)
# - cluster.requested stays high (dead jobs still in self.workers)
# - Adaptive won't scale up new workers
log from Adaptive showing the issue
target=32 len(self.plan)=32 len(self.observed)=15 len(self.requested)=32
target=32 len(self.plan)=32 len(self.observed)=14 len(self.requested)=32
target=32 len(self.plan)=32 len(self.observed)=12 len(self.requested)=32
target=32 len(self.plan)=32 len(self.observed)=11 len(self.requested)=32
target=32 len(self.plan)=32 len(self.observed)=10 len(self.requested)=32
target=32 len(self.plan)=32 len(self.observed)=9 len(self.requested)=32
target=32 len(self.plan)=32 len(self.observed)=8 len(self.requested)=32
target=32 len(self.plan)=32 len(self.observed)=7 len(self.requested)=32
target=32 len(self.plan)=32 len(self.observed)=6 len(self.requested)=32
target=32 len(self.plan)=32 len(self.observed)=5 len(self.requested)=32
target=32 len(self.plan)=32 len(self.observed)=4 len(self.requested)=32
target=32 len(self.plan)=32 len(self.observed)=3 len(self.requested)=32
target=32 len(self.plan)=32 len(self.observed)=2 len(self.requested)=32
target=32 len(self.plan)=32 len(self.observed)=1 len(self.requested)=32
target=32 len(self.plan)=32 len(self.observed)=0 len(self.requested)=32
Context
Setup: JobQueueCluster creates grouped workers when processes > 1:
# In JobQueueCluster.__init__
if "processes" in self._job_kwargs and self._job_kwargs["processes"] > 1:
worker["group"] = ["-" + str(i) for i in range(self._job_kwargs["processes"])]
This creates a worker_spec entry like:
python"cluster-0": {
"cls": SLURMJob,
"options": {...},
"group": ["-0", "-1", "-2", "-3", "-4", "-5", "-6", "-7"]
}
The Bug: When workers die, _update_worker_status receives expanded names ("cluster-0-0") but tries to look them up in self.workers which uses job names ("cluster-0"):
# Current broken implementation
def _update_worker_status(self, op, msg):
if op == "remove":
name = self.scheduler_info["workers"][msg]["name"] # "cluster-0-0"
def f():
if name in self.workers: # self.workers has "cluster-0", not "cluster-0-0"
# This never executes for grouped workers!
del self.workers[name]
Result:
- self.workers keeps dead job objects
- self.worker_spec keeps their specifications
- adaptive doesn't scale up
- Lingua principale
- Python
- Stelle
- 1.7k
- Fork
- 778
- Merge medio
- 1h 36m
- PR unite (30g)
- 1
Preparare l'ambiente
Come iniziare
- Leggi tutta la issue e poi la guida ai contributi del progetto.
- Commenta sulla issue per dire che te ne occupi tu — evita che due persone facciano lo stesso lavoro.
- Fai un fork del repository e lavora su un branch.
- Apri una pull request che faccia riferimento al numero della issue.
Altre issue di dask/distributed
-
needs triage
Difficoltà 2/5 1-3 ore Idoneità per principianti 72/100
dask/distributed#9366 ·
I maintainer di solito rispondono entro 1 giorno
-
needs triage
Difficoltà 2/5 1-3 ore Idoneità per principianti 84/100
dask/distributed#9353 ·
I maintainer di solito rispondono entro 1 giorno
-
documentation
Difficoltà 1/5 1-3 ore Idoneità per principianti 82/100
dask/distributed#8304 ·
I maintainer di solito rispondono entro 1 giorno
-
Difficoltà 2/5 1-3 ore Idoneità per principianti 74/100
dask/distributed#4816 · 2 commenti ·
I maintainer di solito rispondono entro 1 giorno
-
documentation good first issue
Difficoltà 2/5 1-3 ore Idoneità per principianti 74/100
dask/distributed#2378 · 2 commenti ·
I maintainer di solito rispondono entro 1 giorno
Tutte le issue di dask/distributed
Issue simili
-
bug status/needs-triage
Difficoltà 2/5 1-3 ore Idoneità per principianti 86/100
prowler-cloud/prowler#12887 · 1 commento ·
I maintainer di solito rispondono entro 1 giorno
-
area: desktop platform: macos priority: p3 status: ready type: enhancement
Difficoltà 1/5 Meno di un'ora Idoneità per principianti 92/100
use-agent-os/agent-os#3484 ·
I maintainer di solito rispondono entro 2 giorni
-
bug
Difficoltà 2/5 1-3 ore Idoneità per principianti 86/100
open-telemetry/opentelemetry-python-contrib#5113 · 2 commenti · 2 reazioni ·
I maintainer di solito rispondono entro 1 giorno
-
external
Difficoltà 2/5 1-3 ore Idoneità per principianti 68/100
langchain-ai/docs#6255 ·
I maintainer di solito rispondono entro 1 giorno
-
Difficoltà 2/5 1-3 ore Idoneità per principianti 72/100
I maintainer di solito rispondono entro 1 giorno