Hacktoberfest 2026: những issue maintainer đã đánh dấu cho tháng Mười, đang mở và phù hợp người mới. Xem issue Hacktoberfest

`SpecCluster` fails to remove grouped workers when they die, breaking adaptive scaling

Đang mở
#9,102 2 bình luận 1 reaction 0 người được giao Xem trên GitHub

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ả

needs info

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

  1. Đọc hết issue, rồi đọc hướng dẫn đóng góp của dự án.
  2. 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.
  3. Fork repository và làm thay đổi trên một nhánh.
  4. Mở pull request có tham chiếu số hiệu của issue.

Issue khác của dask/distributed

Tất cả issue của dask/distributed

Issue tương tự

Thêm issue về Python

Nhận issue mới trong hộp thư của bạn

Bản tóm tắt ngắn những issue GitHub phù hợp với người mới.