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 bjobs command-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: BatchProcess

Reference 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 bsub call 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 no job_name is 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_script or 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, a failed_jobs.log listing the failed job ids and their log directories is written into the task’s log directory (see write_failed_jobs_log).

Returns:

The aggregated status, or JobStatus.aborted if 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 bsub command-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 bsub call 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 bsub command output.

Steps:
  1. Retrieve optional settings for queue (-q), job_slots (-n), and job_name (-J).

  2. The stdout and stderr log files are created in the task’s log directory. See get_log_file_dir.

  3. The executable is created with create_executable_wrapper.

static _submit_task(task) → str[source]#

Submit one task with bsub and 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 bsub output.

Return type:

str

Raises:

RuntimeError – If the job ID cannot be extracted from the bsub output.

terminate_job()[source]#

Terminate all batch jobs of this task if any were submitted, with a single bkill command listing every job ID. The command’s output is suppressed, and errors during execution are not raised.

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_q command. If no JobId is given as argument, this command shows you the status of all queued jobs (usually only your own by default).

Normally the HTCondor JobID is stated as ClusterId.ProcId. Since only on job is queued per cluster, we can identify jobs by their ClusterId (The ProcId will be 0 for all submitted jobs). With the -json option, the condor_q output is returned in the JSON format. By specifying some attributes, not the entire job ClassAd is returned, but only the necessary information to match a job to its JobStatus. 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 the condor_history file (works mostly in the same way as condor_q). Both commands are used in order to find out the JobStatus.

_fill_from_output(output)[source]#

Processes the output from an HTCondor job query and updates the job statuses.

Parameters:

output (bytes) – The raw output from the HTCondor job query, encoded as bytes.

Returns:

A set of ClusterId values representing the jobs seen in the output.

Return type:

set

class b2luigi.batch.processes.htcondor.HTCondorJobStatus(*values)[source]#

Bases: IntEnum

See https://htcondor.readthedocs.io/en/latest/classad-attributes/job-classad-attributes.html

class b2luigi.batch.processes.htcondor.HTCondorProcess(*args, **kwargs)[source]#

Bases: BatchProcess

Reference 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, an env setting and/or a different executable.

  • 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 the transfer_files mechanism, you need to set the working_dir setting 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.json as well.

  • Via the htcondor_settings setting 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 1 block 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_name setting allows giving a meaningful name to a group of jobs. If you want to be htcondor-specific, you can provide the JobBatchName as an entry in the htcondor_settings dict, which will override the global job_name setting. This is useful for manually checking the status of specific jobs with

    condor_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 b2luigi it makes no difference). No matter if aborted via a call to terminate_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_file method, then submits the job using the condor_submit command.

Raises:

RuntimeError – If the batch submission fails or the job ID cannot be extracted from the condor_submit output.

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_rm command 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_files contains non-absolute file paths or if working_dir is not explicitly set to ‘.’ when using transfer_files.

Note

  • The stdout and stderr log files are created in the task’s log directory. See get_log_file_dir.

  • The HTCondor settings are specified in the htcondor_settings setting, which is a dictionary of key-value pairs.

  • The executable is created with create_executable_wrapper.

  • The transfer_files setting can be used to specify files or directories to be transferred to the worker node.

  • The job_name setting 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 squeue command. If no jobID is 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 as squeue).

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 SlurmJobStatus enumeration value.

Parameters:

state_string (str) – The state string to be converted.

Returns:

The corresponding SlurmJobStatus value.

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 sacct is disabled on the system.

This method determines whether the sacct command is unavailable or disabled by attempting to execute it and analyzing the output. The result is cached in the self.sacct_disabled attribute to avoid repeated checks.

Returns:

True if sacct is disabled on the system, False otherwise.

Return type:

bool

class b2luigi.batch.processes.slurm.SlurmJobStatus(*values)[source]#

Bases: StrEnum

See 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: BatchProcess

Reference 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 an env_script, an env setting, and/or a different executable to create the environment necessary for your task.

  • Via the slurm_settings setting 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 sbatch call 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_name setting allows giving a meaningful name to a group of jobs. If you want to be task-instance-specific, you can provide the job-name as an entry in the slurm_settings dict, which will override the global job_name setting. This is useful for manually checking the status of specific jobs with

    squeue --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, a failed_jobs.log listing the failed job ids and their log directories is written into the task’s log directory (see write_failed_jobs_log).

Returns:

The aggregated status, or JobStatus.aborted if 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 sbatch command. 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 scancel command 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 stdout and stderr log files are created in the task’s log directory. See get_log_file_dir.

  • The Slurm settings are specified in the slurm_settings setting, which is a dictionary of key-value pairs.

  • The job_name setting 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.