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 tooutput_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
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) –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;
0disables size-based rolling.
Examples:
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:
-
IngestProtocolSpec–Matching ingest 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 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:
-
tuple[MappingSpec, ...]–Metadata columns each adapter batch must provide.
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
¶
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
¶
Canonical protocol id for this ingest spec.
Returns:
-
DatasetProtocolId–Shared protocol id from
protocol.
protocol_metadata
property
¶
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:
-
tuple[MappingSpec, ...]–Ingested parquet columns after raw alignment.
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:
-
Iterator[IngestBatch]–Iterator of raw batches. A task may yield multiple batches when one
-
Iterator[IngestBatch]–source contains multiple protocols or logical groups.
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_fileswithraw_spec.is_included_file.
IngestTask
dataclass
¶
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:
-
batch(IngestBatch) –Raw adapter batch.
-
raw_spec(IngestStageSpec) –Dataset ingest configuration.
Returns:
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:
Raises:
-
ValueError–If declared columns are missing, unknown raw columns remain, or protocol
one_of_col_groupsare not satisfied.