`SpecCluster` fails to remove grouped workers when they die, breaking adaptive scaling
Chưa có ai nhận issue này.
Đánh giá
- Độ khó
- 4/5
- Thời gian dự kiến
- 3-5 ngày
- Mức phù hợp với người mới
- 45/100
- Loại issue
- Lỗi
- Độ rõ ràng
- Khá rõ ràng
- Mức độ hoạt động
- Đình trệ
- Công nghệ
- python
- Lĩnh vực
- distributed-systems
Hướng nghiên cứu
Bắt đầu từ entry point _update_worker_status và theo dõi cách các tên worker được nhóm từ scheduler_info ánh xạ tới self.workers và self.worker_spec. Kiểm tra độ bao phủ kiểm thử của adaptive scaling và SpecCluster, sau đó xác minh rằng các job được nhóm đã chết được xóa và capacity được yêu cầu có thể scale up trở lại.
Do mô hình lập chỉ mục viết ra từ nội dung của issue.
Mô tả
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
- Ngôn ngữ chính
- Python
- Star
- 1.7k
- Fork
- 778
- Chỉ số merge pull request
- Không có pull request nào được merge trong 30 ngày
Chuẩn bị môi trường
Bắt đầu từ đâu
- Đọc hết issue, rồi đọc hướng dẫn đóng góp của dự án.
- Bình luận trên issue rằng bạn sẽ nhận — tránh hai người làm cùng một việc.
- Fork repository và làm thay đổi trên một nhánh.
- Mở pull request có tham chiếu số hiệu của issue.
Issue khác của dask/distributed
-
needs triage
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 72/100
dask/distributed#9366 ·
-
needs triage
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 84/100
dask/distributed#9353 ·
-
documentation
Độ khó 1/5 1-3 giờ Mức phù hợp với người mới 82/100
dask/distributed#8304 ·
-
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 74/100
dask/distributed#4816 · 2 bình luận ·
-
documentation good first issue
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 74/100
dask/distributed#2378 · 2 bình luận ·
Tất cả issue của dask/distributed
Issue tương tự
-
bug status/needs-triage
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 86/100
prowler-cloud/prowler#12887 · 1 bình luận ·
Maintainer thường phản hồi trong vòng 1 ngày
-
area: desktop platform: macos priority: p3 status: ready type: enhancement
Độ khó 1/5 Dưới một giờ Mức phù hợp với người mới 92/100
use-agent-os/agent-os#3484 ·
Maintainer thường phản hồi trong vòng 2 ngày
-
bug
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 86/100
open-telemetry/opentelemetry-python-contrib#5113 · 2 bình luận · 2 reaction ·
Maintainer thường phản hồi trong vòng 1 ngày
-
external
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 68/100
langchain-ai/docs#6255 ·
Maintainer thường phản hồi trong vòng 1 ngày
-
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 72/100
Maintainer thường phản hồi trong vòng 1 ngày