From c0971de938f2f9695c88165224d3467b8205bdf3 Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Wed, 22 Jul 2026 16:26:38 +0200 Subject: [PATCH 01/11] remove spurious physical type checks --- parquet/src/encodings/encoding/mod.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/parquet/src/encodings/encoding/mod.rs b/parquet/src/encodings/encoding/mod.rs index eeabcf4ba5ce..88e5e2b00fd0 100644 --- a/parquet/src/encodings/encoding/mod.rs +++ b/parquet/src/encodings/encoding/mod.rs @@ -594,7 +594,7 @@ impl DeltaLengthByteArrayEncoder { impl Encoder for DeltaLengthByteArrayEncoder { fn put(&mut self, values: &[T::T]) -> Result<()> { ensure_phys_ty!( - Type::BYTE_ARRAY | Type::FIXED_LEN_BYTE_ARRAY, + Type::BYTE_ARRAY, "DeltaLengthByteArrayEncoder only supports ByteArrayType" ); @@ -628,7 +628,7 @@ impl Encoder for DeltaLengthByteArrayEncoder { fn flush_buffer(&mut self) -> Result { ensure_phys_ty!( - Type::BYTE_ARRAY | Type::FIXED_LEN_BYTE_ARRAY, + Type::BYTE_ARRAY, "DeltaLengthByteArrayEncoder only supports ByteArrayType" ); From 85c86c5af5b36b9af9a240c6277613ac45ea4c8a Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Wed, 22 Jul 2026 16:27:16 +0200 Subject: [PATCH 02/11] prealloc tmp buffers --- parquet/src/encodings/encoding/mod.rs | 39 +++++++++++++-------------- 1 file changed, 18 insertions(+), 21 deletions(-) diff --git a/parquet/src/encodings/encoding/mod.rs b/parquet/src/encodings/encoding/mod.rs index 88e5e2b00fd0..6b09a13701c0 100644 --- a/parquet/src/encodings/encoding/mod.rs +++ b/parquet/src/encodings/encoding/mod.rs @@ -598,18 +598,18 @@ impl Encoder for DeltaLengthByteArrayEncoder { "DeltaLengthByteArrayEncoder only supports ByteArrayType" ); - let val_it = || { - values - .iter() - .map(|x| x.as_any().downcast_ref::().unwrap()) - }; + let values = values + .iter() + .map(|x| x.as_any().downcast_ref::().unwrap()); - let lengths: Vec = val_it().map(|byte_array| byte_array.len() as i32).collect(); - self.len_encoder.put(&lengths)?; - for byte_array in val_it() { - self.encoded_size += byte_array.len(); + let mut lengths: Vec = Vec::with_capacity(values.len()); + for byte_array in values { + let len = byte_array.len(); + lengths.push(len as i32); + self.encoded_size += len; self.data.push(byte_array.clone()); } + self.len_encoder.put(&lengths)?; Ok(()) } @@ -682,19 +682,16 @@ impl DeltaByteArrayEncoder { impl Encoder for DeltaByteArrayEncoder { fn put(&mut self, values: &[T::T]) -> Result<()> { - let mut prefix_lengths: Vec = vec![]; - let mut suffixes: Vec = vec![]; + let mut prefix_lengths: Vec = Vec::with_capacity(values.len()); + let mut suffixes: Vec = Vec::with_capacity(values.len()); - let values = values - .iter() - .map(|x| x.as_any()) - .map(|x| match T::get_physical_type() { - Type::BYTE_ARRAY => x.downcast_ref::().unwrap(), - Type::FIXED_LEN_BYTE_ARRAY => x.downcast_ref::().unwrap(), - _ => panic!( - "DeltaByteArrayEncoder only supports ByteArrayType and FixedLenByteArrayType" - ), - }); + let values = values.iter().map(|x| match T::get_physical_type() { + Type::BYTE_ARRAY => x.as_any().downcast_ref::().unwrap(), + Type::FIXED_LEN_BYTE_ARRAY => x.as_any().downcast_ref::().unwrap(), + _ => panic!( + "DeltaByteArrayEncoder only supports ByteArrayType and FixedLenByteArrayType" + ), + }); for byte_array in values { let current = byte_array.data(); From 429ecda02c93cc2f63220dc7de8211b27ae17db3 Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Wed, 22 Jul 2026 17:05:20 +0200 Subject: [PATCH 03/11] add Encoder::flush_to --- .../encoding/byte_stream_split_encoder.rs | 40 ++++++------ .../src/encodings/encoding/dict_encoder.rs | 6 ++ parquet/src/encodings/encoding/mod.rs | 61 ++++++++++++------- 3 files changed, 68 insertions(+), 39 deletions(-) diff --git a/parquet/src/encodings/encoding/byte_stream_split_encoder.rs b/parquet/src/encodings/encoding/byte_stream_split_encoder.rs index 0726c6b1c919..3fa0ddbb9a9c 100644 --- a/parquet/src/encodings/encoding/byte_stream_split_encoder.rs +++ b/parquet/src/encodings/encoding/byte_stream_split_encoder.rs @@ -22,7 +22,7 @@ use crate::errors::{ParquetError, Result}; use super::Encoder; -use bytes::{BufMut, Bytes}; +use bytes::BufMut; use std::cmp; use std::marker::PhantomData; @@ -88,12 +88,15 @@ impl Encoder for ByteStreamSplitEncoder { self.buffer.len() } - fn flush_buffer(&mut self) -> Result { - let mut encoded = vec![0; self.buffer.len()]; + fn flush_to(&mut self, out: &mut Vec) -> Result<()> { + let start = out.len(); + out.resize(start + self.buffer.len(), 0); + let encoded = &mut out[start..]; + let type_size = T::get_type_size(); match type_size { - 4 => split_streams_const::<4>(&self.buffer, &mut encoded), - 8 => split_streams_const::<8>(&self.buffer, &mut encoded), + 4 => split_streams_const::<4>(&self.buffer, encoded), + 8 => split_streams_const::<8>(&self.buffer, encoded), _ => { return Err(general_err!( "byte stream split unsupported for data types of size {} bytes", @@ -103,7 +106,7 @@ impl Encoder for ByteStreamSplitEncoder { } self.buffer.clear(); - Ok(encoded.into()) + Ok(()) } /// return the estimated memory size of this encoder. @@ -202,26 +205,29 @@ impl Encoder for VariableWidthByteStreamSplitEncoder { self.buffer.len() } - fn flush_buffer(&mut self) -> Result { - let mut encoded = vec![0; self.buffer.len()]; + fn flush_to(&mut self, out: &mut Vec) -> Result<()> { + let start = out.len(); + out.resize(start + self.buffer.len(), 0); + let encoded = &mut out[start..]; + let type_size = match T::get_physical_type() { Type::FIXED_LEN_BYTE_ARRAY => self.type_width, _ => T::get_type_size(), }; // split_streams_const() is faster up to type_width == 8 match type_size { - 2 => split_streams_const::<2>(&self.buffer, &mut encoded), - 3 => split_streams_const::<3>(&self.buffer, &mut encoded), - 4 => split_streams_const::<4>(&self.buffer, &mut encoded), - 5 => split_streams_const::<5>(&self.buffer, &mut encoded), - 6 => split_streams_const::<6>(&self.buffer, &mut encoded), - 7 => split_streams_const::<7>(&self.buffer, &mut encoded), - 8 => split_streams_const::<8>(&self.buffer, &mut encoded), - _ => split_streams_variable(&self.buffer, &mut encoded, type_size), + 2 => split_streams_const::<2>(&self.buffer, encoded), + 3 => split_streams_const::<3>(&self.buffer, encoded), + 4 => split_streams_const::<4>(&self.buffer, encoded), + 5 => split_streams_const::<5>(&self.buffer, encoded), + 6 => split_streams_const::<6>(&self.buffer, encoded), + 7 => split_streams_const::<7>(&self.buffer, encoded), + 8 => split_streams_const::<8>(&self.buffer, encoded), + _ => split_streams_variable(&self.buffer, out, type_size), } self.buffer.clear(); - Ok(encoded.into()) + Ok(()) } /// return the estimated memory size of this encoder. diff --git a/parquet/src/encodings/encoding/dict_encoder.rs b/parquet/src/encodings/encoding/dict_encoder.rs index 89666bbe7313..cd89a3c8db96 100644 --- a/parquet/src/encodings/encoding/dict_encoder.rs +++ b/parquet/src/encodings/encoding/dict_encoder.rs @@ -184,6 +184,12 @@ impl Encoder for DictEncoder { self.write_indices() } + fn flush_to(&mut self, out: &mut Vec) -> Result<()> { + // TODO: RleEncoder could be reused instead of consumed + out.extend_from_slice(&self.flush_buffer()?); + Ok(()) + } + /// Returns the estimated total memory usage /// /// For this encoder, the indices are unencoded bytes (refer to [`Self::write_indices`]). diff --git a/parquet/src/encodings/encoding/mod.rs b/parquet/src/encodings/encoding/mod.rs index 6b09a13701c0..37e68e5930bc 100644 --- a/parquet/src/encodings/encoding/mod.rs +++ b/parquet/src/encodings/encoding/mod.rs @@ -75,7 +75,15 @@ pub trait Encoder: Send { /// Flushes the underlying byte buffer that's being processed by this encoder, and /// return the immutable copy of it. This will also reset the internal state. - fn flush_buffer(&mut self) -> Result; + fn flush_buffer(&mut self) -> Result { + let mut buffer = Vec::new(); + self.flush_to(&mut buffer)?; + Ok(buffer.into()) + } + + /// Flushes the underlying byte buffer that's being processed by this encoder. + /// This will also reset the internal state. + fn flush_to(&mut self, out: &mut Vec) -> Result<()>; } /// Gets a encoder for the particular data type `T` and encoding `encoding`. Memory usage @@ -164,6 +172,15 @@ impl Encoder for PlainEncoder { Ok(std::mem::take(&mut self.buffer).into()) } + #[inline] + fn flush_to(&mut self, out: &mut Vec) -> Result<()> { + out.extend_from_slice(&self.buffer); + out.extend_from_slice(self.bit_writer.flush_buffer()); + self.buffer.clear(); + self.bit_writer.clear(); + Ok(()) + } + #[inline] fn put(&mut self, values: &[T::T]) -> Result<()> { T::T::encode(values, &mut self.buffer, &mut self.bit_writer)?; @@ -260,6 +277,12 @@ impl Encoder for RleValueEncoder { Ok(buf.into()) } + fn flush_to(&mut self, out: &mut Vec) -> Result<()> { + // TODO: RleEncoder could be reused instead of consumed + out.extend_from_slice(&self.flush_buffer()?); + Ok(()) + } + /// return the estimated memory size of this encoder. fn estimated_memory_size(&self) -> usize { self.encoder @@ -477,18 +500,17 @@ impl Encoder for DeltaBitPackEncoder { } fn estimated_data_encoded_size(&self) -> usize { - self.bit_writer.bytes_written() + self.page_header_writer.bytes_written() + self.bit_writer.bytes_written() } - fn flush_buffer(&mut self) -> Result { + fn flush_to(&mut self, out: &mut Vec) -> Result<()> { // Write remaining values self.flush_block_values()?; // Write page header with total values self.write_page_header(); - let mut buffer = Vec::new(); - buffer.extend_from_slice(self.page_header_writer.flush_buffer()); - buffer.extend_from_slice(self.bit_writer.flush_buffer()); + out.extend_from_slice(self.page_header_writer.flush_buffer()); + out.extend_from_slice(self.bit_writer.flush_buffer()); // Reset state self.page_header_writer.clear(); @@ -498,7 +520,7 @@ impl Encoder for DeltaBitPackEncoder { self.current_value = 0; self.values_in_block = 0; - Ok(buffer.into()) + Ok(()) } /// return the estimated memory size of this encoder. @@ -626,22 +648,22 @@ impl Encoder for DeltaLengthByteArrayEncoder { self.len_encoder.estimated_data_encoded_size() + self.encoded_size } - fn flush_buffer(&mut self) -> Result { + fn flush_to(&mut self, out: &mut Vec) -> Result<()> { ensure_phys_ty!( Type::BYTE_ARRAY, "DeltaLengthByteArrayEncoder only supports ByteArrayType" ); - let mut total_bytes = vec![]; - let lengths = self.len_encoder.flush_buffer()?; - total_bytes.extend_from_slice(&lengths); + out.reserve(self.estimated_data_encoded_size()); + self.len_encoder.flush_to(out)?; + self.data.iter().for_each(|byte_array| { - total_bytes.extend_from_slice(byte_array.data()); + out.extend_from_slice(byte_array.data()); }); self.data.clear(); self.encoded_size = 0; - Ok(total_bytes.into()) + Ok(()) } /// return the estimated memory size of this encoder. @@ -727,21 +749,16 @@ impl Encoder for DeltaByteArrayEncoder { + self.suffix_writer.estimated_data_encoded_size() } - fn flush_buffer(&mut self) -> Result { + fn flush_to(&mut self, out: &mut Vec) -> Result<()> { match T::get_physical_type() { Type::BYTE_ARRAY | Type::FIXED_LEN_BYTE_ARRAY => { - // TODO: investigate if we can merge lengths and suffixes - // without copying data into new vector. - let mut total_bytes = vec![]; // Insert lengths ... - let lengths = self.prefix_len_encoder.flush_buffer()?; - total_bytes.extend_from_slice(&lengths); + self.prefix_len_encoder.flush_to(out)?; // ... followed by suffixes - let suffixes = self.suffix_writer.flush_buffer()?; - total_bytes.extend_from_slice(&suffixes); + self.suffix_writer.flush_to(out)?; self.previous.clear(); - Ok(total_bytes.into()) + Ok(()) } _ => panic!( "DeltaByteArrayEncoder only supports ByteArrayType and FixedLenByteArrayType" From cb3755d1fbf1279b4ba65508db8650cacbb0f61b Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Wed, 22 Jul 2026 17:13:21 +0200 Subject: [PATCH 04/11] remove excessive [cold] attribs --- .../src/encodings/encoding/dict_encoder.rs | 3 --- parquet/src/encodings/encoding/mod.rs | 20 ------------------- 2 files changed, 23 deletions(-) diff --git a/parquet/src/encodings/encoding/dict_encoder.rs b/parquet/src/encodings/encoding/dict_encoder.rs index cd89a3c8db96..2feb77644101 100644 --- a/parquet/src/encodings/encoding/dict_encoder.rs +++ b/parquet/src/encodings/encoding/dict_encoder.rs @@ -164,9 +164,6 @@ impl Encoder for DictEncoder { Ok(()) } - // Performance Note: - // As far as can be seen these functions are rarely called and as such we can hint to the - // compiler that they dont need to be folded into hot locations in the final output. fn encoding(&self) -> Encoding { Encoding::PLAIN_DICTIONARY } diff --git a/parquet/src/encodings/encoding/mod.rs b/parquet/src/encodings/encoding/mod.rs index 37e68e5930bc..ef37c6b852af 100644 --- a/parquet/src/encodings/encoding/mod.rs +++ b/parquet/src/encodings/encoding/mod.rs @@ -152,10 +152,6 @@ impl PlainEncoder { } impl Encoder for PlainEncoder { - // Performance Note: - // As far as can be seen these functions are rarely called and as such we can hint to the - // compiler that they dont need to be folded into hot locations in the final output. - #[cold] fn encoding(&self) -> Encoding { Encoding::PLAIN } @@ -242,10 +238,6 @@ impl Encoder for RleValueEncoder { Ok(()) } - // Performance Note: - // As far as can be seen these functions are rarely called and as such we can hint to the - // compiler that they dont need to be folded into hot locations in the final output. - #[cold] fn encoding(&self) -> Encoding { Encoding::RLE } @@ -491,10 +483,6 @@ impl Encoder for DeltaBitPackEncoder { Ok(()) } - // Performance Note: - // As far as can be seen these functions are rarely called and as such we can hint to the - // compiler that they dont need to be folded into hot locations in the final output. - #[cold] fn encoding(&self) -> Encoding { Encoding::DELTA_BINARY_PACKED } @@ -636,10 +624,6 @@ impl Encoder for DeltaLengthByteArrayEncoder { Ok(()) } - // Performance Note: - // As far as can be seen these functions are rarely called and as such we can hint to the - // compiler that they dont need to be folded into hot locations in the final output. - #[cold] fn encoding(&self) -> Encoding { Encoding::DELTA_LENGTH_BYTE_ARRAY } @@ -736,10 +720,6 @@ impl Encoder for DeltaByteArrayEncoder { Ok(()) } - // Performance Note: - // As far as can be seen these functions are rarely called and as such we can hint to the - // compiler that they dont need to be folded into hot locations in the final output. - #[cold] fn encoding(&self) -> Encoding { Encoding::DELTA_BYTE_ARRAY } From a4e23a1da1ae872a308a8251c6941851ae95d91c Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Wed, 22 Jul 2026 17:19:01 +0200 Subject: [PATCH 05/11] reuse length buffer on encoder --- parquet/src/encodings/encoding/mod.rs | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/parquet/src/encodings/encoding/mod.rs b/parquet/src/encodings/encoding/mod.rs index ef37c6b852af..46719baded3f 100644 --- a/parquet/src/encodings/encoding/mod.rs +++ b/parquet/src/encodings/encoding/mod.rs @@ -576,6 +576,8 @@ impl DeltaBitPackEncoderConversion for DeltaBitPackEncoder { pub struct DeltaLengthByteArrayEncoder { // length encoder len_encoder: DeltaBitPackEncoder, + // length buffer + lengths: Vec, // byte array data data: Vec, // data size in bytes of encoded values @@ -594,6 +596,7 @@ impl DeltaLengthByteArrayEncoder { pub fn new() -> Self { Self { len_encoder: DeltaBitPackEncoder::new(), + lengths: vec![], data: vec![], encoded_size: 0, _phantom: PhantomData, @@ -612,14 +615,15 @@ impl Encoder for DeltaLengthByteArrayEncoder { .iter() .map(|x| x.as_any().downcast_ref::().unwrap()); - let mut lengths: Vec = Vec::with_capacity(values.len()); + self.lengths.reserve(values.len()); for byte_array in values { let len = byte_array.len(); - lengths.push(len as i32); + self.lengths.push(len as i32); self.encoded_size += len; self.data.push(byte_array.clone()); } - self.len_encoder.put(&lengths)?; + self.len_encoder.put(&self.lengths)?; + self.lengths.clear(); Ok(()) } From f72fadefdae92398fda5e6e3c7e20e59d9647be9 Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Wed, 22 Jul 2026 17:33:48 +0200 Subject: [PATCH 06/11] reuse+clone previous array instead of copying --- parquet/src/encodings/encoding/mod.rs | 28 ++++++++++++++------------- 1 file changed, 15 insertions(+), 13 deletions(-) diff --git a/parquet/src/encodings/encoding/mod.rs b/parquet/src/encodings/encoding/mod.rs index 46719baded3f..aa7cdf6c3796 100644 --- a/parquet/src/encodings/encoding/mod.rs +++ b/parquet/src/encodings/encoding/mod.rs @@ -668,7 +668,7 @@ impl Encoder for DeltaLengthByteArrayEncoder { pub struct DeltaByteArrayEncoder { prefix_len_encoder: DeltaBitPackEncoder, suffix_writer: DeltaLengthByteArrayEncoder, - previous: Vec, + previous: ByteArray, _phantom: PhantomData, } @@ -684,7 +684,7 @@ impl DeltaByteArrayEncoder { Self { prefix_len_encoder: DeltaBitPackEncoder::new(), suffix_writer: DeltaLengthByteArrayEncoder::new(), - previous: vec![], + previous: ByteArray::new(), _phantom: PhantomData, } } @@ -703,21 +703,23 @@ impl Encoder for DeltaByteArrayEncoder { ), }); - for byte_array in values { - let current = byte_array.data(); - // Maximum prefix length that is shared between previous value and current - // value - let prefix_len = cmp::min(self.previous.len(), current.len()); + let mut previous_array = &self.previous; + for current_array in values { + let current = current_array.data(); + let previous = previous_array.data(); + // Maximum prefix length that is shared between previous value and current value + let prefix_len = cmp::min(previous.len(), current.len()); let mut match_len = 0; - while match_len < prefix_len && self.previous[match_len] == current[match_len] { + while match_len < prefix_len && previous[match_len] == current[match_len] { match_len += 1; } prefix_lengths.push(match_len as i32); - suffixes.push(byte_array.slice(match_len, byte_array.len() - match_len)); + suffixes.push(current_array.slice(match_len, current_array.len() - match_len)); // Update previous for the next prefix - self.previous.clear(); - self.previous.extend_from_slice(current); + previous_array = current_array; } + self.previous = previous_array.clone(); + self.prefix_len_encoder.put(&prefix_lengths)?; self.suffix_writer.put(&suffixes)?; @@ -741,7 +743,7 @@ impl Encoder for DeltaByteArrayEncoder { // ... followed by suffixes self.suffix_writer.flush_to(out)?; - self.previous.clear(); + self.previous = ByteArray::new(); Ok(()) } _ => panic!( @@ -754,7 +756,7 @@ impl Encoder for DeltaByteArrayEncoder { fn estimated_memory_size(&self) -> usize { self.prefix_len_encoder.estimated_memory_size() + self.suffix_writer.estimated_memory_size() - + (self.previous.capacity() * std::mem::size_of::()) + + (self.previous.len() * std::mem::size_of::()) } } From 20b8ff6021c1fae57d1d3f14383675723fdb968c Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Wed, 22 Jul 2026 17:43:52 +0200 Subject: [PATCH 07/11] reuse buffers in DeltaByteArrayEncoder --- .../src/encodings/encoding/dict_encoder.rs | 2 ++ parquet/src/encodings/encoding/mod.rs | 24 ++++++++++++------- 2 files changed, 18 insertions(+), 8 deletions(-) diff --git a/parquet/src/encodings/encoding/dict_encoder.rs b/parquet/src/encodings/encoding/dict_encoder.rs index 2feb77644101..f99fbd9f4bb7 100644 --- a/parquet/src/encodings/encoding/dict_encoder.rs +++ b/parquet/src/encodings/encoding/dict_encoder.rs @@ -132,6 +132,8 @@ impl DictEncoder { /// Writes out the dictionary values with RLE encoding in a byte buffer, and return /// the result. pub fn write_indices(&mut self) -> Result { + // TODO: Move this into flush_buffer/flush_to? + // That could allow reusing buffer with flush_to. let buffer_len = self.estimated_data_encoded_size(); let mut buffer = Vec::with_capacity(buffer_len); buffer.push(self.bit_width()); diff --git a/parquet/src/encodings/encoding/mod.rs b/parquet/src/encodings/encoding/mod.rs index aa7cdf6c3796..8ab5fc7f3dad 100644 --- a/parquet/src/encodings/encoding/mod.rs +++ b/parquet/src/encodings/encoding/mod.rs @@ -81,7 +81,7 @@ pub trait Encoder: Send { Ok(buffer.into()) } - /// Flushes the underlying byte buffer that's being processed by this encoder. + /// Flushes the underlying byte buffer that's being processed by this encoder. /// This will also reset the internal state. fn flush_to(&mut self, out: &mut Vec) -> Result<()>; } @@ -668,6 +668,8 @@ impl Encoder for DeltaLengthByteArrayEncoder { pub struct DeltaByteArrayEncoder { prefix_len_encoder: DeltaBitPackEncoder, suffix_writer: DeltaLengthByteArrayEncoder, + prefix_lengths: Vec, + suffixes: Vec, previous: ByteArray, _phantom: PhantomData, } @@ -684,6 +686,8 @@ impl DeltaByteArrayEncoder { Self { prefix_len_encoder: DeltaBitPackEncoder::new(), suffix_writer: DeltaLengthByteArrayEncoder::new(), + prefix_lengths: Vec::new(), + suffixes: Vec::new(), previous: ByteArray::new(), _phantom: PhantomData, } @@ -692,8 +696,8 @@ impl DeltaByteArrayEncoder { impl Encoder for DeltaByteArrayEncoder { fn put(&mut self, values: &[T::T]) -> Result<()> { - let mut prefix_lengths: Vec = Vec::with_capacity(values.len()); - let mut suffixes: Vec = Vec::with_capacity(values.len()); + self.prefix_lengths.reserve(values.len()); + self.suffixes.reserve(values.len()); let values = values.iter().map(|x| match T::get_physical_type() { Type::BYTE_ARRAY => x.as_any().downcast_ref::().unwrap(), @@ -713,15 +717,19 @@ impl Encoder for DeltaByteArrayEncoder { while match_len < prefix_len && previous[match_len] == current[match_len] { match_len += 1; } - prefix_lengths.push(match_len as i32); - suffixes.push(current_array.slice(match_len, current_array.len() - match_len)); + self.prefix_lengths.push(match_len as i32); + self.suffixes + .push(current_array.slice(match_len, current.len() - match_len)); // Update previous for the next prefix previous_array = current_array; } self.previous = previous_array.clone(); - - self.prefix_len_encoder.put(&prefix_lengths)?; - self.suffix_writer.put(&suffixes)?; + + self.prefix_len_encoder.put(&self.prefix_lengths)?; + self.suffix_writer.put(&self.suffixes)?; + + self.prefix_lengths.clear(); + self.suffixes.clear(); Ok(()) } From 7c1b5333e75fb66ce3312c03a352c32717ad6b17 Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Wed, 22 Jul 2026 17:48:19 +0200 Subject: [PATCH 08/11] fix empty ByteArray panic --- parquet/src/encodings/encoding/mod.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/parquet/src/encodings/encoding/mod.rs b/parquet/src/encodings/encoding/mod.rs index 8ab5fc7f3dad..a696cb27f047 100644 --- a/parquet/src/encodings/encoding/mod.rs +++ b/parquet/src/encodings/encoding/mod.rs @@ -688,7 +688,7 @@ impl DeltaByteArrayEncoder { suffix_writer: DeltaLengthByteArrayEncoder::new(), prefix_lengths: Vec::new(), suffixes: Vec::new(), - previous: ByteArray::new(), + previous: Bytes::new().into(), _phantom: PhantomData, } } @@ -751,7 +751,7 @@ impl Encoder for DeltaByteArrayEncoder { // ... followed by suffixes self.suffix_writer.flush_to(out)?; - self.previous = ByteArray::new(); + self.previous = Bytes::new().into(); Ok(()) } _ => panic!( From 38cad9195f0f325c8adb37a8fc38bdc70f86739d Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Thu, 23 Jul 2026 14:12:13 +0200 Subject: [PATCH 09/11] avoid keeping ByteArray alive in DeltaByteArrayEncoder --- parquet/src/encodings/encoding/mod.rs | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) diff --git a/parquet/src/encodings/encoding/mod.rs b/parquet/src/encodings/encoding/mod.rs index a696cb27f047..00939349a95f 100644 --- a/parquet/src/encodings/encoding/mod.rs +++ b/parquet/src/encodings/encoding/mod.rs @@ -670,7 +670,7 @@ pub struct DeltaByteArrayEncoder { suffix_writer: DeltaLengthByteArrayEncoder, prefix_lengths: Vec, suffixes: Vec, - previous: ByteArray, + previous: Vec, _phantom: PhantomData, } @@ -688,7 +688,7 @@ impl DeltaByteArrayEncoder { suffix_writer: DeltaLengthByteArrayEncoder::new(), prefix_lengths: Vec::new(), suffixes: Vec::new(), - previous: Bytes::new().into(), + previous: Vec::new(), _phantom: PhantomData, } } @@ -707,10 +707,10 @@ impl Encoder for DeltaByteArrayEncoder { ), }); - let mut previous_array = &self.previous; + let mut previous = self.previous.as_slice(); + let mut previous_array = &ByteArray::from(Bytes::new()); for current_array in values { let current = current_array.data(); - let previous = previous_array.data(); // Maximum prefix length that is shared between previous value and current value let prefix_len = cmp::min(previous.len(), current.len()); let mut match_len = 0; @@ -721,9 +721,11 @@ impl Encoder for DeltaByteArrayEncoder { self.suffixes .push(current_array.slice(match_len, current.len() - match_len)); // Update previous for the next prefix + previous = current; previous_array = current_array; } - self.previous = previous_array.clone(); + self.previous.clear(); + self.previous.extend_from_slice(previous_array.data()); self.prefix_len_encoder.put(&self.prefix_lengths)?; self.suffix_writer.put(&self.suffixes)?; @@ -751,7 +753,7 @@ impl Encoder for DeltaByteArrayEncoder { // ... followed by suffixes self.suffix_writer.flush_to(out)?; - self.previous = Bytes::new().into(); + self.previous.clear(); Ok(()) } _ => panic!( From 24c480e293fb95c81340ff2d2f2bee6dbdf596b7 Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Thu, 23 Jul 2026 14:17:27 +0200 Subject: [PATCH 10/11] use capacity over len again --- parquet/src/encodings/encoding/mod.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/parquet/src/encodings/encoding/mod.rs b/parquet/src/encodings/encoding/mod.rs index 00939349a95f..867a348f6be3 100644 --- a/parquet/src/encodings/encoding/mod.rs +++ b/parquet/src/encodings/encoding/mod.rs @@ -766,7 +766,7 @@ impl Encoder for DeltaByteArrayEncoder { fn estimated_memory_size(&self) -> usize { self.prefix_len_encoder.estimated_memory_size() + self.suffix_writer.estimated_memory_size() - + (self.previous.len() * std::mem::size_of::()) + + (self.previous.capacity() * std::mem::size_of::()) } } From dda796d6e883624d90bf5683515f30070a8f9c2e Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Thu, 30 Jul 2026 13:52:33 +0200 Subject: [PATCH 11/11] use correct output in VariableWidthByteStreamSplitEncoder --- parquet/src/encodings/encoding/byte_stream_split_encoder.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/parquet/src/encodings/encoding/byte_stream_split_encoder.rs b/parquet/src/encodings/encoding/byte_stream_split_encoder.rs index 3fa0ddbb9a9c..ca5ba8b78fce 100644 --- a/parquet/src/encodings/encoding/byte_stream_split_encoder.rs +++ b/parquet/src/encodings/encoding/byte_stream_split_encoder.rs @@ -223,7 +223,7 @@ impl Encoder for VariableWidthByteStreamSplitEncoder { 6 => split_streams_const::<6>(&self.buffer, encoded), 7 => split_streams_const::<7>(&self.buffer, encoded), 8 => split_streams_const::<8>(&self.buffer, encoded), - _ => split_streams_variable(&self.buffer, out, type_size), + _ => split_streams_variable(&self.buffer, encoded, type_size), } self.buffer.clear();