Hacktoberfest 2026: le issue che i maintainer hanno segnato per ottobre, aperte e adatte ai principianti. Sfoglia le issue Hacktoberfest

Restart cluster job on task completion

Aperta
#597 3 commenti 0 reazioni 0 assegnatari Vedi su GitHub

Nessuno ha ancora preso questa issue.

Valutazione

Difficoltà
5/5
Tempo stimato
Più di una settimana
Idoneità per principianti
25/100
Tipo di issue
Funzionalità
Chiarezza
Da chiarire
Stato di attività
Ferma
Stack tecnologico
python

Direzione di ricerca

Inizia dalla configurazione di SLURMCluster e dagli esempi di WorkerPlugin e client.retire_workers nell’issue. Confronta il loro ciclo di vita dei worker e il comportamento di bilanciamento delle attività con l’approccio proposto basato su scheduler-plugin. Il lavoro è completato quando esiste un modo documentato e supportato per ritirare o riavviare i worker dopo il completamento delle attività, senza ripetere le attività né iniziare lavori in prossimità di walltime.

Scritto dal modello di indicizzazione a partire dal testo della issue.

Descrizione

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.

Lingua principale
Python
Stelle
256
Fork
149
Metriche di merge delle PR
Nessuna PR unita negli ultimi 30g

Preparare l'ambiente

Come iniziare

  1. Leggi tutta la issue e poi la guida ai contributi del progetto.
  2. Commenta sulla issue per dire che te ne occupi tu — evita che due persone facciano lo stesso lavoro.
  3. Fai un fork del repository e lavora su un branch.
  4. Apri una pull request che faccia riferimento al numero della issue.

Altre issue di dask/dask-jobqueue

Tutte le issue di dask/dask-jobqueue

Issue simili

Altre issue su Python

Ricevi le nuove issue nella tua casella

Un breve riepilogo di issue GitHub adatte ai principianti.