nemo_curator.stages.interleaved.io.readers.base

View as Markdown

Module Contents

Classes

NameDescription
BaseInterleavedReaderBase contract for interleaved readers.

API

class nemo_curator.stages.interleaved.io.readers.base.BaseInterleavedReader(
read_kwargs: dict[str, typing.Any] = dict(),
schema: pyarrow.Schema | None = None,
schema_overrides: dict[str, pyarrow.DataType] | None = None,
name: str = 'base_interleaved_reader'
)
Dataclass

Bases: ProcessingStage[FileGroupTask, InterleavedBatch]

Base contract for interleaved readers.

By default (schema=None) user-added passthrough columns are preserved and only reserved-column types are reconciled via reconcile_schema.

If schema is set explicitly, every output table is strictly aligned to it (missing columns become typed nulls, extra columns are dropped).

Use schema_overrides to add or override individual field types relative to INTERLEAVED_SCHEMA while keeping strict alignment:

.. code-block:: python

reader = InterleavedParquetReader( “data.parquet”, schema_overrides={“url”: pa.string(), “timestamp”: pa.int64()}, )

name
str = 'base_interleaved_reader'
read_kwargs
dict[str, Any] = field(default_factory=dict)
schema
Schema | None = None
schema_overrides
dict[str, DataType] | None = None
nemo_curator.stages.interleaved.io.readers.base.BaseInterleavedReader.__post_init__() -> None
nemo_curator.stages.interleaved.io.readers.base.BaseInterleavedReader._align_output(
table: pyarrow.Table
) -> pyarrow.Table

Reconcile or align table to the declared schema.

nemo_curator.stages.interleaved.io.readers.base.BaseInterleavedReader._source_files_for_split(
split: pyarrow.Table,
idx: int,
sample_id_to_path: dict[str, str],
all_paths: list[str]
) -> list[str]
staticmethod

Return source_files for one split, annotated with the split index for lineage tracking.

The ::split_NNN suffix is appended so that downstream consumers can correlate each output batch back to the exact split of its source file(s), even when a single source file is split into multiple batches by max_batch_bytes.

nemo_curator.stages.interleaved.io.readers.base.BaseInterleavedReader.inputs() -> tuple[list[str], list[str]]
nemo_curator.stages.interleaved.io.readers.base.BaseInterleavedReader.outputs() -> tuple[list[str], list[str]]