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 tooutput_storeunlessdry_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_runis 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:
-
tuple[NormalizeTask, ...]–Planned normalization tasks.
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
Nonekeyed by protocol id.
Returns:
-
NormalizeStageSpec–A new stage spec with matching protocol specs replaced.
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
1for sequential execution,-1for available CPUs minus one, or a positive count capped by task count. -
worker_polars_max_threads(int | None) –Polars threads per worker.
-1divides CPUs across workers,Noneleaves 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;
0disables size-based rolling. -
max_batch_rows(int | None) –Maximum rows processed in memory per bounded batch. Set to
Noneto 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:
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:
-
NormalizeProtocolSpec–Matching normalize protocol spec.
Raises:
-
ValueError–If no protocol spec matches.
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:
-
dict[MappingSpec, object]–Metadata values for manifest rows and footer resolution.
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
¶
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.