This is an automated email from the ASF dual-hosted git repository.
alamb pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow-rs.git
The following commit(s) were added to refs/heads/main by this push:
new e605de6da5 Make CachedArrayReader cache cleanup incremental instead of
quadratic (#10901)
e605de6da5 is described below
commit e605de6da5200f508c4eb77d40d2c68cd44b02f1
Author: Qi Zhu <[email protected]>
AuthorDate: Mon Sep 28 22:10:45 2026 +0800
Make CachedArrayReader cache cleanup incremental instead of quadratic
(#10901)
## Which issue does this PR close?
- Part of #10774.
## Rationale for this change
`cleanup_consumed_batches` rescans batch ids from 0 on every
`consume_batch` call, re-removing ids that earlier calls already
removed, and takes the shared cache's **write lock** each time even when
there is nothing left to remove:
```rust
for batch_id_to_remove in 0..(current_batch_id.val - 1) {
cache.remove(...);
}
```
Over a row group with N cached batches this is O(N²) `remove` calls; the
consumer-role reader runs it once per output batch per cached column.
## What changes are included in this PR?
Track the already-cleaned frontier (`cleaned_up_to`) and remove only the
new range. When the frontier has not advanced, return without touching
the lock.
## Are these changes tested?
Covered by the existing `cached_array_reader` unit tests and the parquet
`--lib` suite (1331 passed; the one failure, `test_int96_interop`, is a
missing `parquet-testing` file unrelated to this change).
## Performance
Two `arrow_reader_clickbench` bot runs below: **neutral**. With
`batch_size = 8192`, N is small enough (~15 per row group on this
dataset) that the quadratic costs microseconds — so this is a hygiene
fix, not a measurable win at this scale. The pattern grows quadratically
with row-group size, and dropping the per-call write-lock acquisition
also matters more under concurrent readers than in this single-stream
bench.
(Per-run flags that did not reproduce: run 1 showed async Q20/Q21 swings
that vanished in run 2; async/Q22 shows +11% in both runs, but the same
query on the `sync` and `async_object_store` variants — same reader code
— is neutral-to-faster in both runs, and the async baseline's ±16ms
variance points at that variant's flakiness rather than this change.)
## Are there any user-facing changes?
No. `cleanup_consumed_batches` now takes `&mut self`, but it is a
private method of a `pub(crate)` type.
---
.../src/arrow/array_reader/cached_array_reader.rs | 91 +++++++++++++++++++---
1 file changed, 80 insertions(+), 11 deletions(-)
diff --git a/parquet/src/arrow/array_reader/cached_array_reader.rs
b/parquet/src/arrow/array_reader/cached_array_reader.rs
index 012fa25481..980d67f31b 100644
--- a/parquet/src/arrow/array_reader/cached_array_reader.rs
+++ b/parquet/src/arrow/array_reader/cached_array_reader.rs
@@ -87,6 +87,9 @@ pub struct CachedArrayReader {
local_cache: HashMap<BatchID, ArrayRef>,
/// Statistics to report on the Cache behavior
metrics: ArrowReaderMetrics,
+ /// Exclusive upper bound of the batch ids already removed from the shared
+ /// cache by [`Self::cleanup_consumed_batches`].
+ cleaned_up_to: usize,
}
impl CachedArrayReader {
@@ -111,6 +114,7 @@ impl CachedArrayReader {
role,
local_cache: HashMap::new(),
metrics,
+ cleaned_up_to: 0,
}
}
@@ -168,22 +172,28 @@ impl CachedArrayReader {
/// Remove batches from cache that have been completely consumed
/// This is only called for Consumer role readers
- fn cleanup_consumed_batches(&self) {
+ fn cleanup_consumed_batches(&mut self) {
let current_batch_id =
self.get_batch_id_from_position(self.outer_position);
+ // `outer_position` only moves forward, so the watermark can never run
+ // ahead of the batch the reader is currently on.
+ debug_assert!(current_batch_id.val >= self.cleaned_up_to);
// Remove batches that are at least one batch behind the current
position
// This ensures we don't remove batches that might still be needed for
the current batch
// We can safely remove batch_id if current_batch_id > batch_id + 1
- if current_batch_id.val > 1 {
- let mut cache = self.shared_cache.write().unwrap();
- for batch_id_to_remove in 0..(current_batch_id.val - 1) {
- cache.remove(
- self.column_idx,
- BatchID {
- val: batch_id_to_remove,
- },
- );
- }
+ if current_batch_id.val <= 1 {
+ return;
+ }
+ let end = current_batch_id.val - 1;
+ // Everything below `cleaned_up_to` was removed by an earlier call.
+ if end <= self.cleaned_up_to {
+ return;
+ }
+ let start = self.cleaned_up_to;
+ self.cleaned_up_to = end;
+ let mut cache = self.shared_cache.write().unwrap();
+ for val in start..end {
+ cache.remove(self.column_idx, BatchID { val });
}
}
}
@@ -631,6 +641,65 @@ mod tests {
assert!(cache.read().unwrap().get(0, BatchID { val: 2 }).is_some());
}
+ #[test]
+ fn test_consumer_cleanup_after_skip() {
+ // A consumer that skips several batches (e.g. rows filtered out by a
+ // predicate) must still remove the skipped batches from the shared
+ // cache, even though it never fetched them itself: one cleanup call
+ // covers the whole skipped range.
+ let metrics = ArrowReaderMetrics::disabled();
+ let cache = Arc::new(RwLock::new(RowGroupCache::new(3, usize::MAX)));
// Batch size 3
+
+ // Producer populates batches 0..=3 (12 values).
+ let producer_values: Vec<i32> = (1..=12).collect();
+ let mut producer = CachedArrayReader::new(
+ Box::new(MockArrayReader::new(producer_values.clone())),
+ cache.clone(),
+ 0,
+ CacheRole::Producer,
+ metrics.clone(),
+ );
+ for _ in 0..4 {
+ producer.read_records(3).unwrap();
+ producer.consume_batch().unwrap();
+ }
+ for batch in 0..4 {
+ assert!(
+ cache
+ .read()
+ .unwrap()
+ .get(0, BatchID { val: batch })
+ .is_some()
+ );
+ }
+
+ // Consumer skips batches 0..=2 outright and reads batch 3.
+ let mut consumer = CachedArrayReader::new(
+ Box::new(MockArrayReader::new(producer_values)),
+ cache.clone(),
+ 0,
+ CacheRole::Consumer,
+ metrics,
+ );
+ assert_eq!(consumer.skip_records(9).unwrap(), 9);
+ assert_eq!(consumer.read_records(3).unwrap(), 3);
+ let array = consumer.consume_batch().unwrap();
+ assert_eq!(array.len(), 3);
+
+ // current_batch_id = 12 / 3 = 4, so cleanup covers batches 0..3 in a
+ // single call, including the ones the consumer never fetched.
+ for batch in 0..3 {
+ assert!(
+ cache
+ .read()
+ .unwrap()
+ .get(0, BatchID { val: batch })
+ .is_none()
+ );
+ }
+ assert!(cache.read().unwrap().get(0, BatchID { val: 3 }).is_some());
+ }
+
#[test]
fn test_producer_keeps_batches() {
let metrics = ArrowReaderMetrics::disabled();