Skip to content

Ingested

run_ingest

run_ingest(
    adapter: RawDatasetAdapter,
    input_store: DataProcessingStore,
    output_store: DataProcessingStore,
    config: IngestStageConfig,
    *,
    scratch_store: DataProcessingStore | None = None,
    tasks: tuple[IngestTask, ...] | None = None,
) -> None

Run raw-to-ingested parquet processing for a dataset adapter.

The stage plans tasks with the adapter unless tasks is supplied, writes temporary task parquet to the scratch store, then appends results into protocol-sharded output files and an ingested manifest.

Parameters:

  • adapter (RawDatasetAdapter) –

    Dataset-specific raw loader.

  • input_store (DataProcessingStore) –

    Store containing raw source files.

  • output_store (DataProcessingStore) –

    Store receiving ingested shards and manifest.

  • config (IngestStageConfig) –

    Runtime and parquet-writing settings.

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

    Optional store for temporary task outputs. Defaults to output_store.

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

    Optional task subset for retries or tests.

Returns:

  • None

    None. Outputs are written to output_store.

IngestStageConfig dataclass

IngestStageConfig(
    n_jobs: int = 1,
    worker_polars_max_threads: int | None = -1,
    chunk_rows: int = 256000,
    compression: str = "zstd",
    use_content_defined_chunking: bool = True,
    row_group_size: int = 262144,
    max_shard_size_bytes: int = 2 * 1024 * 1024 * 1024,
)

Runtime and parquet-writing settings for the ingest stage.

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

    Rows read from temporary task outputs at a time.

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

Examples:

>>> IngestStageConfig(n_jobs=-1, worker_polars_max_threads=-1)
IngestStageConfig(...)

IngestStageSpec dataclass

IngestStageSpec(
    metadata: StageLayout,
    included_file_patterns: tuple[str, ...],
    excluded_file_patterns: tuple[str, ...] = tuple(),
    protocol_specs: tuple[
        IngestProtocolSpec, ...
    ] = tuple(),
)

Dataset-level ingest configuration.

Include/exclude patterns are used by raw adapters when planning source-file tasks. Stage metadata is expanded with each protocol's task keys and manifest extras before the output manifest is written.

Attributes:

  • metadata (StageLayout) –

    Base stage metadata layout.

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

    Raw files an adapter should consider.

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

    Raw files an adapter should ignore.

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

    Protocol-specific raw mappings.

Examples:

>>> from batgrad.contracts.metadata import INGEST_STAGE_METADATA
>>> from batgrad.contracts.protocols import BatteryProtocols
>>> protocol_spec = IngestProtocolSpec(
...     protocol=BatteryProtocols.cyc,
...     columns=(BaseColumns.time, BaseColumns.curr, BaseColumns.volt),
... )
>>> spec = IngestStageSpec(
...     metadata=INGEST_STAGE_METADATA,
...     included_file_patterns=("*.xlsx",),
...     protocol_specs=(protocol_spec,),
... )
>>> spec.is_included_file("raw/cell.xlsx")
True

Methods:

  • protocol_spec

    Return the protocol spec matching a protocol id or value.

  • output_columns

    Canonical output columns for a protocol.

  • required_metadata

    Batch metadata required for a protocol.

  • manifest_columns

    Protocol metadata columns written to the ingested manifest.

  • is_included_file

    Return whether a raw file path passes include/exclude patterns.

  • output_spec

    Build the output writer configuration for ingested shards.

protocol_spec

protocol_spec(protocol: object) -> IngestProtocolSpec

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 ingested parquet for the protocol.

required_metadata

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

Batch metadata required for a protocol.

Parameters:

  • protocol (object) –

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

Returns:

manifest_columns

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

Protocol metadata columns written to the ingested manifest.

Parameters:

  • protocol (object) –

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

Returns:

  • tuple[MappingSpec, ...]

    Metadata columns used by this protocol's manifest rows.

is_included_file

is_included_file(path: str) -> bool

Return whether a raw file path passes include/exclude patterns.

Parameters:

  • path (str) –

    Raw source path relative to the store.

Returns:

  • bool

    True when path matches include patterns and no exclude pattern.

output_spec

output_spec(dataset_spec: DatasetSpec) -> StageOutputSpec

Build the output writer configuration for ingested shards.

Parameters:

  • dataset_spec (DatasetSpec) –

    Dataset storage configuration.

Returns:

  • StageOutputSpec

    Stage writer configuration for protocol-sharded ingested output.

IngestProtocolSpec dataclass

IngestProtocolSpec(
    protocol: BatteryProtocolSpec,
    columns: tuple[MappingSpec, ...],
    metadata: ProtocolMetadata | None = None,
    dropped_columns: tuple[MappingSpec, ...] = (),
    flip_current_sign: bool = False,
)

Raw-to-ingested mapping for one protocol.

Dataset adapters yield raw batches for a protocol. The ingest stage aligns each batch to columns using MappingSpec aliases and parsers, validates protocol metadata, optionally flips current sign, and writes ingested parquet shards grouped by protocol.

Attributes:

  • protocol (BatteryProtocolSpec) –

    Shared protocol definition.

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

    Canonical output columns expected in adapter output.

  • metadata (ProtocolMetadata | None) –

    Optional protocol metadata override for this dataset.

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

    Declared raw columns that may appear but are omitted.

  • flip_current_sign (bool) –

    Whether to multiply canonical current by -1.

