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]
