srujankgandla opened a new pull request, #4028:
URL: https://github.com/apache/iceberg-python/pull/4028
Closes #2407
# Rationale for this change
`to_arrow_batch_reader()` is documented as the low-memory way to stream scan
results ("a RecordBatch is read one at a time"), but in practice it
materialized every file up front: `to_record_batches()` wraps each file's
iterator in `list()` inside `batches_for_task()`, and `executor.map` submits
all file tasks eagerly, so peak memory scaled with file count rather than batch
size.
Reproduced locally: 5.9 MB across 12 parquet files peaked at ~74 MB RSS
(12.6x the on-disk size), and instrumenting the reader showed all 12 files
fully read even when the consumer took a single batch.
This PR adds `ArrowScan.to_record_batches_lazy()`, which walks scan tasks
sequentially in the calling thread and streams each file's batches directly,
with per-task delete-file reads instead of the eager
`_read_all_delete_files()`. `_to_arrow_batch_reader_via_file_scan_tasks` now
uses the lazy path. `to_table()`, `to_pandas()`, and all other callers keep the
existing threaded path untouched, so there is no throughput regression on the
eager paths.
## Are these changes tested?
Yes — three new tests in `tests/io/test_pyarrow.py`:
- `test_to_arrow_batch_reader_does_not_read_ahead`: consumes one batch from
a 4-file scan and asserts only 1 file was opened for reading (fails on the old
code).
- `test_to_record_batches_lazy_matches_eager`: lazy and threaded paths
return identical rows in identical order for limits None/0/1/100/total/total+10.
- `test_to_record_batches_lazy_applies_positional_deletes`: positional
deletes are applied correctly through the new per-task delete path ([1,2,3,4] →
[1,2,4]).
Also ran the full `tests/io/test_pyarrow.py` suite and the table scan tests
— no regressions versus the clean tree (the only failures are pre-existing
environment issues also present without this change). `ruff check` and `ruff
format` are clean.
## Are there any user-facing changes?
Yes — `to_arrow_batch_reader()` now honors its documented contract: batches
are read one at a time and memory stays bounded by in-flight batches instead of
scaling with the number of files. This is a bug fix aligning behavior with the
documented API, no API signature changes.
--
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]