Skip to main content
Pre-release

This release is a pre-release and may not be stable for production use.

DQM-ML Job

Orchestration engine for DQM-ML V2. Handles data loading, processing, and output writing.

Installation

pip install dqm-ml-job

Note: dqm-ml-job handles data loading and orchestration. To compute Metrics, you also need at least one of: dqm-ml-core, dqm-ml-images, or dqm-ml-pytorch (see Dependencies below).

Quick Start

Using Python

from dqm_ml_job.cli import execute

# Execute a data quality job from a YAML config
execute(["-p", "config.yaml"])

Using Python Module

python -m dqm_ml_job.cli -p config.yaml

Example config.yaml:

features:
  processors:
    - name: image_quality
      type: image_features
      columns:
        input: ["image_data"]
      features: [luminosity, contrast, blur, entropy]
      grayscale: true

metrics:
  processors:
    - name: completeness
      type: completeness
      columns:
        input: [col_a, col_b]

dataloaders:
  loaders:
    - name: my_data
      type: parquet
      path: data/train.parquet

Dependencies

DQM-ML is modular — dqm-ml-job provides the orchestration, but you need additional packages to compute actual metrics:

Interface Package Entry Point Group
Features dqm-ml-images dqm_ml.features
Features (Embeddings) dqm-ml-pytorch dqm_ml.features
Metrics dqm-ml-core dqm_ml.metrics
Gap dqm-ml-pytorch dqm_ml.gap
# For Metrics (Completeness, Representativeness, Diversity)
pip install dqm-ml-job dqm-ml-core

# For Visual Features
pip install dqm-ml-job dqm-ml-images

# For Image Embeddings + Domain Gap
pip install dqm-ml-job dqm-ml-pytorch

# All metrics
pip install dqm-ml-job dqm-ml-core dqm-ml-images dqm-ml-pytorch

Key Components

DatasetPipeline

The main orchestrator that:

  • Loads the configuration
  • Discovers plugins via entry points
  • Executes the streaming loop
  • Manages memory and I/O efficiency

Protocols

Protocol Description
DataLoader Factory for creating Data Selections (e.g., Parquet, CSV loaders)
DataSelection Represents a specific Data Selection and provides an iterator over Batches
OutputWriter Persists computed Features or Metrics to disk

Adding a Custom DataLoader

A DataLoader discovers available selections from a data source. Implement the DataLoader protocol and a companion DataSelection class.

Protocol Overview

DataLoader (dqm_ml_job/dataloaders/proto.py):

class DataLoader(Protocol):
    def get_selections(self) -> list[DataSelection]: ...

DataSelection (dqm_ml_job/dataloaders/proto.py):

class DataSelection(Protocol):
    name: str

    def bootstrap(self, columns_list: list[str]) -> None: ...
    def get_nb_batches(self) -> int: ...
    def __iter__(self) -> Any: ...

Example: JSON Lines Loader

Create dqm_ml_job/dataloaders/jsonl.py:

import json
import logging
from pathlib import Path
from typing import Any

import pyarrow as pa
from dqm_ml_job.dataloaders.proto import DataSelection

logger = logging.getLogger(__name__)


