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

Restart cluster job on task completion

Đang mở
#597 3 bình luận 0 reaction 0 người được giao Xem trên GitHub

Chưa có ai nhận issue này.

Đánh giá

Độ khó
5/5
Thời gian dự kiến
Hơn một tuần
Mức phù hợp với người mới
25/100
Loại issue
Tính năng
Độ rõ ràng
Cần làm rõ
Mức độ hoạt động
Đình trệ
Công nghệ
python
Lĩnh vực
distributed-systems, hpc

Hướng nghiên cứu

Bắt đầu với cấu hình SLURMCluster và các ví dụ về WorkerPlugin và client.retire_workers trong issue. So sánh vòng đời của worker và hành vi cân bằng tác vụ của chúng với cách tiếp cận scheduler-plugin được đề xuất. Được xem là hoàn thành khi có một cách được tài liệu hóa và được hỗ trợ để loại bỏ hoặc khởi động lại worker sau khi hoàn thành tác vụ mà không lặp lại tác vụ hoặc bắt đầu công việc khi gần đến walltime.

Do mô hình lập chỉ mục viết ra từ nội dung của issue.

Mô tả

Use case conditions:

  • tasks have large and variable run times
  • task executes 3rd party software, such that dask cannot migrate the execution state
  • workers have a wall time eg HPC cluster

Current behavior:

Compute is wasted on tasks that cannot finish in remaining walltime, task is restarted from scratch on new worker after worker death.

Originally working with this issue on the forums here

More context:

Task executes a 3rd party software (mine requiring multiple threads, see the issues linked in the above post for examples). Code looks something like the following:

import dask
from dask_jobqueue import SLURMCluster
from distributed import Client
import distributed
import logging
import time
            
def do_one(x):
    worker = distributed.get_worker()
    logger = logging.getLogger('worker')
    logger.setLevel(logging.INFO)
    fh = logging.FileHandler(f'worker_{worker.id}.log', mode='w')
    fh.setLevel(logging.INFO)
    logger.addHandler(fh)

    logger.info(f"I am working on {x}")
    # run third party software
    # this takes a while but is not very consistent in total time
    logger.info(f"I finished {x}")
    return f"Input {x} done"

if __name__ =='__main__':
    cluster = SLURMCluster(
        memory="1g",
        walltime='00:30:00',
        job_extra_directives=['--nodes=1', '--ntasks-per-node=1'],
        cores=1,
        processes=1,
        worker_extra_args=["--lifetime", "28m", "--lifetime-stagger", "50s"],
        job_cpu=6
    )
    cluster.adapt(minimum=2, maximum=10)
    client = Client(cluster)
    
    results = []
    for future in distributed.as_completed(client.map(
        do_one, list(range(100,132))
    )):
        result = future.result()
        results.append(result)

Result of worker_XXX.log

2022-11-07-12:00:00 INFO I am working on 1
2022-11-07-12:22:00 INFO I finished 1
2022-11-07-12:22:03 INFO I am working on 10

Worker XXX is killed at 12:29 due to walltime. 7 minutes of compute is wasted because the state cannot be changed. Task 10 starts from scratch on a new worker.

Attempts to fix:

Short of figuring out a way to move the execution state, I figure the best strategy is to have each task get a brand new SLURM job, so that no compute is wasted and any task that can finish in the walltime works.

  1. I tried a worker plugin like so:
class KillerNannyPlugin(distributed.diagnostics.plugin.WorkerPlugin):
    """Better as a nanny plugin but those are not running transitions properly."""
    def __init__(self, max_stagger_seconds: float = 5):
        self.max_stagger_seconds = max_stagger_seconds
    
    def setup(self, worker):
        self.worker = worker
        
    def transition(self, key, start, finish, *args, **kwargs):
        if start == 'memory' and finish == 'released':
            self.worker.io_loop.call_later(3+random.random() * self.max_stagger_seconds, self.worker.close_gracefully, restart=True)
  • This was successful in ensuring each task got its own job, but caused task repeat to be on the order of 100%, defeating the point of saving compute
  1. Have the client retire the worker that just completed a task when the job is done, like so:
for future in distributed.as_completed(client.map(
        do_one, list(range(100,132))
    )):
        who_has = client.who_has(future)
        closing = list(list(who_has.values())[0])
        client.retire_workers(closing)
  • I also added a small time delay to the worker function such that next tasks did not start (and begin wasting energy) while the client retired the worked.
  • This seems to have the desired effect, any consequences of this are not clear to me as I observe that tasks are not repeated nor do tasks start on a job that is about to time out.

I think this should be codified somehow as the "solution" above is quite hacky. My intuition says that it would fit best as a scheduler plugin, as using the worker plugin above clearly had adverse effects on task balancing. Happy to help contribute with some input on where this would fit best if it would be a useful addition.

Ngôn ngữ chính
Python
Star
256
Fork
149
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/dask-jobqueue

Tất cả issue của dask/dask-jobqueue

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.