Skip to content

Commit 7e3acd3

Browse files
Limit the chunked read memory claim to CSV results
With Polars 1.39.0, scan_parquet(...).collect_batches() read a whole 676 MB UNLOAD object before the first batch, so a result-size-independent memory bound holds only for CSV. The benchmark README note about the Polars file cache now applies only to Polars before 1.39.0. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
1 parent 4f3f711 commit 7e3acd3

3 files changed

Lines changed: 6 additions & 6 deletions

File tree

‎benchmarks/README.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -225,7 +225,7 @@ API row conversion and native DataFrame/Table access appear as separate cases.
225225
Sampled RSS can miss short peaks, and the constructor suite's RSS includes its subsequent validation read.
226226
If the operating system denies access to per-thread CPU times, `thread_cpu_available` is false; thread counts and RSS are still recorded.
227227
Each trial process receives its own `POLARS_TEMP_DIR`, which the parent samples with RSS and removes after the trial.
228-
Polars lazy CSV scans of S3 objects download the whole object into a file cache in this directory before yielding batches; in the recorded runs, those files remained after the process exited.
228+
Polars before 1.39.0 downloaded the whole S3 object of a lazy CSV scan into a file cache in this directory before yielding batches and left the file after the process exited; PyAthena requires Polars 1.39.0 or later, whose scans do not write that cache.
229229
When the temporary directory is a tmpfs, as `/tmp` is on Amazon Linux 2023, this storage uses memory that RSS does not include.
230230

231231
| Capability | Treatment |

‎docs/polars.md‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -311,9 +311,9 @@ for chunk in cursor.iter_chunks():
311311
```
312312

313313
This method uses Polars' `scan_csv()` and `scan_parquet()` with `collect_batches()`,
314-
which stream the result from S3 without writing it to local storage.
315-
Polars reads ahead a bounded amount of data, so memory usage does not grow with the result size,
316-
but it is determined by Polars' read-ahead buffer rather than by `chunksize`.
314+
which read the result from S3 without writing it to local storage.
315+
Memory usage depends on how far Polars reads ahead, not on `chunksize`.
316+
For CSV results, Polars limits the read-ahead by the number of CPU cores, so memory usage does not grow with the result size.
317317
If reading a chunk fails, iteration raises `OperationalError` instead of ending early.
318318

319319
The chunked iteration also works with the unload option:

‎tests/pyathena/polars/test_result_set.py‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,7 @@ def _chunked_result_set() -> AthenaPolarsResultSet:
3030

3131
class TestAthenaPolarsResultSet:
3232
def test_iter_csv_chunks_raises_when_read_fails_partway(self, tmp_path):
33-
"""A CSV read that fails after some batches raises instead of ending the iteration."""
33+
"""A CSV read that fails partway through the data raises instead of ending early."""
3434
path = tmp_path / "result.csv"
3535
path.write_text(
3636
"a\n" + "".join(f"{i}\n" for i in range(_ROWS_BEFORE_FAILURE)) + "not-a-number\n"
@@ -61,7 +61,7 @@ def test_iter_csv_chunks_raises_when_read_fails_partway(self, tmp_path):
6161
list(result_set._iter_csv_chunks())
6262

6363
def test_iter_parquet_chunks_raises_when_read_fails_partway(self, tmp_path):
64-
"""A Parquet read that fails after some batches raises instead of ending the iteration."""
64+
"""A Parquet read that fails partway through the data raises instead of ending early."""
6565
pl.DataFrame({"a": range(_ROWS_BEFORE_FAILURE)}).write_parquet(
6666
tmp_path / "0.parquet", row_group_size=10_000
6767
)

0 commit comments

Comments
 (0)