Skip to content

dqm_ml_job.dataloaders

Data loaders module for DQM ML Job.

This module contains classes for loading data from various sources and protocols. It provides the DataLoader and DataSelection protocols along with concrete implementations for different file formats.

Classes:

Name Description
DataLoader

Protocol for data loader factories.

DataSelection

Protocol for data subsets.

ParquetDataLoader

Loader for Parquet files.

PandasDataLoader

Loader for CSV files using Pandas.

__all__ = ['DataLoader', 'DataSelection', 'PandasDataLoader', 'ParquetDataLoader', 'dqml_dataloaders_registry'] module-attribute

dqml_dataloaders_registry = {'parquet': ParquetDataLoader, 'csv': PandasDataLoader} module-attribute

DataLoader

Bases: Protocol

Protocol for Data Loader factories.

A DataLoader is responsible for scanning a source (disk, DB, S3) and discovering available DataSelections based on its configuration.

Source code in packages/dqm-ml-job/src/dqm_ml_job/dataloaders/proto.py
@runtime_checkable
class DataLoader(Protocol):
    """
    Protocol for Data Loader factories.

    A DataLoader is responsible for scanning a source (disk, DB, S3) and
    discovering available DataSelections based on its configuration.
    """

    def get_selections(self) -> list[DataSelection]:
        """
        Discover and return the list of available selections for this loader.

        Returns:
            A list of initialized DataSelection instances.
        """

get_selections() -> list[DataSelection]

Discover and return the list of available selections for this loader.

Returns:

Type Description
list[DataSelection]

A list of initialized DataSelection instances.

Source code in packages/dqm-ml-job/src/dqm_ml_job/dataloaders/proto.py
def get_selections(self) -> list[DataSelection]:
    """
    Discover and return the list of available selections for this loader.

    Returns:
        A list of initialized DataSelection instances.
    """

DataSelection

Bases: Protocol

Protocol for a specific subset of data discovered by a DataLoader.

A DataSelection represents a concrete set of samples (e.g., a specific folder, a filtered view of a database, or a single file) and provides an iterator over data batches.

Source code in packages/dqm-ml-job/src/dqm_ml_job/dataloaders/proto.py
@runtime_checkable
class DataSelection(Protocol):
    """
    Protocol for a specific subset of data discovered by a DataLoader.

    A DataSelection represents a concrete set of samples (e.g., a
    specific folder, a filtered view of a database, or a single file)
    and provides an iterator over data batches.
    """

    name: str

    def bootstrap(self, columns_list: list[str]) -> None:
        """Perform initial setup for the selection before iteration.

        Args:
            columns_list: List of column names to load.
        """

    def get_nb_batches(self) -> int:
        """
        Return the estimated number of batches in this selection.

        Used primarily for progress bar estimation.
        """

    def __iter__(self) -> Any:
        """
        Iterate over the selection, yielding pyarrow.RecordBatch objects.
        """

name: str instance-attribute

__iter__() -> Any

Iterate over the selection, yielding pyarrow.RecordBatch objects.

Source code in packages/dqm-ml-job/src/dqm_ml_job/dataloaders/proto.py
def __iter__(self) -> Any:
    """
    Iterate over the selection, yielding pyarrow.RecordBatch objects.
    """

bootstrap(columns_list: list[str]) -> None

Perform initial setup for the selection before iteration.

Parameters:

Name Type Description Default
columns_list list[str]

List of column names to load.

required
Source code in packages/dqm-ml-job/src/dqm_ml_job/dataloaders/proto.py
def bootstrap(self, columns_list: list[str]) -> None:
    """Perform initial setup for the selection before iteration.

    Args:
        columns_list: List of column names to load.
    """

get_nb_batches() -> int

Return the estimated number of batches in this selection.

Used primarily for progress bar estimation.

Source code in packages/dqm-ml-job/src/dqm_ml_job/dataloaders/proto.py
def get_nb_batches(self) -> int:
    """
    Return the estimated number of batches in this selection.

    Used primarily for progress bar estimation.
    """

