DAG-Based Workflow Orchestration
Scalable, reproducible machine learning pipelines for particle physics analysis
Overview
Our system uses Directed Acyclic Graphs (DAGs) to represent complex ML workflows, enabling parallel execution and dependency management.
Workflow Architecture
Core Components
- Task Scheduler: The DAG Workflow is encoded by the dependency tree of luigi Tasks.
- Batch Submission: LAW (Luigi Analysis Workflow) is a wrapper for luigi that expands its functionality to batch submission in HEP using htcondor and slurm
- Dynamic Task Management: Changing the number of Tasks at each level of the graph is handled by a .yaml config manager
Key Features
- Automatic parallelization of independent tasks
- Checkpointing of intermediate results for identical trainings
- Support for heterogeneous computing (CPU/GPU)
- Dynamic resource optimisation
Data Ingestion
NEEDLE provides a set of ready-to-use torch dataloaders with different levels of complexity, depending on the computational demand of the training.
The base version loads data eagerly from parquet or root files in awkward arrays and returns a torch base dataloader with all data loaded in memory. The delayed version harnesses dask to lazily load the shape of the array, allowing for read-on-demand. This is especially useful in cases where the memory overhead is a critical limitation, in which case this dataloader caps the maximum concurrent memory usage. Further acceleration is provided by multiprocessing with either dask (IO) or torch (compute).