Skip to content

Local Table Writer

LocalTableWriter

LocalTableWriter(
    path: Path,
    schema: Schema,
    compression: str,
    *,
    use_content_defined_chunking: bool,
)

Write one parquet table from multiple in-memory chunks.

LocalTableWriter owns an open pyarrow.parquet.ParquetWriter for a single output file. Each call to write_table appends one polars.DataFrame as another parquet table chunk. The file is created eagerly, parent directories are created as needed, and existing output files are never overwritten.

Use this writer when a processing stage produces a table incrementally and should not keep the complete result in memory. Call close when all chunks have been written so parquet metadata can be attached and the file can be flushed.

Examples:

>>> import polars as pl
>>> import pyarrow as pa
>>> from pathlib import Path
>>> from batgrad.storage.local import LocalDataProcessingStore
>>> store = LocalDataProcessingStore(Path("/data/batgrad"), create=True)
>>> schema = pa.schema([pa.field("voltage", pa.float64())])
>>> writer = store.open_table_writer("normalized/cell.parquet", schema, "zstd")
>>> writer.write_table(pl.DataFrame({"voltage": [3.7, 3.8]}))
>>> writer.write_table(pl.DataFrame({"voltage": [3.9]}))
>>> writer.close({"stage": "normalized"})

Parameters:

  • path (Path) –

    Absolute output path for the parquet file.

  • schema (Schema) –

    Arrow schema used for all chunks written by this writer.

  • compression (str) –

    Parquet compression codec passed to PyArrow.

  • use_content_defined_chunking (bool) –

    Whether PyArrow should use content-defined chunking when writing parquet data.

Raises:

Source code in batgrad/storage/local.py
def __init__(
    self,
    path: Path,
    schema: pa.Schema,
    compression: str,
    *,
    use_content_defined_chunking: bool,
) -> None:
    """Create a parquet writer for a new local file.

    Args:
        path: Absolute output path for the parquet file.
        schema: Arrow schema used for all chunks written by this writer.
        compression: Parquet compression codec passed to PyArrow.
        use_content_defined_chunking: Whether PyArrow should use content-defined
            chunking when writing parquet data.

    Raises:
        FileExistsError: If `path` already exists.
    """
    path.parent.mkdir(parents=True, exist_ok=True)
    if path.exists():
        raise FileExistsError(f"File exists: {path}")
    self._writer = pq.ParquetWriter(
        path,
        schema,
        compression=compression,
        use_content_defined_chunking=use_content_defined_chunking,
    )

write_table

write_table(
    data: DataFrame, row_group_size: int | None = None
) -> None

Append one dataframe chunk to the open parquet file.

Parameters:

  • data (DataFrame) –

    Dataframe chunk to append. Its schema must match the writer schema.

  • row_group_size (int | None, default: None ) –

    Optional parquet row-group size for this chunk.

Source code in batgrad/storage/local.py
def write_table(self, data: pl.DataFrame, row_group_size: int | None = None) -> None:
    """Append one dataframe chunk to the open parquet file.

    Args:
        data: Dataframe chunk to append. Its schema must match the writer schema.
        row_group_size: Optional parquet row-group size for this chunk.
    """
    self._writer.write_table(data.to_arrow(), row_group_size=row_group_size)

close

close(metadata: dict[str, str] | None = None) -> None

Attach optional footer metadata and close the writer.

Parameters:

  • metadata (dict[str, str] | None, default: None ) –

    Optional parquet key-value footer metadata to write before closing.

Source code in batgrad/storage/local.py
def close(self, metadata: dict[str, str] | None = None) -> None:
    """Attach optional footer metadata and close the writer.

    Args:
        metadata: Optional parquet key-value footer metadata to write before closing.
    """
    if metadata is not None:
        self._writer.add_key_value_metadata(metadata)
    self._writer.close()