PandasDataLoader

Data loader for CSV files using Pandas.

This loader reads CSV files and provides DataSelections for processing by the DQM pipeline.

Attributes:

Name Type Description
type str

The loader type identifier ("csv").

Source code in packages/dqm-ml-job/src/dqm_ml_job/dataloaders/pandas.py
class PandasDataLoader:
    """Data loader for CSV files using Pandas.

    This loader reads CSV files and provides DataSelections for
    processing by the DQM pipeline.

    Attributes:
        type: The loader type identifier ("csv").
    """

    type: str = "csv"

    def __init__(self, name: str, config: dict[str, Any] | None = None):
        """Initialize the Pandas data loader.

        Args:
            name: Unique name for this loader instance.
            config: Configuration dictionary containing:
                - path: Path to CSV file (required)

        Raises:
            ValueError: If required config keys are missing.
        """
        if config is None:
            config = {}
        self.name = name
        self.path = config["path"]
        # Use SplitConfig model fields instead of hardcoded keys
        from dqm_ml_core.models.dataloaders import SplitConfig

        split = config.get("split")
        self.split = SplitConfig.model_validate(split) if split else None
        self.split_by = self.split.by if self.split else None
        self.split_values = self.split.values if self.split else None
        filters = config.get("filters")
        # transform the list of dict into a dict
        self.filters_dict = {}
        if filters is not None:
            for item in filters:
                column = item["column"]
                self.filters_dict[column] = item["values"]
        self.id_column = config.get("id_column")
        self.sample_path = config.get("sample_path", [])
        self.transforms = config.get("transform", [])

        # Storage filesystem configuration - only for S3 paths, not local paths
        self.filesystem = None
        storage_cfg = config.get("storage")
        if storage_cfg:
            # Use StorageConfig model to validate and access fields
            from dqm_ml_core.models.global_ import StorageConfig

            storage_config = StorageConfig.model_validate(storage_cfg)

            if storage_config.type == "s3":
                from dqm_ml_job.utils.s3 import get_s3_filesystem

                self.filesystem = get_s3_filesystem(storage_config)

    def get_selections(self) -> list[DataSelection]:
        """Create one or more PandasDataSelection instances based on split config.

        If split is configured, returns one selection per split value.
        Otherwise returns a single selection for the entire CSV file.

        Returns:
            A list of DataSelection instances.
        """
        if not self.split_by:
            return [
                PandasDataSelection(
                    name=self.name,
                    path=self.path,
                    sample_path=self.sample_path,
                    transforms=self.transforms,
                    filters_dict=self.filters_dict,
                )
            ]

        # Determine split values
        values = self.split_values
        if values is None:
            # Auto-discover unique values from the CSV
            df = pd.read_csv(self.path, sep=",", usecols=[self.split_by])
            values = [str(v) for v in df[self.split_by].unique() if v is not None]
        else:
            # Expand wildcard patterns in values against available data
            if any(has_pattern(v) for v in values):
                df = pd.read_csv(self.path, sep=",", usecols=[self.split_by])
                available = [str(v) for v in df[self.split_by].unique() if v is not None]
                values = resolve_patterns(values, available)

        # Apply split.exclude (including wildcard patterns)
        if self.split and self.split.exclude:
            values = resolve_include_exclude(None, self.split.exclude, values)

        # Create one selection per value
        selections: list[DataSelection] = []
        for val in values:
            selection_name = f"{self.name}_{val}"
            merged_filters = (self.filters_dict or {}).copy()
            merged_filters[self.split_by] = val
            selections.append(
                PandasDataSelection(
                    name=selection_name,
                    path=self.path,
                    sample_path=self.sample_path,
                    transforms=self.transforms,
                    filters_dict=merged_filters,
                )
            )
        return selections

filesystem = None instance-attribute

filters_dict = {} instance-attribute

id_column = config.get('id_column') instance-attribute

name = name instance-attribute

path = config['path'] instance-attribute

sample_path = config.get('sample_path', []) instance-attribute

