nemo_curator.stages.interleaved.io.writers.base

View as Markdown

Module Contents

Classes

NameDescription
BaseInterleavedWriterBase class for interleaved writers.

API

class nemo_curator.stages.interleaved.io.writers.base.BaseInterleavedWriter(
path: str,
file_extension: str,
write_kwargs: dict[str, typing.Any] = dict(),
materialize_on_write: bool = True,
name: str = 'base_interleaved_writer',
mode: typing.Literal['ignore', 'overwrite', 'append', 'error'] = 'ignore',
append_mode_implemented: bool = False,
on_materialize_error: typing.Literal['error', 'warn', 'drop_row', 'drop_sample'] = 'error',
schema: pyarrow.Schema | None = None,
schema_overrides: dict[str, pyarrow.DataType] | None = None
)
DataclassAbstract

Bases: ProcessingStage[InterleavedBatch, FileGroupTask]

Base class for interleaved writers.

Handles filesystem setup, deterministic file naming, optional binary materialization, schema alignment, and process() orchestration. Subclasses implement _write_dataframe for format-specific output.

If schema is set, every output table is aligned to it (missing columns become typed nulls, extra columns are dropped, types are reconciled). By default (schema=None) extra user columns are preserved and only reserved-column types are reconciled via reconcile_schema.

Use schema or schema_overrides only when strict column control is needed (e.g. to prevent heterogeneous-schema crashes).

append_mode_implemented
bool = False
file_extension
str
materialize_on_write
bool = True
mode
Literal['ignore', 'overwrite', 'append', 'error'] = 'ignore'
name
str = 'base_interleaved_writer'
on_materialize_error
Literal['error', 'warn', 'drop_row', 'drop_sample'] = 'error'
path
str
schema
Schema | None = None
schema_overrides
dict[str, DataType] | None = None
write_kwargs
dict[str, Any] = field(default_factory=dict)
nemo_curator.stages.interleaved.io.writers.base.BaseInterleavedWriter.__post_init__() -> None
nemo_curator.stages.interleaved.io.writers.base.BaseInterleavedWriter._align_output(
df: pandas.DataFrame
) -> pandas.DataFrame

Reconcile or align df to the declared schema.

nemo_curator.stages.interleaved.io.writers.base.BaseInterleavedWriter._materialize_dataframe(
task: nemo_curator.tasks.InterleavedBatch
) -> pandas.DataFrame
nemo_curator.stages.interleaved.io.writers.base.BaseInterleavedWriter._write_dataframe(
df: pandas.DataFrame,
file_path: str,
write_kwargs: dict[str, typing.Any]
) -> None

Format-specific DataFrame writer. Subclasses must implement this.

Subclasses that override write_data() or process() directly (e.g. writers that do not follow the one-file-per-task pattern) may override this method as a no-op instead.

nemo_curator.stages.interleaved.io.writers.base.BaseInterleavedWriter.inputs() -> tuple[list[str], list[str]]
nemo_curator.stages.interleaved.io.writers.base.BaseInterleavedWriter.outputs() -> tuple[list[str], list[str]]
nemo_curator.stages.interleaved.io.writers.base.BaseInterleavedWriter.process(
task: nemo_curator.tasks.InterleavedBatch
) -> nemo_curator.tasks.FileGroupTask
nemo_curator.stages.interleaved.io.writers.base.BaseInterleavedWriter.write_data(
task: nemo_curator.tasks.InterleavedBatch,
file_path: str
) -> None