From 104c20829d74cd897a9d05978f70922e6cab6afa Mon Sep 17 00:00:00 2001 From: Nathan Bezualem Date: Sun, 29 Mar 2026 17:57:54 -0400 Subject: [PATCH 1/5] Add API to clear all buffered ranges --- parquet/src/arrow/push_decoder/mod.rs | 71 +++++++++++++++++++ .../arrow/push_decoder/reader_builder/mod.rs | 5 ++ parquet/src/arrow/push_decoder/remaining.rs | 5 ++ parquet/src/util/push_buffers.rs | 7 ++ 4 files changed, 88 insertions(+) diff --git a/parquet/src/arrow/push_decoder/mod.rs b/parquet/src/arrow/push_decoder/mod.rs index cdb0715edb55..27234d716030 100644 --- a/parquet/src/arrow/push_decoder/mod.rs +++ b/parquet/src/arrow/push_decoder/mod.rs @@ -365,6 +365,15 @@ impl ParquetPushDecoder { pub fn buffered_bytes(&self) -> u64 { self.state.buffered_bytes() } + + /// Release any staged ranges currently buffered for future decode work. + /// + /// This clears byte ranges still owned by the decoder's internal + /// [`PushBuffers`]. It does not affect any data that has already been handed + /// off to an active [`ParquetRecordBatchReader`]. + pub fn release_all_ranges(&mut self) { + self.state.release_all_ranges(); + } } /// Internal state machine for the [`ParquetPushDecoder`] @@ -573,6 +582,20 @@ impl ParquetDecoderState { ParquetDecoderState::Finished => 0, } } + + /// Release any staged ranges currently buffered in the decoder. + fn release_all_ranges(&mut self) { + match self { + ParquetDecoderState::ReadingRowGroup { + remaining_row_groups, + } => remaining_row_groups.release_all_ranges(), + ParquetDecoderState::DecodingRowGroup { + record_batch_reader: _, + remaining_row_groups, + } => remaining_row_groups.release_all_ranges(), + ParquetDecoderState::Finished => {} + } + } } #[cfg(test)] @@ -665,6 +688,54 @@ mod test { assert_eq!(all_output, *TEST_BATCH); } + /// Releasing staged ranges should free speculative buffers without affecting + /// the active row group reader. + #[test] + fn test_decoder_release_all_ranges() { + let mut decoder = ParquetPushDecoderBuilder::try_new_decoder(test_file_parquet_metadata()) + .unwrap() + .with_batch_size(100) + .build() + .unwrap(); + + decoder + .push_range(test_file_range(), TEST_FILE_DATA.clone()) + .unwrap(); + assert_eq!(decoder.buffered_bytes(), test_file_len()); + + // The current row group reader is built from the prefetched bytes, but + // the speculative full-file range remains staged in the decoder. + let batch1 = expect_data(decoder.try_decode()); + assert_eq!(batch1, TEST_BATCH.slice(0, 100)); + assert_eq!(decoder.buffered_bytes(), test_file_len()); + + decoder.release_all_ranges(); + assert_eq!(decoder.buffered_bytes(), 0); + + // The active reader still owns the current row group's bytes, so it can + // continue decoding without consulting PushBuffers. + let batch2 = expect_data(decoder.try_decode()); + assert_eq!(batch2, TEST_BATCH.slice(100, 100)); + assert_eq!(decoder.buffered_bytes(), 0); + + // Moving to the next row group now requires the decoder to ask for data + // again because the staged speculative ranges were released. + let ranges = expect_needs_data(decoder.try_decode()); + let num_bytes_requested: u64 = ranges.iter().map(|r| r.end - r.start).sum(); + push_ranges_to_decoder(&mut decoder, ranges); + assert_eq!(decoder.buffered_bytes(), num_bytes_requested); + + let batch3 = expect_data(decoder.try_decode()); + assert_eq!(batch3, TEST_BATCH.slice(200, 100)); + assert_eq!(decoder.buffered_bytes(), 0); + + let batch4 = expect_data(decoder.try_decode()); + assert_eq!(batch4, TEST_BATCH.slice(300, 100)); + assert_eq!(decoder.buffered_bytes(), 0); + + expect_finished(decoder.try_decode()); + } + /// Decode the entire file incrementally, simulating partial reads #[test] fn test_decoder_partial() { diff --git a/parquet/src/arrow/push_decoder/reader_builder/mod.rs b/parquet/src/arrow/push_decoder/reader_builder/mod.rs index d3d78ca7c263..6833f79404bf 100644 --- a/parquet/src/arrow/push_decoder/reader_builder/mod.rs +++ b/parquet/src/arrow/push_decoder/reader_builder/mod.rs @@ -212,6 +212,11 @@ impl RowGroupReaderBuilder { self.buffers.buffered_bytes() } + /// Release any staged ranges currently buffered for future decode work. + pub fn release_all_ranges(&mut self) { + self.buffers.clear_all_ranges(); + } + /// take the current state, leaving None in its place. /// /// Returns an error if there the state wasn't put back after the previous diff --git a/parquet/src/arrow/push_decoder/remaining.rs b/parquet/src/arrow/push_decoder/remaining.rs index 4613fda08749..65e79aa50bb2 100644 --- a/parquet/src/arrow/push_decoder/remaining.rs +++ b/parquet/src/arrow/push_decoder/remaining.rs @@ -70,6 +70,11 @@ impl RemainingRowGroups { self.row_group_reader_builder.buffered_bytes() } + /// Release any staged ranges currently buffered for future decode work + pub fn release_all_ranges(&mut self) { + self.row_group_reader_builder.release_all_ranges(); + } + /// returns [`ParquetRecordBatchReader`] suitable for reading the next /// group of rows from the Parquet data, or the list of data ranges still /// needed to proceed diff --git a/parquet/src/util/push_buffers.rs b/parquet/src/util/push_buffers.rs index 0c00cf9bd57f..eb4982fb3c6f 100644 --- a/parquet/src/util/push_buffers.rs +++ b/parquet/src/util/push_buffers.rs @@ -154,6 +154,13 @@ impl PushBuffers { self.ranges = new_ranges; self.buffers = new_buffers; } + + /// Clear all buffered ranges and their corresponding data + #[cfg(feature = "arrow")] + pub fn clear_all_ranges(&mut self) { + self.ranges.clear(); + self.buffers.clear(); + } } impl Length for PushBuffers { From 20acef257b2872885bc686810e6c5083fd8c408c Mon Sep 17 00:00:00 2001 From: Nathan Bezualem Date: Sun, 29 Mar 2026 18:14:32 -0400 Subject: [PATCH 2/5] Clarify release_all_ranges docs --- parquet/src/arrow/push_decoder/mod.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/parquet/src/arrow/push_decoder/mod.rs b/parquet/src/arrow/push_decoder/mod.rs index 27234d716030..a01b00153169 100644 --- a/parquet/src/arrow/push_decoder/mod.rs +++ b/parquet/src/arrow/push_decoder/mod.rs @@ -366,7 +366,7 @@ impl ParquetPushDecoder { self.state.buffered_bytes() } - /// Release any staged ranges currently buffered for future decode work. + /// Release any staged byte ranges currently buffered for future decode work. /// /// This clears byte ranges still owned by the decoder's internal /// [`PushBuffers`]. It does not affect any data that has already been handed From c61729988307aeb14e4c0b04b1c5419e2344563a Mon Sep 17 00:00:00 2001 From: Nathan Bezualem Date: Sun, 29 Mar 2026 18:52:27 -0400 Subject: [PATCH 3/5] Tighten release_all_ranges test comment --- parquet/src/arrow/push_decoder/mod.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/parquet/src/arrow/push_decoder/mod.rs b/parquet/src/arrow/push_decoder/mod.rs index a01b00153169..9cb4d4242ac7 100644 --- a/parquet/src/arrow/push_decoder/mod.rs +++ b/parquet/src/arrow/push_decoder/mod.rs @@ -709,6 +709,7 @@ mod test { assert_eq!(batch1, TEST_BATCH.slice(0, 100)); assert_eq!(decoder.buffered_bytes(), test_file_len()); + // All of the buffer is released decoder.release_all_ranges(); assert_eq!(decoder.buffered_bytes(), 0); From 13da3d5113fea2af647d579929f82d71d25785de Mon Sep 17 00:00:00 2001 From: Nathan Bezualem Date: Mon, 6 Apr 2026 17:02:27 -0400 Subject: [PATCH 4/5] rename release_all_ranges to clear_all_ranges --- parquet/src/arrow/push_decoder/mod.rs | 18 +++++++++--------- .../arrow/push_decoder/reader_builder/mod.rs | 4 ++-- parquet/src/arrow/push_decoder/remaining.rs | 6 +++--- 3 files changed, 14 insertions(+), 14 deletions(-) diff --git a/parquet/src/arrow/push_decoder/mod.rs b/parquet/src/arrow/push_decoder/mod.rs index 9cb4d4242ac7..04707ec46777 100644 --- a/parquet/src/arrow/push_decoder/mod.rs +++ b/parquet/src/arrow/push_decoder/mod.rs @@ -366,13 +366,13 @@ impl ParquetPushDecoder { self.state.buffered_bytes() } - /// Release any staged byte ranges currently buffered for future decode work. + /// Clear any staged byte ranges currently buffered for future decode work. /// /// This clears byte ranges still owned by the decoder's internal /// [`PushBuffers`]. It does not affect any data that has already been handed /// off to an active [`ParquetRecordBatchReader`]. - pub fn release_all_ranges(&mut self) { - self.state.release_all_ranges(); + pub fn clear_all_ranges(&mut self) { + self.state.clear_all_ranges(); } } @@ -583,16 +583,16 @@ impl ParquetDecoderState { } } - /// Release any staged ranges currently buffered in the decoder. - fn release_all_ranges(&mut self) { + /// Clear any staged ranges currently buffered in the decoder. + fn clear_all_ranges(&mut self) { match self { ParquetDecoderState::ReadingRowGroup { remaining_row_groups, - } => remaining_row_groups.release_all_ranges(), + } => remaining_row_groups.clear_all_ranges(), ParquetDecoderState::DecodingRowGroup { record_batch_reader: _, remaining_row_groups, - } => remaining_row_groups.release_all_ranges(), + } => remaining_row_groups.clear_all_ranges(), ParquetDecoderState::Finished => {} } } @@ -691,7 +691,7 @@ mod test { /// Releasing staged ranges should free speculative buffers without affecting /// the active row group reader. #[test] - fn test_decoder_release_all_ranges() { + fn test_decoder_clear_all_ranges() { let mut decoder = ParquetPushDecoderBuilder::try_new_decoder(test_file_parquet_metadata()) .unwrap() .with_batch_size(100) @@ -710,7 +710,7 @@ mod test { assert_eq!(decoder.buffered_bytes(), test_file_len()); // All of the buffer is released - decoder.release_all_ranges(); + decoder.clear_all_ranges(); assert_eq!(decoder.buffered_bytes(), 0); // The active reader still owns the current row group's bytes, so it can diff --git a/parquet/src/arrow/push_decoder/reader_builder/mod.rs b/parquet/src/arrow/push_decoder/reader_builder/mod.rs index 6833f79404bf..922d8070c064 100644 --- a/parquet/src/arrow/push_decoder/reader_builder/mod.rs +++ b/parquet/src/arrow/push_decoder/reader_builder/mod.rs @@ -212,8 +212,8 @@ impl RowGroupReaderBuilder { self.buffers.buffered_bytes() } - /// Release any staged ranges currently buffered for future decode work. - pub fn release_all_ranges(&mut self) { + /// Clear any staged ranges currently buffered for future decode work. + pub fn clear_all_ranges(&mut self) { self.buffers.clear_all_ranges(); } diff --git a/parquet/src/arrow/push_decoder/remaining.rs b/parquet/src/arrow/push_decoder/remaining.rs index 65e79aa50bb2..2986ca0da8d8 100644 --- a/parquet/src/arrow/push_decoder/remaining.rs +++ b/parquet/src/arrow/push_decoder/remaining.rs @@ -70,9 +70,9 @@ impl RemainingRowGroups { self.row_group_reader_builder.buffered_bytes() } - /// Release any staged ranges currently buffered for future decode work - pub fn release_all_ranges(&mut self) { - self.row_group_reader_builder.release_all_ranges(); + /// Clear any staged ranges currently buffered for future decode work + pub fn clear_all_ranges(&mut self) { + self.row_group_reader_builder.clear_all_ranges(); } /// returns [`ParquetRecordBatchReader`] suitable for reading the next From f8bd0f35424bc9ba387fcd671bf999f6823f1767 Mon Sep 17 00:00:00 2001 From: Andrew Lamb Date: Tue, 7 Apr 2026 11:33:31 -0400 Subject: [PATCH 5/5] fix doc build --- parquet/src/arrow/push_decoder/mod.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/parquet/src/arrow/push_decoder/mod.rs b/parquet/src/arrow/push_decoder/mod.rs index 04707ec46777..24384471a4ec 100644 --- a/parquet/src/arrow/push_decoder/mod.rs +++ b/parquet/src/arrow/push_decoder/mod.rs @@ -369,7 +369,7 @@ impl ParquetPushDecoder { /// Clear any staged byte ranges currently buffered for future decode work. /// /// This clears byte ranges still owned by the decoder's internal - /// [`PushBuffers`]. It does not affect any data that has already been handed + /// `PushBuffers`. It does not affect any data that has already been handed /// off to an active [`ParquetRecordBatchReader`]. pub fn clear_all_ranges(&mut self) { self.state.clear_all_ranges();