split = SplitConfig.model_validate(split) if split else None instance-attribute

split_by = self.split.by if self.split else None instance-attribute

split_values = self.split.values if self.split else None instance-attribute

transforms = config.get('transform', []) instance-attribute

type: str = 'csv' class-attribute instance-attribute

__init__(name: str, config: dict[str, Any] | None = None)

Initialize the Pandas data loader.

Parameters:

Name Type Description Default
name str

Unique name for this loader instance.

required
config dict[str, Any] | None

Configuration dictionary containing: - path: Path to CSV file (required)

None

Raises:

Type Description
ValueError

If required config keys are missing.

Source code in packages/dqm-ml-job/src/dqm_ml_job/dataloaders/pandas.py
def __init__(self, name: str, config: dict[str, Any] | None = None):
    """Initialize the Pandas data loader.

    Args:
        name: Unique name for this loader instance.
        config: Configuration dictionary containing:
            - path: Path to CSV file (required)

    Raises:
        ValueError: If required config keys are missing.
    """
    if config is None:
        config = {}
    self.name = name
    self.path = config["path"]
    # Use SplitConfig model fields instead of hardcoded keys
    from dqm_ml_core.models.dataloaders import SplitConfig

    split = config.get("split")
    self.split = SplitConfig.model_validate(split) if split else None
    self.split_by = self.split.by if self.split else None
    self.split_values = self.split.values if self.split else None
    filters = config.get("filters")
    # transform the list of dict into a dict
    self.filters_dict = {}
    if filters is not None:
        for item in filters:
            column = item["column"]
            self.filters_dict[column] = item["values"]
    self.id_column = config.get("id_column")
    self.sample_path = config.get("sample_path", [])
    self.transforms = config.get("transform", [])

    # Storage filesystem configuration - only for S3 paths, not local paths
    self.filesystem = None
    storage_cfg = config.get("storage")
    if storage_cfg:
        # Use StorageConfig model to validate and access fields
        from dqm_ml_core.models.global_ import StorageConfig

        storage_config = StorageConfig.model_validate(storage_cfg)

        if storage_config.type == "s3":
            from dqm_ml_job.utils.s3 import get_s3_filesystem

            self.filesystem = get_s3_filesystem(storage_config)

get_selections() -> list[DataSelection]

Create one or more PandasDataSelection instances based on split config.

If split is configured, returns one selection per split value. Otherwise returns a single selection for the entire CSV file.

Returns:

Type Description
list[DataSelection]

A list of DataSelection instances.

Source code in packages/dqm-ml-job/src/dqm_ml_job/dataloaders/pandas.py
def get_selections(self) -> list[DataSelection]:
    """Create one or more PandasDataSelection instances based on split config.

    If split is configured, returns one selection per split value.
    Otherwise returns a single selection for the entire CSV file.

    Returns:
        A list of DataSelection instances.
    """
    if not self.split_by:
        return [
            PandasDataSelection(
                name=self.name,
                path=self.path,
                sample_path=self.sample_path,
                transforms=self.transforms,
                filters_dict=self.filters_dict,
            )
        ]

    # Determine split values
    values = self.split_values
    if values is None:
        # Auto-discover unique values from the CSV
        df = pd.read_csv(self.path, sep=",", usecols=[self.split_by])
        values = [str(v) for v in df[self.split_by].unique() if v is not None]
    else:
        # Expand wildcard patterns in values against available data
        if any(has_pattern(v) for v in values):
            df = pd.read_csv(self.path, sep=",", usecols=[self.split_by])
            available = [str(v) for v in df[self.split_by].unique() if v is not None]
            values = resolve_patterns(values, available)

    # Apply split.exclude (including wildcard patterns)
    if self.split and self.split.exclude:
        values = resolve_include_exclude(None, self.split.exclude, values)

    # Create one selection per value
    selections: list[DataSelection] = []
    for val in values:
        selection_name = f"{self.name}_{val}"
        merged_filters = (self.filters_dict or {}).copy()
        merged_filters[self.split_by] = val
        selections.append(
            PandasDataSelection(
                name=selection_name,
                path=self.path,
                sample_path=self.sample_path,
                transforms=self.transforms,
                filters_dict=merged_filters,
            )
        )
    return selections