class JsonLinesDataSelection(DataSelection):
    def __init__(self, name: str, path: str, batch_size: int = 10_000):
        self.name = name
        self.path = path
        self.batch_size = batch_size
        self._lines: list[dict[str, Any]] | None = None

    def bootstrap(self, columns_list: list[str]) -> None:
        with open(self.path) as f:
            self._lines = [json.loads(line) for line in f]

    def get_nb_batches(self) -> int:
        n = len(self._lines) if self._lines else 0
        return (n // self.batch_size) + (1 if n % self.batch_size else 0)

    def __iter__(self) -> Any:
        if self._lines is None:
            return
        for i in range(0, len(self._lines), self.batch_size):
            yield pa.RecordBatch.from_pylist(self._lines[i:i + self.batch_size])


class JsonLinesDataLoader:
    type: str = "jsonl"

    def __init__(self, name: str, config: dict[str, Any] | None = None):
        config = config or {}
        self.name = name
        self.path = config["path"]
        self.batch_size = config.get("batch_size", 10_000)

    def get_selections(self) -> list[DataSelection]:
        return [JsonLinesDataSelection(self.name, self.path, self.batch_size)]

Registration

Add to the registry in dqm_ml_job/dataloaders/__init__.py:

from dqm_ml_job.dataloaders.jsonl import JsonLinesDataLoader

dqml_dataloaders_registry["jsonl"] = JsonLinesDataLoader

For external packages, register via entry points in pyproject.toml:

[project.entry-points."dqm_ml.dataloaders"]
jsonl = "my_package.dataloaders:JsonLinesDataLoader"

Adding a Custom OutputWriter

An OutputWriter persists computed features or metrics to a storage backend.

Protocol Overview

OutputWriter (dqm_ml_job/outputwriter/__init__.py):

class OutputWriter(Protocol):
    columns: list[str]
    name: str

    def write_metrics_dict(self, metrics_dict: dict[str, dict[str, Any]]) -> None: ...
    def write_table(self, name: str, table: Any, part_index: int | None = None) -> None: ...

Example: CSV Output Writer

Create dqm_ml_job/outputwriter/csv.py:

import csv
import logging
from pathlib import Path
from typing import Any

import pyarrow as pa
from dqm_ml_core.models.global_ import StorageConfig
from dqm_ml_core.models.outputs import ParquetOutputConfig

logger = logging.getLogger(__name__)


class CsvOutputWriter:
    def __init__(self, name: str, config: dict[str, Any] | None = None):
        cfg = ParquetOutputConfig.model_validate(config or {})  # reuse path_pattern, columns
        self.name = name
        self.path_pattern = cfg.path_pattern
        self.columns = list(cfg.columns)

    def write_metrics_dict(self, metrics_dict: dict[str, dict[str, Any]]) -> None:
        pass  # simplified — write all rows per selection

    def write_table(self, name: str, table: Any, part_index: int | None = None) -> None:
        if isinstance(table, dict):
            table = pa.table(table)
        if isinstance(table, pa.Table):
            df = table.to_pandas()
            path = Path(self.path_pattern.format(name=name))
            path.parent.mkdir(parents=True, exist_ok=True)
            df.to_csv(path, index=False)
            logger.info(f"Wrote CSV to {path}")

Registration

Add to dqm_ml_job/outputwriter/__init__.py:

from dqm_ml_job.outputwriter.csv import CsvOutputWriter

dqml_outputs_registry["csv"] = CsvOutputWriter

For external packages, use entry points:

[project.entry-points."dqm_ml.outputwriter"]
csv = "my_package.outputwriters:CsvOutputWriter"

Built-in Loaders

Loader Description
parquet Optimized loading using PyArrow
csv Flexible loading using Pandas

See Also

Release files for dqm-ml-job 2.0.0rc4

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for dqm-ml-job 2.0.0rc4
File Size Uploaded
dqm_ml_job-2.0.0rc4.tar.gz 27.5 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for dqm-ml-job 2.0.0rc4
File Interpreter ABI Platform
dqm_ml_job-2.0.0rc4-py3-none-any.whl Python 3 none any Details

Total release size: 59.0 kB

Release files / dqm_ml_job-2.0.0rc4.tar.gz

Download URL dqm_ml_job-2.0.0rc4.tar.gz
Size 27.5 kB
Tags Source
SHA-256 checksum
How to use checksums
b2f572e3511d29233b68a04940e2fdef6811c1a28d058c5d0fdcb70e3244b6fe
BLAKE2b-256 checksum
How to use checksums
b5880debd4a1d2cf1a27a06539650df52d9cb784c90adf80ed6f8a8f3bcb3099
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via uv/0.8.17

Release files / dqm_ml_job-2.0.0rc4-py3-none-any.whl

Download URL dqm_ml_job-2.0.0rc4-py3-none-any.whl
Size 31.5 kB
Tags Python 3
SHA-256 checksum
How to use checksums
5b02aa479153f86a178bf8b6e22ce332e9375739514a53ff47beafa07d2958b6
BLAKE2b-256 checksum
How to use checksums
516f80a6bcd4e59ae5b4635ad6c93364c535fcab22f5b0bbfc1326d81b7274ad
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via uv/0.8.17
Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page