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).

Back to Project Overview