ParquetDataLoader

Data loader for Parquet files that generates one or more DataSelections.

This loader can read from a single Parquet file or a directory of Parquet files, optionally splitting the data by a column value to create multiple selections.

Attributes:

Name Type Description
type str

The loader type identifier ("parquet").

filesystem

Optional PyArrow filesystem for reading.

Source code in packages/dqm-ml-job/src/dqm_ml_job/dataloaders/parquet.py
class ParquetDataLoader:
    """Data loader for Parquet files that generates one or more DataSelections.

    This loader can read from a single Parquet file or a directory of Parquet
    files, optionally splitting the data by a column value to create multiple
    selections.

    Attributes:
        type: The loader type identifier ("parquet").
        filesystem: Optional PyArrow filesystem for reading.
    """

    type: str = "parquet"

    def __init__(self, name: str, config: dict[str, Any] | None = None):
        """Initialize the Parquet data loader.

        Args:
            name: Unique name for this loader instance.
            config: Configuration dictionary containing:
                - path: Path to Parquet file or directory (required)
                - batch_size: Rows per batch (default: 100000)
                - threads: Number of threads (default: 4)
                - split_by: Column name to split selections by
                - split_values: Specific values to split on
                - filter.: list of filters
                - storage: Storage configuration (bool or dict)

        Raises:
            ValueError: If required config keys are missing.
        """
        if config is None:
            config = {}
        self.name = name
        self.config = config
        self.path: str = config["path"]
        self.batch_size = config.get("batch_size", 100_000)
        self.threads = config.get("threads", 4)
        # Use SplitConfig model fields instead of hardcoded keys
        from dqm_ml_core.models.dataloaders import SplitConfig

        split = config.get("split")
        self.split = SplitConfig.model_validate(split) if split else None
        self.split_by = self.split.by if self.split else None
        self.split_values = self.split.values if self.split else None
        filters = config.get("filters")
        # transform the list of dict into a dict
        self.filters_dict = {}
        if filters is not None:
            for item in filters:
                column = item["column"]
                self.filters_dict[column] = item["values"]
        logger.debug(f"[DEBUG] ParquetDataLoader.__init__: filters_dict = {self.filters_dict}")

        self.id_column = config.get("id_column")
        self.sample_path = config.get("sample_path", [])
        self.transforms = config.get("transform", [])

        # Storage filesystem configuration - only for S3 paths, not local paths
        self.filesystem = None
        storage_config = None
        storage_cfg = config.get("storage")
        if storage_cfg:
            # Use StorageConfig model to validate and access fields
            from dqm_ml_core.models.global_ import StorageConfig

            storage_config = StorageConfig.model_validate(storage_cfg)

            self.storage_config = storage_config

            if storage_config.type == "s3":
                from dqm_ml_job.utils.s3 import get_s3_filesystem

                self.filesystem = get_s3_filesystem(storage_config)

    def _resolve_selection_path(self) -> str:
        """Resolve the full path, prepending S3 bucket if applicable."""
        path = self.path
        if self.filesystem is not None and isinstance(self.filesystem, fs.S3FileSystem):
            bucket_name = self.storage_config.bucket if self.storage_config else os.getenv("S3_BUCKET_NAME", "")
            if bucket_name and not path.startswith(bucket_name + "/"):
                path = f"{bucket_name}/{path}"
        return path

    def get_selections(self) -> list[DataSelection]:
        """Create one or more ParquetDataSelection instances based on configuration.

        Returns:
            A list of DataSelection instances. If split_by is configured,
            returns one selection per unique value. Otherwise, returns a
            single selection for the entire dataset.
        """
        path = self._resolve_selection_path()

        if not self.split_by:
            # Single selection
            return [
                ParquetDataSelection(
                    name=self.name,
                    path=path,
                    batch_size=self.batch_size,
                    threads=self.threads,
                    filters_dict=self.filters_dict,
                    filesystem=self.filesystem,
                    sample_path=self.sample_path,
                    transforms=self.transforms,
                )
            ]

        # Splitting logic
        values = self.split_values
        if values is None:
            # Automatic discovery if split_values not provided
            logger.info(f"Discovering unique values for split_by='{self.split_by}' in {path}")
            table = pq.read_table(path, columns=[self.split_by], filesystem=self.filesystem)
            values = [str(v) for v in pc.unique(table.column(0)).to_pylist() if v is not None]
        else:
            # Expand wildcard patterns in values against available data
            if any(has_pattern(v) for v in values):
                logger.info(f"Expanding wildcard values for split_by='{self.split_by}' in {path}")
                table = pq.read_table(path, columns=[self.split_by], filesystem=self.filesystem)
                available = [str(v) for v in pc.unique(table.column(0)).to_pylist() if v is not None]
                values = resolve_patterns(values, available)

        # Apply split.exclude (including wildcard patterns)
        if self.split and self.split.exclude:
            values = resolve_include_exclude(None, self.split.exclude, values)

        selections: list[DataSelection] = []
        for val in values:
            selection_name = f"{self.name}_{val}"
            # Merge existing filters with the split filter
            merged_filters = (self.filters_dict or {}).copy()
            # TODO: filters shouldn't be on the same column than split by
            # raise error here ? or add this check in pydantic model ? is it possible ?
            merged_filters[self.split_by] = val

            selections.append(
                ParquetDataSelection(
                    name=selection_name,
                    path=path,
                    batch_size=self.batch_size,
                    threads=self.threads,
                    filters_dict=merged_filters,
                    filesystem=self.filesystem,
                    sample_path=self.sample_path,
                    transforms=self.transforms,
                )
            )
        return selections

