azwanzuharimi opened a new pull request, #4033:
URL: https://github.com/apache/iceberg-python/pull/4033

   Closes #3388
   
   # Rationale for this change
   
   `Table.append(reader)` and `Table.overwrite(reader)` with a 
`pa.RecordBatchReader` grouped batches by `RecordBatch.nbytes` in 
`bin_pack_record_batches` before they were written. So 
`write.target-file-size-bytes` was compared against uncompressed Arrow bytes. 
Files on disk landed 3 to 10 times smaller than the target, and peak memory was 
one whole group of batches. Both caveats are in the current docstrings.
   
   #3336 and #3661 fixed this with a direct `pq.ParquetWriter` and closed as 
stale. Since then the file format writer API landed in #3119 and #3381. This 
change goes through that API.
   
   - `FileFormatWriter.length()` is a new abstract method. It returns the 
estimated file size in bytes, including buffered rows. This mirrors 
`FileAppender.length()` in Java. It is a breaking change for a third party 
subclass of `FileFormatWriter`; the repo has none.
   - `ParquetFormatWriter` buffers the tables it receives and flushes a row 
group when the buffer reaches `write.parquet.row-group-limit` rows or 
`write.parquet.row-group-size-bytes`. `length()` returns `OutputStream.tell()` 
plus the buffered bytes, like `ParquetWriter.length()` in Java. Buffered bytes 
are uncompressed Arrow bytes, so a file can land up to 
`write.parquet.row-group-size-bytes` under the target.
   - The `pa.RecordBatchReader` branch of `_dataframe_to_data_files` writes 
each batch through the format writer. It rolls to a new file when `length()` 
reaches the target, like `RollingFileWriter` in Java. It writes one file at a 
time. Memory is bounded by one row group plus one input batch and the writer's 
buffers.
   - The `pa.Table` path does not change. It calls `write()` once and `close()` 
once, which gives the same row groups as before. A test compares the file bytes 
with a direct `pq.ParquetWriter` write for six settings of the two row group 
properties.
   - `bin_pack_record_batches` stays, because it is public. Nothing in 
`pyiceberg` calls it now.
   - `write.parquet.row-group-size-bytes` is removed from the unsupported 
property warning, because the streaming writer reads it. The configuration doc 
says it applies to streaming writes only and that the size is uncompressed 
Arrow bytes.
   
   Measured locally on macOS with pyarrow 25.0.1, 400 batches of 10,000 rows 
(int64, uuid string, random float64), target 64 MiB, default row group 
properties. Before: 3 files at 0.69, 0.69 and 0.19 of the target. After: the 
same, because one default row group of this data is 56 MiB uncompressed and the 
row limit fires first. With `write.parquet.row-group-size-bytes` at 8 MiB: 2 
files at 0.97 and 0.79 of the target. With a 1 MiB target and a 10,000 row 
limit: 13 files at about 1.15 of the target.
   
   With 5,000 batches of 100 rows, the file has 1 row group and 1,043,980 
bytes, the same as `main`. With 5,000 batches of 1 row, the write takes 0.09 s.
   
   ## Are these changes tested?
   
   Yes.
   
   - `tests/io/test_format_writers.py`: `length()` is 0 before the first write, 
grows after each write, and equals the file size after close. A failed flush on 
the exception path closes the stream and keeps the original error.
   - `tests/io/test_pyarrow.py`: the streaming path rolls files at the on-disk 
target, 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 linear 
with 5,000 one row batches. `write_file` writes the same bytes as a direct 
`pq.ParquetWriter` write.
   - `tests/catalog/test_catalog_behaviors.py`: the existing 
`RecordBatchReader` tests pass unchanged.
   
   `make lint` passes. `tests/io`, `tests/catalog/test_catalog_behaviors.py` 
and `tests/table` pass with 923 tests, with the S3, ADLS, GCS and integration 
markers deselected.
   
   ## Are there any user-facing changes?
   
   Yes.
   
   - Streaming writes now produce files close to `write.target-file-size-bytes` 
on disk, and memory no longer grows with the target.
   - `FileFormatWriter` has a new abstract method `length()`.
   - `write.parquet.row-group-size-bytes` no longer warns as unsupported. It 
applies to streaming writes.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to