Examples:

>>> from batgrad.contracts.protocols import BatteryProtocols
>>> IngestProtocolSpec(
...     protocol=BatteryProtocols.cyc,
...     columns=(BaseColumns.time, BaseColumns.curr, BaseColumns.volt),
...     flip_current_sign=True,
... )
IngestProtocolSpec(...)

protocol_id property

protocol_id: DatasetProtocolId

Canonical protocol id for this ingest spec.

Returns:

  • DatasetProtocolId

    Shared protocol id from protocol.

protocol_metadata property

protocol_metadata: ProtocolMetadata

Protocol metadata override, or the shared protocol metadata.

Returns:

  • ProtocolMetadata

    Metadata used to validate adapter batch metadata and expand manifests.

output_columns property

output_columns: tuple[MappingSpec, ...]

Canonical columns written for this protocol.

Returns:

manifest_columns property

manifest_columns: tuple[MappingSpec, ...]

Protocol task and manifest metadata columns expected from batches.

Returns:

  • MappingSpec

    Task-key and protocol manifest metadata columns. Adapter batches may

  • ...

    carry these as metadata instead of data columns.

required_metadata property

required_metadata: tuple[MappingSpec, ...]

Metadata columns that each adapter batch must provide.

Returns:

  • tuple[MappingSpec, ...]

    Task-key and required protocol manifest metadata columns.

RawDatasetAdapter

Bases: Protocol

Interface implemented by dataset-specific raw loaders.

Implementations discover raw source files in plan_raw_tasks and load them in load_raw_task. The loader is responsible for inferring protocol and task metadata; the ingest stage handles column alignment, validation, and parquet writing.

Methods:

  • load_raw_task

    Load a planned raw task and yield protocol batches.

  • plan_raw_tasks

    Discover raw source work units for the ingest stage.

load_raw_task

load_raw_task(
    task: IngestTask,
    input_store: DataProcessingStore,
    raw_spec: IngestStageSpec,
) -> Iterator[IngestBatch]

Load a planned raw task and yield protocol batches.

Parameters:

  • task (IngestTask) –

    Planned raw task.

  • input_store (DataProcessingStore) –

    Store containing raw files.

  • raw_spec (IngestStageSpec) –

    Dataset ingest configuration.

Returns:

plan_raw_tasks

plan_raw_tasks(
    input_store: DataProcessingStore,
    raw_spec: IngestStageSpec,
) -> tuple[IngestTask, ...]

Discover raw source work units for the ingest stage.

Parameters:

  • input_store (DataProcessingStore) –

    Store containing raw files.

  • raw_spec (IngestStageSpec) –

    Dataset ingest configuration.

Returns:

  • IngestTask

    Planned ingest tasks. Implementations usually combine

  • ...

    input_store.list_files with raw_spec.is_included_file.

IngestTask dataclass

IngestTask(task_id: str, source_paths: tuple[str, ...])

One adapter-planned raw ingest unit, usually one source file.

Attributes:

  • task_id (str) –

    Stable id used in logs and temporary output paths.

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

    Raw source paths consumed by this task.

IngestBatch dataclass

IngestBatch(
    data: DataFrame | LazyFrame,
    protocol_id: DatasetProtocolId,
    source_paths: tuple[str, ...],
    metadata: dict[MappingSpec, object],
)

Raw data and metadata yielded by a raw dataset adapter.

metadata must include BaseColumns.proto and all columns required by the selected protocol metadata. The ingest stage writes these values to manifest rows and parquet footer metadata where declared by the stage layout.

Attributes:

  • data (DataFrame | LazyFrame) –

    Raw frame yielded by the adapter.

  • protocol_id (DatasetProtocolId) –

    Protocol represented by this batch.

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

    Raw source paths that produced the batch.

  • metadata (dict[MappingSpec, object]) –

    Task and protocol metadata for manifest/footer writing.

prepare_raw_batch

prepare_raw_batch(
    batch: IngestBatch, raw_spec: IngestStageSpec
) -> tuple[DataFrame, tuple[str, ...]]

Validate, align, and materialize one adapter batch.

Parameters:

Returns:

  • DataFrame

    Materialized canonical frame and non-fatal warnings, such as declared

  • tuple[str, ...]

    dropped columns or duplicate canonical output names.

align_to_protocol_spec

align_to_protocol_spec(
    data: DataFrame | LazyFrame,
    protocol_spec: IngestProtocolSpec,
    source_paths: tuple[str, ...],
) -> tuple[DataFrame | LazyFrame, tuple[str, ...]]

Select raw columns into canonical protocol columns.

Source columns are matched through MappingSpec aliases. If a mapping has a parser, the parser builds the output expression; otherwise the source column is cast to the mapping dtype. Unknown raw columns are errors unless declared in protocol columns, manifest metadata columns, or dropped_columns. Missing declared columns are errors; adapters should add null optional columns before yielding batches. Duplicate canonical mappings are kept with suffixed output names such as "column 1" and returned as warnings.

Returns:

  • tuple[DataFrame | LazyFrame, tuple[str, ...]]

    Canonically selected frame and non-fatal warnings.

Raises:

  • ValueError

    If declared columns are missing, unknown raw columns remain, or protocol one_of_col_groups are not satisfied.