batch_size = config.get('batch_size', 100000) instance-attribute

config = config instance-attribute

filesystem = None instance-attribute

filters_dict = {} instance-attribute

id_column = config.get('id_column') instance-attribute

name = name instance-attribute

path: str = config['path'] instance-attribute

sample_path = config.get('sample_path', []) instance-attribute

split = SplitConfig.model_validate(split) if split else None instance-attribute

split_by = self.split.by if self.split else None instance-attribute

split_values = self.split.values if self.split else None instance-attribute

storage_config = storage_config instance-attribute

threads = config.get('threads', 4) instance-attribute

transforms = config.get('transform', []) instance-attribute

type: str = 'parquet' class-attribute instance-attribute

__init__(name: str, config: dict[str, Any] | None = None)

Initialize the Parquet data loader.

Parameters:

Name Type Description Default
name str

Unique name for this loader instance.

required
config dict[str, Any] | None

Configuration dictionary containing: - path: Path to Parquet file or directory (required) - batch_size: Rows per batch (default: 100000) - threads: Number of threads (default: 4) - split_by: Column name to split selections by - split_values: Specific values to split on - filter.: list of filters - storage: Storage configuration (bool or dict)

None

Raises:

Type Description
ValueError

If required config keys are missing.

Source code in packages/dqm-ml-job/src/dqm_ml_job/dataloaders/parquet.py
def __init__(self, name: str, config: dict[str, Any] | None = None):
    """Initialize the Parquet data loader.

    Args:
        name: Unique name for this loader instance.
        config: Configuration dictionary containing:
            - path: Path to Parquet file or directory (required)
            - batch_size: Rows per batch (default: 100000)
            - threads: Number of threads (default: 4)
            - split_by: Column name to split selections by
            - split_values: Specific values to split on
            - filter.: list of filters
            - storage: Storage configuration (bool or dict)

    Raises:
        ValueError: If required config keys are missing.
    """
    if config is None:
        config = {}
    self.name = name
    self.config = config
    self.path: str = config["path"]
    self.batch_size = config.get("batch_size", 100_000)
    self.threads = config.get("threads", 4)
    # Use SplitConfig model fields instead of hardcoded keys
    from dqm_ml_core.models.dataloaders import SplitConfig

    split = config.get("split")
    self.split = SplitConfig.model_validate(split) if split else None
    self.split_by = self.split.by if self.split else None
    self.split_values = self.split.values if self.split else None
    filters = config.get("filters")
    # transform the list of dict into a dict
    self.filters_dict = {}
    if filters is not None:
        for item in filters:
            column = item["column"]
            self.filters_dict[column] = item["values"]
    logger.debug(f"[DEBUG] ParquetDataLoader.__init__: filters_dict = {self.filters_dict}")

    self.id_column = config.get("id_column")
    self.sample_path = config.get("sample_path", [])
    self.transforms = config.get("transform", [])

    # Storage filesystem configuration - only for S3 paths, not local paths
    self.filesystem = None
    storage_config = None
    storage_cfg = config.get("storage")
    if storage_cfg:
        # Use StorageConfig model to validate and access fields
        from dqm_ml_core.models.global_ import StorageConfig

        storage_config = StorageConfig.model_validate(storage_cfg)

        self.storage_config = storage_config

        if storage_config.type == "s3":
            from dqm_ml_job.utils.s3 import get_s3_filesystem

            self.filesystem = get_s3_filesystem(storage_config)

