From 87ce023f671ef8265b1037a3597e14b312ce894b Mon Sep 17 00:00:00 2001 From: Hippolyte Barraud Date: Fri, 20 Feb 2026 00:02:55 -0500 Subject: [PATCH 1/5] parquet: stream-encode definition/repetition levels incrementally Previously, the column writer accumulated raw definition and repetition levels in `Vec` sinks (`def_levels_sink` / `rep_levels_sink`) and only RLE-encoded them in bulk at page-flush time. Replace the two sinks with streaming `LevelEncoder` fields. Levels are now encoded incrementally as each `write_batch` call arrives, so only the compact encoded bytes are held in memory at all times. At page flush, the encoder is consumed and its bytes are written directly into the page buffer; a fresh encoder is swapped in for the next page. Signed-off-by: Hippolyte Barraud --- parquet/src/column/writer/mod.rs | 83 +++++++++++++++----------------- parquet/src/encodings/levels.rs | 40 ++++++++++++++- 2 files changed, 77 insertions(+), 46 deletions(-) diff --git a/parquet/src/column/writer/mod.rs b/parquet/src/column/writer/mod.rs index c014397f132e..f74211deccf3 100644 --- a/parquet/src/column/writer/mod.rs +++ b/parquet/src/column/writer/mod.rs @@ -24,6 +24,7 @@ use crate::bloom_filter::Sbbf; use crate::file::page_index::column_index::ColumnIndexMetaData; use crate::file::page_index::offset_index::OffsetIndexMetaData; use std::collections::{BTreeSet, VecDeque}; +use std::mem; use std::str; use crate::basic::{ @@ -342,9 +343,9 @@ pub struct GenericColumnWriter<'a, E: ColumnValueEncoder> { /// but we use a BTreeSet so that the output is deterministic encodings: BTreeSet, encoding_stats: Vec, - // Reused buffers - def_levels_sink: Vec, - rep_levels_sink: Vec, + // Streaming level encoders for definition/repetition levels. + def_levels_encoder: LevelEncoder, + rep_levels_encoder: LevelEncoder, data_pages: VecDeque, // column index and offset index column_index_builder: ColumnIndexBuilder, @@ -401,6 +402,9 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { _ => None, }; + let def_levels_encoder = Self::create_level_encoder(descr.max_def_level(), &props); + let rep_levels_encoder = Self::create_level_encoder(descr.max_rep_level(), &props); + Self { descr, props, @@ -409,8 +413,8 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { codec, compressor, encoder, - def_levels_sink: vec![], - rep_levels_sink: vec![], + def_levels_encoder, + rep_levels_encoder, data_pages: VecDeque::new(), page_metrics, column_metrics, @@ -638,6 +642,14 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { }) } + /// Creates a new streaming level encoder appropriate for the writer version. + fn create_level_encoder(max_level: i16, props: &WriterProperties) -> LevelEncoder { + match props.writer_version() { + WriterVersion::PARQUET_1_0 => LevelEncoder::v1_streaming(Encoding::RLE, max_level), + WriterVersion::PARQUET_2_0 => LevelEncoder::v2_streaming(max_level), + } + } + /// Writes mini batch of values, definition and repetition levels. /// This allows fine-grained processing of values and maintaining a reasonable /// page size. @@ -668,7 +680,7 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { // Update histogram self.page_metrics.update_definition_level_histogram(levels); - self.def_levels_sink.extend_from_slice(levels); + self.def_levels_encoder.put(levels); values_to_write } else { num_levels @@ -699,7 +711,7 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { // Update histogram self.page_metrics.update_repetition_level_histogram(levels); - self.rep_levels_sink.extend_from_slice(levels); + self.rep_levels_encoder.put(levels); } else { // Each value is exactly one row. // Equals to the number of values, we count nulls as well. @@ -1051,23 +1063,19 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { let mut buffer = vec![]; if max_rep_level > 0 { - buffer.extend_from_slice( - &self.encode_levels_v1( - Encoding::RLE, - &self.rep_levels_sink[..], - max_rep_level, - )[..], + let encoder = mem::replace( + &mut self.rep_levels_encoder, + Self::create_level_encoder(max_rep_level, &self.props), ); + buffer.extend_from_slice(&encoder.consume()); } if max_def_level > 0 { - buffer.extend_from_slice( - &self.encode_levels_v1( - Encoding::RLE, - &self.def_levels_sink[..], - max_def_level, - )[..], + let encoder = mem::replace( + &mut self.def_levels_encoder, + Self::create_level_encoder(max_def_level, &self.props), ); + buffer.extend_from_slice(&encoder.consume()); } buffer.extend_from_slice(&values_data.buf); @@ -1097,15 +1105,23 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { let mut buffer = vec![]; if max_rep_level > 0 { - let levels = self.encode_levels_v2(&self.rep_levels_sink[..], max_rep_level); + let encoder = mem::replace( + &mut self.rep_levels_encoder, + Self::create_level_encoder(max_rep_level, &self.props), + ); + let levels = encoder.consume(); rep_levels_byte_len = levels.len(); - buffer.extend_from_slice(&levels[..]); + buffer.extend_from_slice(&levels); } if max_def_level > 0 { - let levels = self.encode_levels_v2(&self.def_levels_sink[..], max_def_level); + let encoder = mem::replace( + &mut self.def_levels_encoder, + Self::create_level_encoder(max_def_level, &self.props), + ); + let levels = encoder.consume(); def_levels_byte_len = levels.len(); - buffer.extend_from_slice(&levels[..]); + buffer.extend_from_slice(&levels); } let uncompressed_size = @@ -1155,10 +1171,6 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { // Update total number of rows. self.column_metrics.total_rows_written += self.page_metrics.num_buffered_rows as u64; - - // Reset state. - self.rep_levels_sink.clear(); - self.def_levels_sink.clear(); self.page_metrics.new_page(); Ok(()) @@ -1235,23 +1247,6 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { Ok(metadata) } - /// Encodes definition or repetition levels for Data Page v1. - #[inline] - fn encode_levels_v1(&self, encoding: Encoding, levels: &[i16], max_level: i16) -> Vec { - let mut encoder = LevelEncoder::v1(encoding, max_level, levels.len()); - encoder.put(levels); - encoder.consume() - } - - /// Encodes definition or repetition levels for Data Page v2. - /// Encoding is always RLE. - #[inline] - fn encode_levels_v2(&self, levels: &[i16], max_level: i16) -> Vec { - let mut encoder = LevelEncoder::v2(max_level, levels.len()); - encoder.put(levels); - encoder.consume() - } - /// Writes compressed data page into underlying sink and updates global metrics. #[inline] fn write_data_page(&mut self, page: CompressedPage) -> Result<()> { diff --git a/parquet/src/encodings/levels.rs b/parquet/src/encodings/levels.rs index 8c95e5ca51aa..c1c687fb2c8c 100644 --- a/parquet/src/encodings/levels.rs +++ b/parquet/src/encodings/levels.rs @@ -26,6 +26,7 @@ use crate::util::bit_util::{BitWriter, ceil, num_required_bits}; /// Computes max buffer size for level encoder/decoder based on encoding, max /// repetition/definition level and number of total buffered values (includes null /// values). +#[allow(dead_code)] #[inline] pub fn max_buffer_size(encoding: Encoding, max_level: i16, num_buffered_values: usize) -> usize { let bit_width = num_required_bits(max_level as u64); @@ -53,6 +54,7 @@ impl LevelEncoder { /// Used to encode levels for Data Page v1. /// /// Panics, if encoding is not supported. + #[allow(dead_code)] pub fn v1(encoding: Encoding, max_level: i16, capacity: usize) -> Self { let capacity_bytes = max_buffer_size(encoding, max_level, capacity); let mut buffer = Vec::with_capacity(capacity_bytes); @@ -76,6 +78,7 @@ impl LevelEncoder { /// Creates new level encoder based on RLE encoding. Used to encode Data Page v2 /// repetition and definition levels. + #[allow(dead_code)] pub fn v2(max_level: i16, capacity: usize) -> Self { let capacity_bytes = max_buffer_size(Encoding::RLE, max_level, capacity); let buffer = Vec::with_capacity(capacity_bytes); @@ -83,9 +86,44 @@ impl LevelEncoder { LevelEncoder::RleV2(RleEncoder::new_from_buf(bit_width, buffer)) } + /// Creates a new streaming level encoder for Data Page v1. + /// + /// Unlike [`v1`](Self::v1), this does not require knowing the number of values + /// upfront, making it suitable for incremental encoding where levels are fed in + /// as they arrive via [`put`](Self::put). + pub fn v1_streaming(encoding: Encoding, max_level: i16) -> Self { + let bit_width = num_required_bits(max_level as u64); + match encoding { + Encoding::RLE => { + // Reserve space for length header + let buffer = vec![0u8; 4]; + LevelEncoder::Rle(RleEncoder::new_from_buf(bit_width, buffer)) + } + #[allow(deprecated)] + Encoding::BIT_PACKED => { + LevelEncoder::BitPacked(bit_width, BitWriter::new_from_buf(Vec::new())) + } + _ => panic!("Unsupported encoding type {encoding}"), + } + } + + /// Creates a new streaming RLE level encoder for Data Page v2. + /// + /// Unlike [`v2`](Self::v2), this does not require knowing the number of values + /// upfront, making it suitable for incremental encoding where levels are fed in + /// as they arrive via [`put`](Self::put). + pub fn v2_streaming(max_level: i16) -> Self { + let bit_width = num_required_bits(max_level as u64); + LevelEncoder::RleV2(RleEncoder::new_from_buf(bit_width, Vec::new())) + } + /// Put/encode levels vector into this level encoder. /// Returns number of encoded values that are less than or equal to length of the /// input buffer. + /// + /// This method does **not** flush the underlying encoder, so it can be called + /// incrementally across multiple batches without forcing run boundaries. + /// The encoder is flushed automatically when [`consume`](Self::consume) is called. #[inline] pub fn put(&mut self, buffer: &[i16]) -> usize { let mut num_encoded = 0; @@ -95,14 +133,12 @@ impl LevelEncoder { encoder.put(*value as u64); num_encoded += 1; } - encoder.flush(); } LevelEncoder::BitPacked(bit_width, ref mut encoder) => { for value in buffer { encoder.put_value(*value as u64, bit_width as usize); num_encoded += 1; } - encoder.flush(); } } num_encoded From 10b49911afb6324ef5f04c58e309827da2214329 Mon Sep 17 00:00:00 2001 From: Hippolyte Barraud Date: Fri, 20 Feb 2026 11:11:45 -0500 Subject: [PATCH 2/5] inline level encoder initialization in column writer Signed-off-by: Hippolyte Barraud --- parquet/src/column/writer/mod.rs | 7 ++----- 1 file changed, 2 insertions(+), 5 deletions(-) diff --git a/parquet/src/column/writer/mod.rs b/parquet/src/column/writer/mod.rs index f74211deccf3..53daf1e1aa8b 100644 --- a/parquet/src/column/writer/mod.rs +++ b/parquet/src/column/writer/mod.rs @@ -402,10 +402,9 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { _ => None, }; - let def_levels_encoder = Self::create_level_encoder(descr.max_def_level(), &props); - let rep_levels_encoder = Self::create_level_encoder(descr.max_rep_level(), &props); - Self { + def_levels_encoder: Self::create_level_encoder(descr.max_def_level(), &props), + rep_levels_encoder: Self::create_level_encoder(descr.max_rep_level(), &props), descr, props, statistics_enabled, @@ -413,8 +412,6 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { codec, compressor, encoder, - def_levels_encoder, - rep_levels_encoder, data_pages: VecDeque::new(), page_metrics, column_metrics, From 0e28fba0ed81492f0535d9a7a1baa0c22140a6b5 Mon Sep 17 00:00:00 2001 From: Hippolyte Barraud Date: Wed, 25 Mar 2026 15:28:17 -0400 Subject: [PATCH 3/5] parquet: reuse level encoder allocations across page flushes Previously, flushing a data page would `mem::replace` each level encoder with a freshly allocated one, consuming the old encoder to get its buffer. This allocated new internal `Vec`s on every page boundary. We now preserve the internal state of the encoder and reuse memory across pages. Signed-off-by: Hippolyte Barraud --- parquet/src/column/writer/mod.rs | 35 +++++++++----------------------- parquet/src/encodings/levels.rs | 31 ++++++++++++++++++++++++++++ parquet/src/encodings/rle.rs | 17 ++++++++++++++-- parquet/src/util/bit_util.rs | 7 +++++++ 4 files changed, 63 insertions(+), 27 deletions(-) diff --git a/parquet/src/column/writer/mod.rs b/parquet/src/column/writer/mod.rs index 53daf1e1aa8b..eaf0a6f74816 100644 --- a/parquet/src/column/writer/mod.rs +++ b/parquet/src/column/writer/mod.rs @@ -24,7 +24,6 @@ use crate::bloom_filter::Sbbf; use crate::file::page_index::column_index::ColumnIndexMetaData; use crate::file::page_index::offset_index::OffsetIndexMetaData; use std::collections::{BTreeSet, VecDeque}; -use std::mem; use std::str; use crate::basic::{ @@ -1060,19 +1059,13 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { let mut buffer = vec![]; if max_rep_level > 0 { - let encoder = mem::replace( - &mut self.rep_levels_encoder, - Self::create_level_encoder(max_rep_level, &self.props), - ); - buffer.extend_from_slice(&encoder.consume()); + self.rep_levels_encoder + .flush_to(|data| buffer.extend_from_slice(data)); } if max_def_level > 0 { - let encoder = mem::replace( - &mut self.def_levels_encoder, - Self::create_level_encoder(max_def_level, &self.props), - ); - buffer.extend_from_slice(&encoder.consume()); + self.def_levels_encoder + .flush_to(|data| buffer.extend_from_slice(data)); } buffer.extend_from_slice(&values_data.buf); @@ -1102,23 +1095,15 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { let mut buffer = vec![]; if max_rep_level > 0 { - let encoder = mem::replace( - &mut self.rep_levels_encoder, - Self::create_level_encoder(max_rep_level, &self.props), - ); - let levels = encoder.consume(); - rep_levels_byte_len = levels.len(); - buffer.extend_from_slice(&levels); + self.rep_levels_encoder + .flush_to(|data| buffer.extend_from_slice(data)); + rep_levels_byte_len = buffer.len(); } if max_def_level > 0 { - let encoder = mem::replace( - &mut self.def_levels_encoder, - Self::create_level_encoder(max_def_level, &self.props), - ); - let levels = encoder.consume(); - def_levels_byte_len = levels.len(); - buffer.extend_from_slice(&levels); + self.def_levels_encoder + .flush_to(|data| buffer.extend_from_slice(data)); + def_levels_byte_len = buffer.len() - rep_levels_byte_len; } let uncompressed_size = diff --git a/parquet/src/encodings/levels.rs b/parquet/src/encodings/levels.rs index c1c687fb2c8c..02b758e2c2c5 100644 --- a/parquet/src/encodings/levels.rs +++ b/parquet/src/encodings/levels.rs @@ -147,6 +147,7 @@ impl LevelEncoder { /// Finalizes level encoder, flush all intermediate buffers and return resulting /// encoded buffer. Returned buffer is already truncated to encoded bytes only. #[inline] + #[allow(unused)] pub fn consume(self) -> Vec { match self { LevelEncoder::Rle(encoder) => { @@ -162,4 +163,34 @@ impl LevelEncoder { LevelEncoder::BitPacked(_, encoder) => encoder.consume(), } } + + /// Flushes all intermediate buffers, passes the encoded data to `f`, then + /// resets the encoder for reuse while retaining the buffer allocation. + #[inline] + pub fn flush_to(&mut self, f: F) -> R + where + F: FnOnce(&[u8]) -> R, + { + let result = match self { + LevelEncoder::Rle(encoder) => { + let data = encoder.flush_buffer_mut(); + // Patch the 4-byte length header reserved at the start of the buffer + let encoded_len = (data.len() - mem::size_of::()) as i32; + data[..4].copy_from_slice(&encoded_len.to_le_bytes()); + f(data) + } + LevelEncoder::RleV2(encoder) => f(encoder.flush_buffer()), + LevelEncoder::BitPacked(_, encoder) => f(encoder.flush_buffer()), + }; + match self { + LevelEncoder::Rle(encoder) => { + encoder.clear(); + // Re-reserve the 4-byte length header for the next page + encoder.skip(mem::size_of::()); + } + LevelEncoder::RleV2(encoder) => encoder.clear(), + LevelEncoder::BitPacked(_, encoder) => encoder.clear(), + } + result + } } diff --git a/parquet/src/encodings/rle.rs b/parquet/src/encodings/rle.rs index c95a46c634d2..2815c20dab56 100644 --- a/parquet/src/encodings/rle.rs +++ b/parquet/src/encodings/rle.rs @@ -177,16 +177,22 @@ impl RleEncoder { /// Borrow equivalent of the `consume` method. /// Call `clear()` after invoking this method. #[inline] - #[allow(unused)] pub fn flush_buffer(&mut self) -> &[u8] { self.flush(); self.bit_writer.flush_buffer() } + /// Like `flush_buffer`, but returns mutable access to the internal buffer. + /// Call `clear()` after invoking this method. + #[inline] + pub fn flush_buffer_mut(&mut self) -> &mut [u8] { + self.flush(); + self.bit_writer.flush_buffer_mut() + } + /// Clears the internal state so this encoder can be reused (e.g., after becoming /// full). #[inline] - #[allow(unused)] pub fn clear(&mut self) { self.bit_writer.clear(); self.num_buffered_values = 0; @@ -196,6 +202,13 @@ impl RleEncoder { self.indicator_byte_pos = -1; } + /// Advances the buffer by `num_bytes` zero bytes, delegating to the + /// underlying [`BitWriter::skip`]. + #[inline] + pub fn skip(&mut self, num_bytes: usize) { + self.bit_writer.skip(num_bytes); + } + /// Flushes all remaining values and return the final byte buffer maintained by the /// internal writer. #[inline] diff --git a/parquet/src/util/bit_util.rs b/parquet/src/util/bit_util.rs index 3a26603fabc4..1c5ec19517bd 100644 --- a/parquet/src/util/bit_util.rs +++ b/parquet/src/util/bit_util.rs @@ -219,6 +219,13 @@ impl BitWriter { self.buffer() } + /// Like `flush_buffer`, but returns mutable access to the buffer. + #[inline] + pub fn flush_buffer_mut(&mut self) -> &mut [u8] { + self.flush(); + &mut self.buffer + } + /// Clears the internal state so the buffer can be reused. #[inline] pub fn clear(&mut self) { From f2767de5f8a62a3b8521c4efb53b1ec82cae5f1b Mon Sep 17 00:00:00 2001 From: Hippolyte Barraud Date: Tue, 31 Mar 2026 22:43:46 -0400 Subject: [PATCH 4/5] parquet: remove pre-allocated level encoders The `v1`, `v2`, and `max_buffer_size` functions required knowing the number of values upfront and pre-allocated buffers. All callers have been migrated to the streaming variants (`v1_streaming`, `v2_streaming`), so remove the dead code and switch the last remaining caller in test utils. Signed-off-by: Hippolyte Barraud --- parquet/src/encodings/levels.rs | 54 ----------------------- parquet/src/util/test_common/page_util.rs | 2 +- 2 files changed, 1 insertion(+), 55 deletions(-) diff --git a/parquet/src/encodings/levels.rs b/parquet/src/encodings/levels.rs index 02b758e2c2c5..288c3d19edae 100644 --- a/parquet/src/encodings/levels.rs +++ b/parquet/src/encodings/levels.rs @@ -23,21 +23,6 @@ use crate::basic::Encoding; use crate::data_type::AsBytes; use crate::util::bit_util::{BitWriter, ceil, num_required_bits}; -/// Computes max buffer size for level encoder/decoder based on encoding, max -/// repetition/definition level and number of total buffered values (includes null -/// values). -#[allow(dead_code)] -#[inline] -pub fn max_buffer_size(encoding: Encoding, max_level: i16, num_buffered_values: usize) -> usize { - let bit_width = num_required_bits(max_level as u64); - match encoding { - Encoding::RLE => RleEncoder::max_buffer_size(bit_width, num_buffered_values), - #[allow(deprecated)] - Encoding::BIT_PACKED => ceil(num_buffered_values * bit_width as usize, 8), - _ => panic!("Unsupported encoding type {encoding}"), - } -} - /// Encoder for definition/repetition levels. /// Currently only supports Rle and BitPacked (dev/null) encoding, including v2. pub enum LevelEncoder { @@ -47,45 +32,6 @@ pub enum LevelEncoder { } impl LevelEncoder { - /// Creates new level encoder based on encoding, max level and underlying byte buffer. - /// For bit packed encoding it is assumed that buffer is already allocated with - /// `levels::max_buffer_size` method. - /// - /// Used to encode levels for Data Page v1. - /// - /// Panics, if encoding is not supported. - #[allow(dead_code)] - pub fn v1(encoding: Encoding, max_level: i16, capacity: usize) -> Self { - let capacity_bytes = max_buffer_size(encoding, max_level, capacity); - let mut buffer = Vec::with_capacity(capacity_bytes); - let bit_width = num_required_bits(max_level as u64); - match encoding { - Encoding::RLE => { - // Reserve space for length header - buffer.extend_from_slice(&[0; 4]); - LevelEncoder::Rle(RleEncoder::new_from_buf(bit_width, buffer)) - } - #[allow(deprecated)] - Encoding::BIT_PACKED => { - // Here we set full byte buffer without adjusting for num_buffered_values, - // because byte buffer will already be allocated with size from - // `max_buffer_size()` method. - LevelEncoder::BitPacked(bit_width, BitWriter::new_from_buf(buffer)) - } - _ => panic!("Unsupported encoding type {encoding}"), - } - } - - /// Creates new level encoder based on RLE encoding. Used to encode Data Page v2 - /// repetition and definition levels. - #[allow(dead_code)] - pub fn v2(max_level: i16, capacity: usize) -> Self { - let capacity_bytes = max_buffer_size(Encoding::RLE, max_level, capacity); - let buffer = Vec::with_capacity(capacity_bytes); - let bit_width = num_required_bits(max_level as u64); - LevelEncoder::RleV2(RleEncoder::new_from_buf(bit_width, buffer)) - } - /// Creates a new streaming level encoder for Data Page v1. /// /// Unlike [`v1`](Self::v1), this does not require knowing the number of values diff --git a/parquet/src/util/test_common/page_util.rs b/parquet/src/util/test_common/page_util.rs index 5b64eb54133c..6a99beaea151 100644 --- a/parquet/src/util/test_common/page_util.rs +++ b/parquet/src/util/test_common/page_util.rs @@ -75,7 +75,7 @@ impl DataPageBuilderImpl { if max_level <= 0 { return 0; } - let mut level_encoder = LevelEncoder::v1(Encoding::RLE, max_level, levels.len()); + let mut level_encoder = LevelEncoder::v1_streaming(Encoding::RLE, max_level); level_encoder.put(levels); let encoded_levels = level_encoder.consume(); // Actual encoded bytes (without length offset) From d55ca607945662d8d47e500ddf1bf3d750ef4a1e Mon Sep 17 00:00:00 2001 From: Andrew Lamb Date: Wed, 1 Apr 2026 12:17:46 -0400 Subject: [PATCH 5/5] Fix clippy/compile --- parquet/src/column/reader.rs | 4 ++-- parquet/src/encodings/levels.rs | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/parquet/src/column/reader.rs b/parquet/src/column/reader.rs index 29cb50185a58..d1a328237f26 100644 --- a/parquet/src/column/reader.rs +++ b/parquet/src/column/reader.rs @@ -1401,11 +1401,11 @@ mod tests { // Helper: build a DataPage v2 for this list column. let make_v2_page = |rep_levels: &[i16], def_levels: &[i16], values: &[i32], num_rows: u32| -> Page { - let mut rep_enc = LevelEncoder::v2(max_rep_level, rep_levels.len()); + let mut rep_enc = LevelEncoder::v2_streaming(max_rep_level); rep_enc.put(rep_levels); let rep_bytes = rep_enc.consume(); - let mut def_enc = LevelEncoder::v2(max_def_level, def_levels.len()); + let mut def_enc = LevelEncoder::v2_streaming(max_def_level); def_enc.put(def_levels); let def_bytes = def_enc.consume(); diff --git a/parquet/src/encodings/levels.rs b/parquet/src/encodings/levels.rs index 288c3d19edae..b761a5ac5daf 100644 --- a/parquet/src/encodings/levels.rs +++ b/parquet/src/encodings/levels.rs @@ -21,7 +21,7 @@ use super::rle::RleEncoder; use crate::basic::Encoding; use crate::data_type::AsBytes; -use crate::util::bit_util::{BitWriter, ceil, num_required_bits}; +use crate::util::bit_util::{BitWriter, num_required_bits}; /// Encoder for definition/repetition levels. /// Currently only supports Rle and BitPacked (dev/null) encoding, including v2.