diff --git a/parquet/src/encodings/encoding/byte_stream_split_encoder.rs b/parquet/src/encodings/encoding/byte_stream_split_encoder.rs index 0726c6b1c919..ca5ba8b78fce 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, encoded, 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..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()); @@ -164,9 +166,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 } @@ -184,6 +183,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 eeabcf4ba5ce..867a348f6be3 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 @@ -144,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 } @@ -164,6 +168,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)?; @@ -225,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 } @@ -260,6 +269,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 @@ -468,27 +483,22 @@ 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 } 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 +508,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. @@ -566,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 @@ -584,6 +596,7 @@ impl DeltaLengthByteArrayEncoder { pub fn new() -> Self { Self { len_encoder: DeltaBitPackEncoder::new(), + lengths: vec![], data: vec![], encoded_size: 0, _phantom: PhantomData, @@ -594,30 +607,27 @@ 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" ); - 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(); + self.lengths.reserve(values.len()); + for byte_array in values { + let len = byte_array.len(); + self.lengths.push(len as i32); + self.encoded_size += len; self.data.push(byte_array.clone()); } + self.len_encoder.put(&self.lengths)?; + self.lengths.clear(); 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 } @@ -626,22 +636,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 | Type::FIXED_LEN_BYTE_ARRAY, + 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. @@ -658,6 +668,8 @@ impl Encoder for DeltaLengthByteArrayEncoder { pub struct DeltaByteArrayEncoder { prefix_len_encoder: DeltaBitPackEncoder, suffix_writer: DeltaLengthByteArrayEncoder, + prefix_lengths: Vec, + suffixes: Vec, previous: Vec, _phantom: PhantomData, } @@ -674,7 +686,9 @@ impl DeltaByteArrayEncoder { Self { prefix_len_encoder: DeltaBitPackEncoder::new(), suffix_writer: DeltaLengthByteArrayEncoder::new(), - previous: vec![], + prefix_lengths: Vec::new(), + suffixes: Vec::new(), + previous: Vec::new(), _phantom: PhantomData, } } @@ -682,45 +696,46 @@ 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![]; + self.prefix_lengths.reserve(values.len()); + self.suffixes.reserve(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(); - // 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 = self.previous.as_slice(); + let mut previous_array = &ByteArray::from(Bytes::new()); + for current_array in values { + let current = current_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)); + 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 - self.previous.clear(); - self.previous.extend_from_slice(current); + previous = current; + previous_array = current_array; } - self.prefix_len_encoder.put(&prefix_lengths)?; - self.suffix_writer.put(&suffixes)?; + 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)?; + + self.prefix_lengths.clear(); + self.suffixes.clear(); 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 } @@ -730,21 +745,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"