Skip to content

Normalized

run_normalize

run_normalize(
    dataset_spec: DatasetSpec,
    input_store: DataProcessingStore,
    output_store: DataProcessingStore,
    config: NormalizeStageConfig,
    *,
    scratch_store: DataProcessingStore | None = None,
    tasks: tuple[NormalizeTask, ...] | None = None,
    dry_run: bool = False,
) -> None

Run ingested-to-normalized parquet processing for a dataset.

The stage reads the ingested manifest, plans protocol/group tasks, applies configured transforms, checks, and resampling, then writes normalized shards and a normalized manifest. dry_run=True validates tasks and checks without writing outputs.

Parameters:

  • dataset_spec (DatasetSpec) –

    Dataset configuration containing a normalized stage spec.

  • input_store (DataProcessingStore) –

    Store containing ingested manifest and shards.

  • output_store (DataProcessingStore) –

    Store receiving normalized shards and manifest.

  • config (NormalizeStageConfig) –

    Runtime, validation, and parquet-writing settings.

  • scratch_store (DataProcessingStore | None, default: None ) –

    Optional store for temporary task outputs. Defaults to output_store.

  • tasks (tuple[NormalizeTask, ...] | None, default: None ) –

    Optional task subset for retries or tests.

  • dry_run (bool, default: False ) –

    Validate and log check violations without writing output.

Returns:

  • None

    None. Outputs are written to output_store unless dry_run=True.

run_normalize_interactive

run_normalize_interactive(
    dataset_spec: DatasetSpec,
    input_store: DataProcessingStore,
    scratch_store: DataProcessingStore,
    config: NormalizeStageConfig,
    *,
    protocols: object = None,
    group_values: object = None,
    annotate: bool = True,
    source_run: InteractiveStageRun | None = None,
    normalize_spec: NormalizeStageSpec | None = None,
) -> InteractiveStageRun

Run selected normalization work into a scratch-backed interactive run.

protocols and group_values filter planned tasks. Passing an unresampled normalized source_run lets developers try different resampling settings without repeating the full normalize preparation path.

Parameters:

  • dataset_spec (DatasetSpec) –

    Dataset configuration containing a normalized stage spec.

  • input_store (DataProcessingStore) –

    Store containing ingested manifest and shards, unless source_run is provided.

  • scratch_store (DataProcessingStore) –

    Store receiving temporary interactive outputs.

  • config (NormalizeStageConfig) –

    Normalize runtime, validation, and parquet-writing settings.

  • protocols (object, default: None ) –

    Optional protocol selector, such as one id or a list of ids.

  • group_values (object, default: None ) –

    Optional task group selector or list of selectors.

  • annotate (bool, default: True ) –

    Whether public annotation columns are included in outputs.

  • source_run (InteractiveStageRun | None, default: None ) –

    Optional prior unresampled normalized run used as the source for resampling experiments.

  • normalize_spec (NormalizeStageSpec | None, default: None ) –

    Optional normalization spec override.

Returns:

  • InteractiveStageRun

    Interactive run handle for scanning outputs and cleaning scratch files.

plan_normalize_tasks

plan_normalize_tasks(
    dataset_spec: DatasetSpec,
    input_store: DataProcessingStore,
    normalize_spec: NormalizeStageSpec,
    *,
    protocols: object = None,
    group_values: object = None,
) -> tuple[NormalizeTask, ...]

Plan normalization tasks from the ingested manifest.

The ingested manifest is optionally filtered by protocol and group values, then rows are grouped by each protocol's task key. Each output task merges raw source paths and parquet segment references for one group.

Parameters:

  • dataset_spec (DatasetSpec) –

    Dataset configuration.

  • input_store (DataProcessingStore) –

    Store containing the ingested manifest.

  • normalize_spec (NormalizeStageSpec) –

    Normalization configuration.

  • protocols (object, default: None ) –

    Optional protocol selector.

  • group_values (object, default: None ) –

    Optional task group selector or list of selectors.

Returns:

Raises:

  • ValueError

    If required manifest metadata or group selectors are invalid.

normalize_spec_with_resampling

normalize_spec_with_resampling(
    normalize_spec: NormalizeStageSpec,
    resampling_by_protocol: dict[
        DatasetProtocolId, ResamplingSpec | None
    ],
) -> NormalizeStageSpec

Return a copy of a normalize spec with protocol resampling overrides.

Parameters:

  • normalize_spec (NormalizeStageSpec) –

    Base normalization spec.

  • resampling_by_protocol (dict[DatasetProtocolId, ResamplingSpec | None]) –

    Resampling spec or None keyed by protocol id.

Returns:

Examples:

