Batch Processing#
As shown in Quick Start, using a batch system instead of local processing is really just a --batch
on the command line or calling process with batch=True.
However, there is more to discover!
Choosing the batch system#
Using b2luigi’s settings mechanism (described here b2luigi.get_setting()) you can choose which
batch system should be used.
Currently, htcondor and lsf are supported, with lsf``being the default setting.
There is also a wrapper for ``gbasf2, the Belle II
submission tool for the LHC Worldwide Computing Grid, which works for Basf2PathTask tasks
(see gbasf2).
In addition, it is possible to set the batch_system setting to auto which tries to detect which batch system is
available on your system. The automated discovery checks for the submission tools present on the system and sets the
BatchProcess() accordingly. This functionality works for lsf and htcondor systems. gbasf2 will not
be detected by auto but needs to be set explicitly.
Choosing the Environment#
If you are doing a local calculation, all calculated tasks will use the same environment (e.g. $PATH setting, libraries etc.)
as you have currently set up when calling your script(s).
This makes it predictable and simple.
Things get a bit more complicated when using a batch farm, as the workers might not have the same environment set up. The batch submission does not copy the environment (or the local site administrators have forbidden that) or the system on the workers is so different that copying the environment from the scheduling machine does not make sense.
Therefore b2luigi provides you with three mechanism to set the environment for each task:
You can give a bash script in the
env_scriptsetting (viaset_setting(),settings.jsonor for each task as usual, seeb2luigi.get_setting()), which will be called even before anything else on the worker. Use it to set up things like the path variables or the libraries (e.g. when you are using a virtual environment) and your batch system does not support environment copy from the scheduler to the workers. For example a useful script might look like this:# Source my virtual environment source venv/bin/activate # Set some specific settings export MY_IMPORTANT_SETTING 10
You can set the
envsetting to a dictionary, which contains additional variables to be set up before your job runs. By using the mechanism described inb2luigi.get_setting(), it is possible to make this task- or even parameter-dependent.By default,
b2luigire-uses the samepythonexecutable on the workers as you used to schedule the tasks (by calling your script). In some cases, this specific python executable is not present on the worker or is not usable (e.g. because of different operation systems or architectures). You can choose a new executable with theexecutablesetting (it is also possible to just usepython3as the executable assuming it is in the path). The executable needs to be callable after yourenv_scriptor your specificenvsettings are used. Please note, that theenvironmentsetting is a list, so you need to pass your python executable with possible arguments like this:b2luigi.set_setting("executable", ["python3"])
File System#
Depending on your batch system, the filesystem on the worker processing the task and the scheduler machine can be different or even unrelated. Different batch systems and batch systems implementations treat this fact differently. In the following, the basic procedure and assumption is explained. Any deviation from this is described in the next section.
By default, b2luigi needs at least two folders to be accessible from the scheduling as well as worker machine:
the result folder and the folder of your script(s).
If possible, use absolute paths for the result and log directory to prevent any problems.
Some batch systems (e.g. htcondor) support file copy mechanisms from the scheduler to the worker systems.
Please checkout the specifics below.
Hint
All relative paths given to e.g. the result_dir or the log_dir are always evaluated
relative to the folder where your script lives.
To prevent any disambiguities, try to use absolute paths whenever possible.
Some batch system starts the job in an arbitrary folder on the workers instead of the current folder on the scheduler.
That is why b2luigi will change the directory into the path of your called script before starting the job.
In case your script is accessible from a different location on the worker than on the scheduling machine, you can give the setting working_dir
to specify where the job should run.
Your script needs to be in this folder and every relative path (e.g. for results or log files) will be evaluated from there.
Note
When submitting through the b2luigi command line interface, the task file is sent to the worker
exactly as you gave it: a relative path stays relative and is resolved against working_dir,
an absolute path stays absolute.
Leave it relative (the default, tasks.py) if you want a relocating working_dir to work: the
worker changes into working_dir and finds the task file at the same position inside the copy of
the project living there. Two common deployments need this — staging the project into a scratch
directory on the node, where the submission host’s paths do not exist at all, and running jobs from
a production checkout while you develop in another, where an absolute path would send the worker
back to your development copy.
Pass an absolute --task-file when you mean a fixed location that is valid on the worker too.
The parameters file follows the same rule. The worker imports it before running the task, so a
b2luigi.set_setting(...) call made in parameters.py is in force on the worker exactly as
it is on the submission host. The config dict itself is not read there; every parameter value
travels on the command line.
Drawbacks of the batch mode#
Although the batch mode has many benefits, it would be unfair to not mention its downsides:
You have to choose the queue/batch settings/etc. depending in your requirements (e.g. wall clock time) by yourself. So you need to make sure that the tasks will actually finish before the batch system terminates them because of timeout. There is just no way for
b2luigito know this beforehand.There is currently no resubmission implemented. This means dying jobs because of batch system failures are just dead. But because of the dependency checking mechanism of
luigiit is simple to just redo the calculation and re-calculate what is missing.The
luigifeature to request new dependencies while task running (viayield) is not implemented for the batch mode so far.
Implementing your own Batch Process#
If you want to implement your own batch process, please have a look at the existing implementations.
You can create your own class inheriting from b2luigi.batch.processes.BatchProcess and implement the required methods.
Then set the batch_system setting to custom and make sure your task has the process_class property implemented.
The property should return your custom class.
As we consider this to be a rather advanced feature, we do not provide a step-by-step guide here.
The main goal of this feature is to allow experts to develop new batch system implementations easily.
Feel free to reach out to us, if the batch system you want is not yet implemented and you need help with the implementation.
Batch System Specific Settings#
Every batch system has special settings. You can look them up here:
LSF#
- 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.
HTCondor#
- 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)
Slurm#
- 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)
GBasf2 Wrapper for LCG#
The gbasf2 batch process and its settings are documented in gbasf2.
Apptainer#
- class b2luigi.batch.processes.apptainer.ApptainerProcess(*args, **kwargs)[source]
Bases:
BatchProcessSimple implementation of a batch process for running jobs in an Apptainer container. Strictly speaking, this is not a batch process, but it is a simple way to run jobs in a container environment.
This process inherits the basic properties from
b2luigi.batch.processes.BatchProcessbut does not need to be executed in thebatchcontext. However, running inbatchmode is possible for thelsfand thehtcondorbatch systems. Although, for the latter batch system it is not recommended to use apptainer images since HTCondor is already running in a container environment.The core principle of this process is to run the task in an Apptainer container. To achieve the execution of tasks, an
apptainer execcommand is build within this class and executed in a subprocess. To steer the execution, one can use the following settings:apptainer_image: The image to use for the Apptainer container.sThis parameter is mandatory and needs to be set if the task should be executed in an Apptainer container. The image needs to be accessible from the machine where the task is executed. There are no further checks if the image is available or valid. When using custom images, it may be helpful to first check the image with
apptainer inspect. For people with access to the Belle II own/cvmfsdirectory, images are provided in the/cvmfs/belle.cern.ch/imagesdirectory. The description of the images (the repository contains the docker images which are transformed to Apptainer images) and instructions on how to create them can be found in https://gitlab.desy.de/belle2/software/docker-images.
apptainer_mounts: A list of directories to mount into the Apptainer container.This parameter is optional and can be used to mount directories into the Apptainer container. The directories need to be accessible from the machine where the task is executed. The directories are mounted under the exact same path as they are provided/on the host machine. For most usecases mounts need to be provided to access software or data locations. For people using for example
basf2software in the Apptainer container, the/cvmfsdirectory needs to be mounted. Caution is required when system specific directories are mounted.
apptainer_mount_defaults: Boolean parameter to mountlog_dirandresult_dirby default.The default value is
Truemeaning theresult_dirandlog_dirare automatically created and mounted if they are not accessible from the execution location. When using custom targets with non local output directories, this parameter should be set toFalseto avoid mounting non-existing directories.
apptainer_additional_params: Additional parameters to pass to theapptainer execcommand.This parameter should be a string and will be directly appended to the
apptainer execcommand. It can be used to pass additional parameters to theapptainer execcommand as they would be added in the CLI. A very useful parameter is the--cleanenvparameter which will clean the environment before executing the task in the Apptainer container. This can be useful to avoid conflicts with the environment in the container. A prominent usecase is the usage of software which depends on the operating system.
A simple example of how an Apptainer based task can be defined is shown below:
class MyApptainerTask(luigi.Task): apptainer_image = "/cvmfs/belle.cern.ch/images/belle2-base-el9" apptainer_mounts = ["/cvmfs"] apptainer_mount_defaults = True apptainer_additional_params = "--cleanenv" <rest of the task definition>
Add your own batch system#
If you want to add a new batch system, all you need to do is to implement the
abstract functions of BatchProcess for your system:
- class b2luigi.batch.processes.BatchProcess(task, scheduler, result_queue, worker_timeout)[source]
This is the base class for all batch algorithms that allow luigi to run on a specific batch system. This is an abstract base class and inheriting classes need to supply functionalities for
starting a job using the commands in
self.task_cmdgetting the job status of a running, finished or failed job
and terminating a job
All those commands are called from the main process, which is not running on the batch system. Every batch system that is capable of these functions can in principle work together with
b2luigi.- Implementation note:
In principle, using the batch system is transparent to the user. In case of problems, it may however be useful to understand how it is working.
When you start your
luigidependency tree withprocess(..., batch=True), the normalluigiprocess is started looking for unfinished tasks and running them etc. Normally, luigi creates a process for each running task and runs them either directly or on a different core (if you have enabled more than one worker). In the batch case, this process is not a normal python multiprocessing process, but thisBatchProcess, which has the same interface (one can check the status of the process, start or terminate it). The process does not need to wait for the batch job to finish but is asked repeatedly for the job status. By this, most of the core functionality ofluigiis kept and reused. This also means, that every batch job only includes a single task and is finished whenever this task is done decreasing the batch runtime. You will need exactly as many batch jobs as you have tasks and no batch job will idle waiting for input data as all are scheduled only when the task they should run is actually runnable (the input files are there).What is the batch command now? In each job, we call a specific executable bash script only created for this task. It contains the setup of the environment (if given by the user via the settings), the change of the working directory (the directory of the python script or a specified directory by the user) and a call of this script with the current python interpreter (the one you used to call this main file or given by the setting
executable) . However, we give this call an additional parameter, which tells it to only run one single task. Task can be identified by their task id. A typical task command may look like:/<path-to-your-exec>/python /your-project/some-file.py --batch-runner --task-id MyTask_38dsf879w3
if the batch job should run the
MyTask. The implementation of the abstract functions is responsible for creating an running the executable file and writing the log of the job into appropriate locations. You can use the functionscreate_executable_wrapperandget_log_file_dirto get the needed information.Checkout the implementation of the
lsftask for some implementation example.
- 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]
Override this function in your child class to start a job on the batch system. It is called exactly once. You need to store any information identifying your batch job on your own.
You can use the
b2luigi.core.utils.get_log_file_dirand theb2luigi.core.executable.create_executable_wrapperfunctions to get the log base name and to create the executable script which you should call in your batch job.After the
start_jobfunction is called by the framework (and no exception is thrown), it is assumed that a batch job is started or scheduled.After the job is finished (no matter if aborted or successful) we assume the stdout and stderr is written into the two files given by
b2luigi.core.utils.get_log_file_dir(self.task).
- terminate_job()[source]
This command is used to abort a job started by the
start_jobfunction. It is only called once to abort a job, so make sure to either block until the job is really gone or be sure that it will go down soon. Especially, do not wait until the job is finished. It is called for example when the user pressesCtrl-C.In some strange corner cases it may happen that this function is called even before the job is started (the
start_jobfunction is called). In this case, you do not need to do anything (but also not raise an exception).