Skip to content

fix(streaming-write): use rolling ParquetWriter + OutputStream.tell() for spec-correct file sizes and bounded memory #3388

Description

@paultmathew

Background

PR #3335 added pa.RecordBatchReader as a valid input to Table.append/Table.overwrite using a buffered bin-pack approach (bin_pack_record_batches). That implementation has two acknowledged caveats called out in its docstrings:

  1. Memory bound: peak memory is N_workers × write.target-file-size-bytes (~4 GiB at defaults) — better than materialising everything, but not constant.
  2. Byte semantics: write.target-file-size-bytes is interpreted as uncompressed in-memory Arrow bytes, not on-disk compressed Parquet bytes. Resulting files are typically 3–10× smaller than the property suggests — diverging from the Java/Spark/Flink writers.

Proposed fix

Replace the bin-pack approach with a rolling pq.ParquetWriter driven by OutputStream.tell() (added in #2998 specifically for this purpose):

with output_file.create(overwrite=True) as fos:
    with pq.ParquetWriter(fos, schema=..., ...) as writer:
        writer.write_batch(first_batch)
        while fos.tell() < target_file_size:   # ← compressed on-disk bytes
            batch = next(batches)
            writer.write_batch(batch)

This delivers:

  • Spec-correct file sizes: tell() reports compressed on-disk bytes, so write.target-file-size-bytes finally means what the Iceberg spec intends — consistent with the Java/Spark/Flink writers.
  • Truly bounded memory: peak RSS is bounded by one input batch + Parquet page buffer (~1 MiB × columns) + S3 multipart pool (~5 MiB × ~8 parts), regardless of target_file_size, dataset size, or number of files produced.
  • No public API change: same tbl.append(reader) / tbl.overwrite(reader) interface.

Fix

#3336

Activity

  1. azwanzuharimi commented on Sep 30, 2026

    @azwanzuharimi

    I am working on this and will open a PR this week.

    Both earlier PRs, #3336 and #3661, closed as stale before a review. Since then the file format writer API landed in #3119 and #3381. So a new PR must go through FileFormatWriter and not call pq.ParquetWriter directly.

    Plan:

    • Add a length() method to FileFormatWriter. It returns the bytes written so far, from OutputStream.tell(). This mirrors FileAppender.length() in Java.
    • In the streaming path, write each batch through the format writer. Roll to a new file when length() reaches write.target-file-size-bytes. This mirrors RollingFileWriter in Java.
    • Keep the pa.Table path unchanged.
    • Add tests that check file sizes on disk and that a large stream does not hold more than one batch in memory.

    @paultmathew tell me if you plan to revive #3336, so we do not do the work twice.

  2. paultmathew commented on Oct 1, 2026

    @paultmathew
    ContributorAuthor

    @azwanzuharimi no plans to revive it. Thanks for picking it up!

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions