Implemented Batch Processes#
LSF#
- class b2luigi.batch.processes.lsf.LSFJobStatusCache[source]#
Bases:
BatchJobStatusCache- _ask_for_job_status(job_id=None)[source]#
Queries the job status from the LSF batch system and updates the internal job status mapping.
- Parameters:
job_id (str, optional) – The ID of the job to query. If not provided, the status of all jobs will be queried.
Notes
This method uses the
bjobscommand-line tool to fetch job statuses in JSON format.The output is expected to contain a “RECORDS” key with a list of job records.
Each job record should have “JOBID” and “STAT” keys, which are used to update the internal mapping.
- class b2luigi.batch.processes.lsf.LSFProcess(*args, **kwargs)[source]#
Bases:
BatchProcessReference implementation of the batch process for a LSF batch system. Heavily inspired by this post.
Additional to the basic batch setup (see Batch Processing), there are LSF-specific
settings. These are:the LSF queue:
queue.the number of slots for the job. On KEKCC this increases the memory available to the job:
job_slots.the LSF job name:
job_name.
Parameter grouping (see Parameter Grouping) is supported: a grouped task is submitted as one
bsubcall per scalar sub-task, and the task is only reported successful once every one of these jobs has finished successfully.For example:
class MyLongTask(b2luigi.Task): queue = "l" job_name = "my_long_task"
The default queue is the short queue
"s". If nojob_nameis set the task will appear as<result_dir>/parameter1=value/.../executable_wrapper.sh"
when running
bjobs.By default, the environment variables from the scheduler are copied to the workers. This also implies we start in the same working directory, can reuse the same executable, etc. Normally, you do not need to supply
env_scriptor alike.- get_job_status() JobStatus[source]#
Determine the status of the task from the LSF status of all its jobs.
A plain task has exactly one job. A parameter-grouped task has one job per scalar sub-task, and their statuses are combined with
aggregate_job_status: running while any job runs, aborted if any job failed, successful only when all did. When the task is aborted, afailed_jobs.loglisting the failed job ids and their log directories is written into the task’s log directory (seewrite_failed_jobs_log).- Returns:
The aggregated status, or
JobStatus.abortedif no job was submitted.- Return type:
JobStatus
- static _get_job_status_for_id(job_id: str) JobStatus[source]#
Retrieves the current status of one batch job.
- Returns:
- The status of the job, which can be one of the following:
JobStatus.successful: If the job has completed successfully (“DONE”).JobStatus.aborted: If the job has been aborted or is not found in the cache (“EXIT” or missing ID).JobStatus.running: If the job is still in progress.
- Return type:
JobStatus
- start_job()[source]#
Submits a batch job to the LSF system.
This method constructs a command to submit a job using the
bsubcommand-line tool. It dynamically configures the job submission parameters based on task-specific settings and creates necessary log files for capturing standard output and error.For a parameter-grouped task (see Parameter Grouping) one
bsubcall is made per scalar sub-task that is not complete yet. If every sub-task is already complete, nothing is submitted and the task is reported as done.- Raises:
RuntimeError – If the batch submission fails or the job ID cannot be extracted from the
bsubcommand output.
- Steps:
Retrieve optional settings for
queue(-q),job_slots(-n), andjob_name(-J).The
stdoutandstderrlog files are created in the task’s log directory. Seeget_log_file_dir.The executable is created with
create_executable_wrapper.
- static _submit_task(task) str[source]#
Submit one task with
bsuband return its LSF job ID.- Parameters:
task – The task to submit. Its settings, log directory and executable wrapper are used.
- Returns:
The LSF job ID parsed from the
bsuboutput.- Return type:
str
- Raises:
RuntimeError – If the job ID cannot be extracted from the
bsuboutput.
HTCondor#
- class b2luigi.batch.processes.htcondor.HTCondorJobStatusCache[source]#
Bases:
BatchJobStatusCache- _ask_for_job_status(job_id: int = None)[source]#
With HTCondor, you can check the progress of your jobs using the
condor_qcommand. If noJobIdis given as argument, this command shows you the status of all queued jobs (usually only your own by default).Normally the HTCondor
JobIDis stated asClusterId.ProcId. Since only on job is queued per cluster, we can identify jobs by theirClusterId(TheProcIdwill be0for all submitted jobs). With the-jsonoption, thecondor_qoutput is returned in the JSON format. By specifying some attributes, not the entire jobClassAdis returned, but only the necessary information to match a job to itsJobStatus. The output is given as string and cannot be directly parsed into a json dictionary. It has the following form:[ {...} , {...} , {...} ]The
{...}are the different dictionaries including the specified attributes. Sometimes it might happen that a job is completed in between the status checks. Then its final status can be found in thecondor_historyfile (works mostly in the same way ascondor_q). Both commands are used in order to find out theJobStatus.
- class b2luigi.batch.processes.htcondor.HTCondorJobStatus(*values)[source]#
Bases:
IntEnumSee https://htcondor.readthedocs.io/en/latest/classad-attributes/job-classad-attributes.html
- class b2luigi.batch.processes.htcondor.HTCondorProcess(*args, **kwargs)[source]#
Bases:
BatchProcessReference implementation of the batch process for a HTCondor batch system.
Additional to the basic batch setup (see Batch Processing), additional HTCondor-specific things are:
Please note that most of the HTCondor batch farms do not have the same environment setup on submission and worker machines, so you probably want to give an
env_script, anenvsettingand/or a differentexecutable.HTCondor supports copying files from submission to workers. This means if the folder of your script(s)/python project/etc. is not accessible on the worker, you can copy it from the submission machine by adding it to the setting
transfer_files. This list can host both folders and files. Please note that due to HTCondors file transfer mechanism, all specified folders and files will be copied into the worker node flattened, so if you specify a/b/c.txt you will end up with a file c.txt. If you use thetransfer_filesmechanism, you need to set theworking_dirsetting to “.” as the files will end up in the current worker scratch folder. All specified files/folders should be absolute paths.Hint
Please do not specify any parts or the full results folder. This will lead to unexpected behavior. We are working on a solution to also copy results, but until this the results folder is still expected to be shared.
If you copy your python project using this setting to the worker machine, do not forget to actually set it up in your setup script. Additionally, you might want to copy your
settings.jsonas well.Via the
htcondor_settingssetting you can provide a dict as a for additional options, such as requested memory etc. Its value has to be a dictionary containing HTCondor settings as key/value pairs. These options will be written into the job submission file. For an overview of possible settings refer to the HTCondor documentation.Parameter grouping (see Parameter Grouping) is supported: a grouped task is expanded into one
queue 1block per scalar sub-task in a single submit file, and the task is only reported successful once every one of these jobs has finished successfully.Same as for the LSF, the
job_namesetting allows giving a meaningful name to a group of jobs. If you want to be htcondor-specific, you can provide theJobBatchNameas an entry in thehtcondor_settingsdict, which will override the globaljob_namesetting. This is useful for manually checking the status of specific jobs withcondor_q -batch <job name>
Example
1import b2luigi 2import random 3 4 5class MyNumberTask(b2luigi.Task): 6 some_parameter = b2luigi.IntParameter() 7 8 htcondor_settings = {"request_cpus": 1, "request_memory": "100 MB"} 9 10 def output(self): 11 yield self.add_to_output("output_file.txt") 12 13 def run(self): 14 print("I am now starting a task") 15 random_number = random.random() 16 17 if self.some_parameter == 3: 18 raise ValueError 19 20 with open(self.get_output_file_name("output_file.txt"), "w") as f: 21 f.write(f"{random_number}\n") 22 23 24class MyAverageTask(b2luigi.Task): 25 htcondor_settings = {"request_cpus": 1, "request_memory": "200 MB"} 26 27 def requires(self): 28 for i in range(10): 29 yield self.clone(MyNumberTask, some_parameter=i) 30 31 def output(self): 32 yield self.add_to_output("average.txt") 33 34 def run(self): 35 print("I am now starting the average task") 36 37 # Build the mean 38 summed_numbers = 0 39 counter = 0 40 for input_file in self.get_input_file_names("output_file.txt"): 41 with open(input_file, "r") as f: 42 summed_numbers += float(f.read()) 43 counter += 1 44 45 average = summed_numbers / counter 46 47 with open(self.get_output_file_name("average.txt"), "w") as f: 48 f.write(f"{average}\n") 49 50 51if __name__ == "__main__": 52 b2luigi.process(MyAverageTask(), workers=200, batch=True)
- static get_job_status_for_id(job_id)[source]#
Determines the status of a batch job based on its HTCondor job status.
- Returns:
- The status of the job, which can be one of the following:
JobStatus.successful: If the HTCondor job status is ‘completed’.JobStatus.running: If the HTCondor job status is one of ‘idle’, ‘running’, ‘transferring_output’, or ‘suspended’.JobStatus.aborted: If the HTCondor job status is ‘removed’, ‘held’, ‘failed’, or if the job ID is not found in the cache.
- Return type:
JobStatus
- Raises:
ValueError – If the HTCondor job status is unknown.
- get_job_status()[source]#
Implement this function to return the current job status. How you identify exactly your job is dependent on the implementation and needs to be handled by your own child class.
Must return one item of the JobStatus enumeration: running, aborted, successful or idle. Will only be called after the job is started but may also be called when the job is finished already. If the task status is unknown, return aborted. If the task has not started already but is scheduled, return running nevertheless (for
b2luigiit makes no difference). No matter if aborted via a call toterminate_job, by the batch system or by an exception in the job itself, you should return aborted if the job is not finished successfully (maybe you need to check the exit code of your job).
- start_job()[source]#
Starts a job by creating and submitting an HTCondor submit file.
This method generates an HTCondor submit file using the
_create_htcondor_submit_filemethod, then submits the job using thecondor_submitcommand.- Raises:
RuntimeError – If the batch submission fails or the job ID cannot be extracted from the
condor_submitoutput.
- terminate_job()[source]#
Terminates a batch job managed by HTCondor.
This method checks if a batch job ID is available. If a valid job ID exists, it executes the
condor_rmcommand to remove the job from the HTCondor queue.
- static _create_submit_file_content(task)[source]#
Creates an HTCondor submit file for the current task.
This method generates the content of an HTCondor submit file based on the task’s configuration and writes it to a file.
- Returns:
The path to the generated HTCondor submit file.
- Return type:
str
- Raises:
ValueError – If
transfer_filescontains non-absolute file paths or ifworking_diris not explicitly set to ‘.’ when usingtransfer_files.
Note
The
stdoutandstderrlog files are created in the task’s log directory. Seeget_log_file_dir.The HTCondor settings are specified in the
htcondor_settingssetting, which is a dictionary of key-value pairs.The executable is created with
create_executable_wrapper.The
transfer_filessetting can be used to specify files or directories to be transferred to the worker node.The
job_namesetting can be used to specify a meaningful name for the job.The submit file is named job.submit and is created in the task’s output directory (
get_task_file_dir).
Slurm#
- class b2luigi.batch.processes.slurm.SlurmJobStatusCache[source]#
Bases:
BatchJobStatusCache- _ask_for_job_status(job_id: int = None)[source]#
With Slurm, you can check the progress of your jobs using the
squeuecommand. If nojobIDis given as argument, this command shows you the status of all queued jobs.Sometimes it might happen that a job is completed in between the status checks. Then its final status can be found using
sacct(works mostly in the same way assqueue).If in the unlikely case the server has the Slurm accounting disabled, then the scontrol command is used as a last resort to access the jobs history. This is the fail safe command as the scontrol by design only holds onto a jobs information for a short period of time after completion. The time between status checks is sufficiently short however so the scontrol command should still have the jobs information on hand.
All three commands are used in order to find out the
SlurmJobStatus.
- _fill_from_output(output: str) set[source]#
Parses the output of a Slurm command to extract job IDs and their states, updating the internal job status mapping and returning a set of seen job IDs.
- Parameters:
output (str) – The output string from a Slurm command, expected to be formatted as ‘<job id> <state>’ per line.
- Returns:
A set of job IDs that were parsed from the output.
- Return type:
set
- Raises:
AssertionError – If a line in the output does not contain exactly two entries (job ID and state).
Notes
If the output is empty, an empty set is returned.
Lines in the output that are empty or contain unexpected formatting are skipped.
Job states with a ‘+’ suffix (e.g., ‘CANCELLED+’) are normalized by stripping the ‘+’ character.
- _get_SlurmJobStatus_from_string(state_string: str) str[source]#
Converts a state string into a
SlurmJobStatusenumeration value.- Parameters:
state_string (str) – The state string to be converted.
- Returns:
The corresponding
SlurmJobStatusvalue.- Return type:
str
- Raises:
KeyError – If the provided state string does not match any valid
SlurmJobStatus.
- _check_if_sacct_is_disabled_on_server() bool[source]#
Checks if the Slurm accounting command
sacctis disabled on the system.This method determines whether the
sacctcommand is unavailable or disabled by attempting to execute it and analyzing the output. The result is cached in theself.sacct_disabledattribute to avoid repeated checks.- Returns:
True if
sacctis disabled on the system,Falseotherwise.- Return type:
bool
- class b2luigi.batch.processes.slurm.SlurmJobStatus(*values)[source]#
Bases:
StrEnumSee https://slurm.schedmd.com/job_state_codes.html
- completed#
The job has completed successfully.
- Type:
str
- pending#
The job is waiting to be scheduled.
- Type:
str
- running#
The job is currently running.
- Type:
str
- configuring#
The job is allocated resources but waiting for nodes to be prepared.
- Type:
str
- suspended#
The job has been suspended.
- Type:
str
- preempted#
The job has been preempted by another job.
- Type:
str
- completing#
The job is in the process of completing.
- Type:
str
- boot_fail#
The job failed during the boot process.
- Type:
str
- cancelled#
The job was cancelled by the user or system.
- Type:
str
- deadline#
The job missed its deadline.
- Type:
str
- node_fail#
The job failed due to a node failure.
- Type:
str
- out_of_memory#
The job ran out of memory.
- Type:
str
- failed#
The job failed for an unspecified reason.
- Type:
str
- timeout#
The job exceeded its time limit.
- Type:
str
- static _generate_next_value_(name, start, count, last_values)#
Return the lower-cased version of the member name.
- class b2luigi.batch.processes.slurm.SlurmProcess(*args, **kwargs)[source]#
Bases:
BatchProcessReference implementation of the batch process for a Slurm batch system.
Additional to the basic batch setup (see Batch Processing), additional Slurm-specific things are:
Please note that most of the Slurm batch farms by default copy the user environment from the submission node to the worker machine. As this can lead to different results when running the same tasks depending on your active environment, you probably want to pass the argument
export=NONE. This ensures that a reproducible environment is used. You can provide anenv_script, anenvsetting, and/or a differentexecutableto create the environment necessary for your task.Via the
slurm_settingssetting you can provide a dict for additional options, such as requested memory etc. Its value has to be a dictionary containing Slurm settings as key/value pairs. These options will be written into the job submission file. For an overview of possible settings refer to the Slurm documentation <https://slurm.schedmd.com/sbatch.html#>_ and the documentation of the cluster you are using.Parameter grouping (see Parameter Grouping) is supported: a grouped task is submitted as one
sbatchcall per scalar sub-task, and the task is only reported successful once every one of these jobs has finished successfully.Same as for the LSF and HTCondor, the
job_namesetting allows giving a meaningful name to a group of jobs. If you want to be task-instance-specific, you can provide thejob-nameas an entry in theslurm_settingsdict, which will override the globaljob_namesetting. This is useful for manually checking the status of specific jobs withsqueue --name <job name>
Example
1import b2luigi 2import random 3 4 5class MyNumberTask(b2luigi.Task): 6 some_parameter = b2luigi.IntParameter() 7 batch_system = "slurm" 8 slurm_settings = {"export": "NONE", "ntasks": 1, "mem": "100MB"} 9 10 def output(self): 11 yield self.add_to_output("output_file.txt") 12 13 def run(self): 14 print("I am now starting a task") 15 random_number = random.random() 16 17 with open(self.get_output_file_name("output_file.txt"), "w") as f: 18 f.write(f"{random_number}\n") 19 20 21class MyAverageTask(b2luigi.Task): 22 batch_system = "slurm" 23 slurm_settings = {"export": "NONE", "ntasks": 1, "mem": "100MB"} 24 25 def requires(self): 26 for i in range(10): 27 yield self.clone(MyNumberTask, some_parameter=i) 28 29 def output(self): 30 yield self.add_to_output("average.txt") 31 32 def run(self): 33 print("I am now starting the average task") 34 35 # Build the mean 36 summed_numbers = 0 37 counter = 0 38 for input_file in self.get_input_file_names("output_file.txt"): 39 with open(input_file, "r") as f: 40 summed_numbers += float(f.read()) 41 counter += 1 42 43 average = summed_numbers / counter 44 45 with open(self.get_output_file_name("average.txt"), "w") as f: 46 f.write(f"{average}\n") 47 48 49if __name__ == "__main__": 50 b2luigi.process(MyAverageTask(), workers=200, batch=True)
- get_job_status() JobStatus[source]#
Determine the status of the task from the Slurm status of all its jobs.
A plain task has exactly one job. A parameter-grouped task has one job per scalar sub-task, and their statuses are combined with
aggregate_job_status: running while any job runs, aborted if any job failed, successful only when all did. When the task is aborted, afailed_jobs.loglisting the failed job ids and their log directories is written into the task’s log directory (seewrite_failed_jobs_log).- Returns:
The aggregated status, or
JobStatus.abortedif no job was submitted.- Return type:
JobStatus
- static _get_job_status_for_id(job_id: int) JobStatus[source]#
Determine the status of one batch job based on its Slurm job status.
- Returns:
- The status of the job, which can be one of the following:
JobStatus.successful: If the job has completed successfully.JobStatus.running: If the job is currently running, pending, suspended, preempted, or completing.JobStatus.aborted: If the job has failed, been cancelled, exceeded its deadline, encountered a node failure, ran out of memory, timed out, or if the job ID is not found.
- Return type:
JobStatus
- Raises:
ValueError – If the Slurm job status is unknown or not handled.
- start_job()[source]#
Starts a job by submitting the Slurm submission script.
This method creates a Slurm submit file and submits it using the
sbatchcommand. After submission, it parses the output to extract the batch job ID.For a parameter-grouped task (see Parameter Grouping) one submit file is created and submitted per scalar sub-task that is not complete yet. If every sub-task is already complete, nothing is submitted and the task is reported as done.
- Raises:
RuntimeError – If the batch submission fails or the job ID cannot be extracted.
- self._batch_job_ids#
The IDs of the submitted Slurm batch jobs.
- Type:
list[int]
- terminate_job()[source]#
Terminates all batch jobs of this task if any were submitted.
This method executes a single
scancelcommand with every submitted job ID (one for a plain task, several for a parameter-grouped task). The command’s output is suppressed.
- _create_slurm_submit_file(task=None)[source]#
Creates a Slurm submit file for a task.
This method generates a Slurm batch script that specifies the necessary configurations for submitting a job to a Slurm workload manager.
- Parameters:
task – The task to write the submit file for. Defaults to
self.task; a parameter-grouped task passes each scalar sub-task in turn.- Returns:
The path to the generated Slurm submit file.
- Return type:
pathlib.Path
Note
The
stdoutandstderrlog files are created in the task’s log directory. Seeget_log_file_dir.The Slurm settings are specified in the
slurm_settingssetting, which is a dictionary of key-value pairs.The
job_namesetting can be used to specify a meaningful name for the job.The executable is created with
create_executable_wrapper.The submit file is named slurm_parameters.sh and is created in the task’s output directory (
get_task_file_dir).
The gbasf2 batch process is documented in gbasf2.