`SpecCluster` fails to remove grouped workers when they die, breaking adaptive scaling
Nadie ha tomado este issue todavía.
Evaluación
- Dificultad
- 4/5
- Tiempo estimado
- 3-5 días
- Aptitud para principiantes
- 45/100
- Tipo de issue
- Error
- Claridad
- Bastante claro
- Estado de actividad
- Estancado
- Stack tecnológico
- python
- Área
- distributed-systems
Línea de trabajo
Comienza en el punto de entrada _update_worker_status y sigue cómo los nombres de worker agrupados de scheduler_info se asignan a self.workers y self.worker_spec. Comprueba la cobertura de pruebas de adaptive scaling y SpecCluster, y verifica después que los trabajos agrupados muertos se eliminen y que la capacidad solicitada pueda volver a escalarse hacia arriba.
Escrito por el modelo de indexación a partir del texto del issue.
Descripción
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
- Lenguaje dominante
- Python
- Estrellas
- 1.7k
- Forks
- 780
- Métricas de merge de PR
- Sin PR fusionados en 30 d
Preparar el entorno
- Sin Dockerfile ni archivo de Docker Compose
- Tiene una plantilla de pull request
- Leer la guía de contribución
Primeros pasos
- Lee el issue completo y luego la guía de contribución del proyecto.
- Comenta en el issue que vas a ocuparte — evita que dos personas hagan lo mismo.
- Haz un fork del repositorio y trabaja en una rama.
- Abre un pull request que haga referencia al número del issue.
Más de dask/distributed
-
needs triage
Dificultad 2/5 1-3 horas Aptitud para principiantes 72/100
dask/distributed#9366 ·
-
needs triage
Dificultad 2/5 1-3 horas Aptitud para principiantes 84/100
dask/distributed#9353 ·
-
documentation
Dificultad 1/5 1-3 horas Aptitud para principiantes 82/100
dask/distributed#8304 ·
-
Dificultad 2/5 1-3 horas Aptitud para principiantes 74/100
dask/distributed#4816 · 2 comentarios ·
-
documentation good first issue
Dificultad 2/5 1-3 horas Aptitud para principiantes 74/100
dask/distributed#2378 · 2 comentarios ·
Todos los issues de dask/distributed
Issues similares
-
Dificultad 1/5 Menos de una hora Aptitud para principiantes 72/100
letsencrypt/cp-cps#353 ·
-
Marble Madness II is missingAbierto
Dificultad 2/5 1-3 horas Aptitud para principiantes 68/100
-
Dificultad 2/5 1-3 horas Aptitud para principiantes 84/100
PedestrianDynamics/pyFDS-Evac#394 ·
Los mantenedores suelen responder en 1 día
-
Dificultad 2/5 1-3 horas Aptitud para principiantes 78/100
DOI-USGS/pywatershed#421 ·
-
Dificultad 2/5 1-3 horas Aptitud para principiantes 78/100
python-pillow/Pillow#10087 · 1 comentario ·
Los mantenedores suelen responder en 1 día