Roll streaming writes at the on-disk target file size - #4033
Open
azwanzuharimi wants to merge 3 commits into
Open
azwanzuharimi wants to merge 3 commits into
azwanzuharimi wants to merge 3 commits into
Conversation
Table.append and Table.overwrite with a pa.RecordBatchReader grouped batches by their Arrow memory size before they were written. So the target file size was compared against uncompressed bytes, and files on disk landed 3 to 10 times smaller than write.target-file-size-bytes. Peak memory was one whole group of batches. Add length() to FileFormatWriter, like FileAppender.length() in Java. ParquetFormatWriter buffers rows up to write.parquet.row-group-limit or write.parquet.row-group-size-bytes and reports the bytes on disk plus the buffered bytes. The streaming path writes each batch through the format writer and rolls to a new file when length() reaches the target, like RollingFileWriter in Java. The pa.Table path does not change. Closes apache#3388
Author
|
@Fokko @kevinjqliu this picks up #3336 on the file format writer API from #3381. The two earlier PRs closed as stale without a review, so a first look on the design would help. |
write_file and the streaming branch of _dataframe_to_data_files resolved the file format, the format model, the location provider and the file schema with the same block. Move it into one helper.
length() counted buffered rows at their Arrow size. Parquet compresses them 3 to 10 times, so a rolled file landed under the target by up to one row group. Measure the disk bytes each flush adds and divide by the Arrow bytes it wrote. Scale the buffered bytes by that ratio in length(). The ratio is 1.0 until the first flush in each file.
Fokko
self-requested a review
October 1, 2026 19:20
This branch has not been deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Closes #3388
Rationale for this change
When you append a
pa.RecordBatchReader, PyIceberg groups batches by their size in memory before it writes them. Parquet then compresses the data, so files on disk come out 3 to 10 times smaller thanwrite.target-file-size-bytes. Memory also grows with the target, because a whole group waits in memory before the write starts.This change writes each batch as it arrives and rolls to a new file when the file on disk reaches the target. It is the same idea as
RollingFileWriterin Java.What changes:
FileFormatWritergets a new abstract method,length(). It returns the estimated file size so far, likeFileAppender.length()in Java. This is a breaking change for anyone who subclassedFileFormatWriteroutside this repo. There are no subclasses in the repo.ParquetFormatWriterbuffers rows and flushes a row group when the buffer reacheswrite.parquet.row-group-limitrows orwrite.parquet.row-group-size-bytes. Without this buffer every input batch would become its own row group.length()counts the bytes on disk plus the buffered rows. Buffered rows are scaled by the compression ratio seen on the last flush, so files land close to the target._dataframe_to_data_fileswrites one file at a time through the format writer and rolls whenlength()reaches the target.pa.Tablepath does not change. It still callswrite()once andclose()once, so it writes the same row groups and the same bytes as before.bin_pack_record_batchesstays, because it is public. Nothing inpyicebergcalls it now.write.parquet.row-group-size-bytesno longer warns as unsupported, because the streaming writer reads it. The configuration doc says it applies to streaming writes only.Measured on macOS with pyarrow 25.0.1, 64 MiB target, default properties. The value is file size on disk divided by the target.
Memory during a streaming write is bounded by one row group plus one input batch and the writer's buffers. It no longer depends on the target.
Are these changes tested?
Yes.
tests/io/test_format_writers.py:length()starts at 0, grows with each write, and equals the file size after close. Buffered rows count in full before the first flush and are scaled after it. A failed flush on the exception path closes the stream and keeps the original error.tests/io/test_pyarrow.py: the streaming path rolls at the on-disk target, lands within a few percent of the target on compressible data, writes one file for a large target, packs small batches into row groups, flushes row groups on bytes, yields the first file before it reads the rest of the reader, closes the stream when the reader raises, and stays fast with 5,000 one row batches.write_filewrites the same bytes as a directpq.ParquetWriterwrite for six settings of the row group properties.tests/catalog/test_catalog_behaviors.py: the existingRecordBatchReadertests pass unchanged.make lintpasses.tests/io,tests/catalog/test_catalog_behaviors.pyandtests/tablepass, with the S3, ADLS, GCS and integration markers deselected.Are there any user-facing changes?
Yes.
write.target-file-size-byteson disk, and memory no longer grows with the target.FileFormatWriterhas a new abstract method,length().write.parquet.row-group-size-bytesno longer warns as unsupported. It applies to streaming writes.