>>> from batgrad.contracts.mapping import DatasetProtocolId
>>> from batgrad.contracts.protocols import BatteryProtocols
>>> from batgrad.data.transforms.resampling import MinMaxLTTBResamplingSpec
>>> base = NormalizeStageSpec(
...     protocol_specs=(NormalizeProtocolSpec(protocol=BatteryProtocols.cyc),),
... )
>>> updated = normalize_spec_with_resampling(
...     base,
...     {
...         DatasetProtocolId.cycling: MinMaxLTTBResamplingSpec(
...             x_col=BaseColumns.time,
...             y_col=BaseColumns.volt,
...             points_ratio=0.1,
...         )
...     },
... )

NormalizeStageConfig dataclass

NormalizeStageConfig(
    n_jobs: int = 1,
    worker_polars_max_threads: int | None = -1,
    chunk_rows: int = 500000,
    compression: str = "zstd",
    use_content_defined_chunking: bool = True,
    row_group_size: int = 262144,
    max_shard_size_bytes: int = 2 * 1024 * 1024 * 1024,
    max_batch_rows: int | None = 500000,
    apply_resampling: bool = True,
    apply_physics_compensation: bool = True,
)

Runtime, validation, and parquet-writing settings for normalization.

Tasks are collected in memory when max_batch_rows is None or the task row count fits within that limit. Larger tasks use bounded chunk processing; if a protocol also has resampling, its resampling spec must support bounded execution.

Attributes:

  • n_jobs (int) –

    Worker count. Use 1 for sequential execution, -1 for available CPUs minus one, or a positive count capped by task count.

  • worker_polars_max_threads (int | None) –

    Polars threads per worker. -1 divides CPUs across workers, None leaves Polars unrestricted, and a positive value sets an exact thread count.

  • chunk_rows (int) –

    Output chunk size written to final shards.

  • compression (str) –

    Parquet compression codec.

  • use_content_defined_chunking (bool) –

    Whether table writers may use content defined chunking.

  • row_group_size (int) –

    Parquet row group size for written files.

  • max_shard_size_bytes (int) –

    Roll a protocol shard after this approximate size; 0 disables size-based rolling.

  • max_batch_rows (int | None) –

    Maximum rows processed in memory per bounded batch. Set to None to collect each task fully.

  • apply_resampling (bool) –

    Whether protocol resampling specs are applied.

  • apply_physics_compensation (bool) –

    Whether bounded MinMaxLTTB recomputes dt, rebuilt time, and averaged current/C-rate when possible.

Examples:

>>> NormalizeStageConfig(n_jobs=-1, max_batch_rows=500_000)
NormalizeStageConfig(...)

NormalizeStageSpec dataclass

NormalizeStageSpec(
    metadata: StageLayout = NORMALIZE_STAGE_METADATA,
    protocol_specs: tuple[
        NormalizeProtocolSpec, ...
    ] = tuple(),
    time_convention: str = "start_of_interval",
)

Dataset-level normalization configuration.

The stage expands manifest metadata with raw paths, ingested parquet segment references, resampling method/arguments, time convention, and protocol group keys. time_convention is also written to generated parquet footers.

Attributes:

  • metadata (StageLayout) –

    Base normalized-stage metadata layout.

  • protocol_specs (tuple[NormalizeProtocolSpec, ...]) –

    Protocol normalization recipes.

  • time_convention (str) –

    Label written to manifest/footer metadata describing how task time values are interpreted.

Examples:

>>> from batgrad.contracts.protocols import BatteryProtocols
>>> NormalizeStageSpec(
...     protocol_specs=(NormalizeProtocolSpec(protocol=BatteryProtocols.cyc),),
...     time_convention="start_of_interval",
... )
NormalizeStageSpec(...)

Methods:

  • protocol_spec

    Return the protocol spec matching a protocol id or value.

  • output_columns

    Canonical output columns for a protocol.

  • required_input_columns

    Ingested columns required for a protocol's normalize tasks.

  • task_metadata

    Metadata written for one normalized task's manifest rows and footers.

  • output_spec

    Build the output writer configuration for normalized shards.

protocol_spec

protocol_spec(protocol: object) -> NormalizeProtocolSpec

Return the protocol spec matching a protocol id or value.

Parameters:

  • protocol (object) –

    Protocol id, enum value, or string-like value.

Returns:

Raises:

output_columns

output_columns(protocol: object) -> tuple[MappingSpec, ...]

Canonical output columns for a protocol.

Parameters:

  • protocol (object) –

    Protocol id, enum value, or string-like value.

Returns:

  • tuple[MappingSpec, ...]

    Columns written to normalized parquet before annotations.

required_input_columns

required_input_columns(
    protocol: object,
) -> tuple[MappingSpec, ...]

Ingested columns required for a protocol's normalize tasks.

Parameters:

  • protocol (object) –

    Protocol id, enum value, or string-like value.

