adriangb opened a new issue, #11233:
URL: https://github.com/apache/arrow-rs/issues/11233
**Describe the bug**
`ParquetPushDecoder` releases a pushed buffer only if its range is equal to
a range that the decoder requested. If the caller fetches the requested ranges
with one larger request and pushes one buffer, the decoder does not release
that buffer. The buffered bytes increase with each row group, until they hold
all of the data that the scan read. The bytes are freed only when the scan
finishes.
Callers often merge nearby ranges into one request to object storage, so
this pattern is common.
**To Reproduce**
A file with 4 row groups and 2 columns. The caller pushes one buffer per
`NeedsData`, from the start of the first requested range to the end of the last
one:
<details><summary>Test (public API only)</summary>
```rust
use std::sync::Arc;
use arrow_array::{ArrayRef, Int64Array, RecordBatch};
use bytes::Bytes;
use parquet::DecodeResult;
use parquet::arrow::ArrowWriter;
use parquet::arrow::push_decoder::ParquetPushDecoderBuilder;
use parquet::file::metadata::ParquetMetaDataReader;
use parquet::file::properties::WriterProperties;
#[test]
fn mre() {
// A file with 4 row groups of 10_000 rows and 2 columns.
let a: ArrayRef = Arc::new(Int64Array::from_iter_values(0..40_000));
let b: ArrayRef = Arc::new(Int64Array::from_iter_values(0..40_000));
let batch = RecordBatch::try_from_iter([("a", a), ("b", b)]).unwrap();
let props = WriterProperties::builder()
.set_max_row_group_row_count(Some(10_000))
.build();
let mut file = vec![];
let mut writer = ArrowWriter::try_new(&mut file, batch.schema(),
Some(props)).unwrap();
writer.write(&batch).unwrap();
writer.close().unwrap();
let file = Bytes::from(file);
let metadata =
ParquetMetaDataReader::new().parse_and_finish(&file).unwrap();
let mut decoder =
ParquetPushDecoderBuilder::try_new_decoder(Arc::new(metadata))
.unwrap()
.build()
.unwrap();
loop {
match decoder.try_decode().unwrap() {
DecodeResult::NeedsData(ranges) => {
// Fetch the requested ranges with one request, as a caller
// that coalesces nearby ranges does, and push one buffer.
let start = ranges.iter().map(|r| r.start).min().unwrap();
let end = ranges.iter().map(|r| r.end).max().unwrap();
decoder
.push_range(start..end, file.slice(start as usize..end
as usize))
.unwrap();
}
DecodeResult::Data(batch) => {
println!(
"decoded {} rows, buffered_bytes = {}",
batch.num_rows(),
decoder.buffered_bytes()
);
}
DecodeResult::Finished => break,
}
}
}
```
</details>
| After row group | `buffered_bytes()`, one buffer per request |
`buffered_bytes()`, one buffer per requested range |
|---|---:|---:|
| 0 | 195,128 | 0 |
| 1 | 390,256 | 0 |
| 2 | 585,384 | 0 |
| 3 | 780,512 | 0 |
Tested on `main` (92b5042818).
**Expected behavior**
When the decoder is done with a row group, it releases all pushed bytes of
that row group, whatever the shape of the pushed buffers. `buffered_bytes()`
does not increase from one row group to the next.
**Additional context**
The cause is `PushBuffers::clear_ranges`, which `reader_builder/data.rs`
calls after it builds a row group. It removes a buffer only if its range is
exactly equal to a requested range.
This is related to https://github.com/apache/arrow-rs/issues/6946 (buffered
memory in the push decoder), but it is a separate problem: it occurs with the
default row-group granularity.
--
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]