get_selections() -> list[DataSelection]

Create one or more ParquetDataSelection instances based on configuration.

Returns:

Type Description
list[DataSelection]

A list of DataSelection instances. If split_by is configured,

list[DataSelection]

returns one selection per unique value. Otherwise, returns a

list[DataSelection]

single selection for the entire dataset.

Source code in packages/dqm-ml-job/src/dqm_ml_job/dataloaders/parquet.py
def get_selections(self) -> list[DataSelection]:
    """Create one or more ParquetDataSelection instances based on configuration.

    Returns:
        A list of DataSelection instances. If split_by is configured,
        returns one selection per unique value. Otherwise, returns a
        single selection for the entire dataset.
    """
    path = self._resolve_selection_path()

    if not self.split_by:
        # Single selection
        return [
            ParquetDataSelection(
                name=self.name,
                path=path,
                batch_size=self.batch_size,
                threads=self.threads,
                filters_dict=self.filters_dict,
                filesystem=self.filesystem,
                sample_path=self.sample_path,
                transforms=self.transforms,
            )
        ]

    # Splitting logic
    values = self.split_values
    if values is None:
        # Automatic discovery if split_values not provided
        logger.info(f"Discovering unique values for split_by='{self.split_by}' in {path}")
        table = pq.read_table(path, columns=[self.split_by], filesystem=self.filesystem)
        values = [str(v) for v in pc.unique(table.column(0)).to_pylist() if v is not None]
    else:
        # Expand wildcard patterns in values against available data
        if any(has_pattern(v) for v in values):
            logger.info(f"Expanding wildcard values for split_by='{self.split_by}' in {path}")
            table = pq.read_table(path, columns=[self.split_by], filesystem=self.filesystem)
            available = [str(v) for v in pc.unique(table.column(0)).to_pylist() if v is not None]
            values = resolve_patterns(values, available)

    # Apply split.exclude (including wildcard patterns)
    if self.split and self.split.exclude:
        values = resolve_include_exclude(None, self.split.exclude, values)

    selections: list[DataSelection] = []
    for val in values:
        selection_name = f"{self.name}_{val}"
        # Merge existing filters with the split filter
        merged_filters = (self.filters_dict or {}).copy()
        # TODO: filters shouldn't be on the same column than split by
        # raise error here ? or add this check in pydantic model ? is it possible ?
        merged_filters[self.split_by] = val

        selections.append(
            ParquetDataSelection(
                name=selection_name,
                path=path,
                batch_size=self.batch_size,
                threads=self.threads,
                filters_dict=merged_filters,
                filesystem=self.filesystem,
                sample_path=self.sample_path,
                transforms=self.transforms,
            )
        )
    return selections