Returns:

  • tuple[MappingSpec, ...]

    Ingested columns needed to prepare and validate the protocol.

task_metadata

task_metadata(
    dataset_id: str, task: NormalizeTask
) -> dict[MappingSpec, object]

Metadata written for one normalized task's manifest rows and footers.

Parameters:

  • dataset_id (str) –

    Dataset id for generated output metadata.

  • task (NormalizeTask) –

    Planned normalization task.

Returns:

output_spec

output_spec(
    dataset_spec: DatasetSpec,
    *,
    output_root: str | None = None,
    manifest_path: str | None = None,
) -> StageOutputSpec

Build the output writer configuration for normalized shards.

Parameters:

  • dataset_spec (DatasetSpec) –

    Dataset storage configuration.

  • output_root (str | None, default: None ) –

    Optional output root override for interactive runs.

  • manifest_path (str | None, default: None ) –

    Optional manifest path override for interactive runs.

Returns:

  • StageOutputSpec

    Stage writer configuration for protocol-sharded normalized output.

NormalizeProtocolSpec dataclass

NormalizeProtocolSpec(
    protocol: BatteryProtocolSpec,
    order_by: tuple[MappingSpec, ...] = (),
    columns: tuple[MappingSpec, ...] = (),
    constant_columns: dict[MappingSpec, object] = dict(),
    transforms: tuple[
        NormalizeTransformSpec, ...
    ] = tuple(),
    checks: tuple[CheckSpec, ...] = tuple(),
    resampling: ResamplingSpec | None = None,
)

Normalization recipe for one protocol.

Required input columns are inferred from order columns, transform inputs, output columns not produced by transforms/constants/checks, check inputs, and resampling inputs. When multiple aliases for a requested column are available in ingested shards, normalization coalesces them into the canonical output column.

Resampling runs after transforms and checks. Large tasks only work with a resampler that implements bounded execution; otherwise keep max_batch_rows unset or high enough for full-task processing.

Attributes:

  • protocol (BatteryProtocolSpec) –

    Shared protocol definition, including task grouping metadata.

  • order_by (tuple[MappingSpec, ...]) –

    Columns used to sort each task. Defaults to protocol.axis_col.

  • columns (tuple[MappingSpec, ...]) –

    Canonical output columns to write.

  • constant_columns (dict[MappingSpec, object]) –

    Fixed columns added to every output row, such as ambient temperature when it is known from dataset metadata instead of raw measurements.

  • transforms (tuple[NormalizeTransformSpec, ...]) –

    Column derivations run before checks.

  • checks (tuple[CheckSpec, ...]) –

    Validation or annotation checks run after transforms.

  • resampling (ResamplingSpec | None) –

    Optional row reduction or interpolation step.

Examples:

>>> from batgrad.contracts.protocols import BatteryProtocols
>>> from batgrad.data.transforms.checks import MissingCheckSpec, TimeCheckSpec
>>> NormalizeProtocolSpec(
...     protocol=BatteryProtocols.cyc,
...     columns=(
...         BaseColumns.time,
...         BaseColumns.volt,
...         BaseColumns.curr,
...         BaseColumns.amb_temp,
...     ),
...     constant_columns={BaseColumns.amb_temp: 20.0},
...     checks=(MissingCheckSpec(), TimeCheckSpec(BaseColumns.time, BaseColumns.dt)),
... )
NormalizeProtocolSpec(...)

protocol_id property

protocol_id: DatasetProtocolId

Canonical protocol id for this normalize spec.

group_by property

group_by: tuple[MappingSpec, ...]

Task grouping columns inherited from protocol metadata.

output_columns property

output_columns: tuple[MappingSpec, ...]

Canonical columns written for this protocol before annotations.

required_input_columns property

required_input_columns: tuple[MappingSpec, ...]

Ingested columns needed to prepare this protocol's normalize tasks.

NormalizeTask dataclass

NormalizeTask(
    task_id: str,
    protocol_id: str,
    group_values: dict[MappingSpec, object],
    raw_paths: tuple[str, ...],
    parquet_segments: tuple[ParquetSegment, ...],
    row_count: int,
)

One protocol/group normalization unit planned from the ingested manifest.

A task contains the ingested parquet segments for one protocol task key, such as one cell/cycle pair, plus source-path and row-count metadata used for the normalized manifest.

Attributes:

  • task_id (str) –

    Stable id used in logs and temporary output paths.

  • protocol_id (str) –

    Protocol represented by the task.

  • group_values (dict[MappingSpec, object]) –

    Task-key metadata values, such as cell and cycle.

  • raw_paths (tuple[str, ...]) –

    Raw source paths represented by the task.

  • parquet_segments (tuple[ParquetSegment, ...]) –

    Ingested parquet segments to read.

  • row_count (int) –

    Total input row count across segments.