zhuqi-lucas commented on code in PR #10901:
URL: https://github.com/apache/arrow-rs/pull/10901#discussion_r4119103008
##########
parquet/src/arrow/array_reader/cached_array_reader.rs:
##########
@@ -168,22 +172,33 @@ 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);
Review Comment:
Added in c3e3a4ec5, as `debug_assert!(current_batch_id.val >=
self.cleaned_up_to)`. That is the invariant the whole incremental scan rests on
— `outer_position` only moves forward, so the watermark can never run ahead of
the batch the reader is on.
##########
parquet/src/arrow/array_reader/cached_array_reader.rs:
##########
@@ -168,22 +172,33 @@ 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);
// 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.
+ // Rescanning from 0 each time made this quadratic in the number of
Review Comment:
Removed, you are right — it described the loop that is no longer there.
Moved to the commit message.
##########
parquet/src/arrow/array_reader/cached_array_reader.rs:
##########
@@ -168,22 +172,33 @@ 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);
// 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.
+ // Rescanning from 0 each time made this quadratic in the number of
+ // batches in the row group, and took the shared cache's write lock on
+ // every `consume_batch` even when there was nothing left to remove.
+ 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 batch_id_to_remove in start..end {
Review Comment:
Done, it collapses to `cache.remove(self.column_idx, BatchID { val });`.
--
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]