Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 23 additions & 17 deletions parquet/src/encodings/encoding/byte_stream_split_encoder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -88,12 +88,15 @@ impl<T: DataType> Encoder<T> for ByteStreamSplitEncoder<T> {
self.buffer.len()
}

fn flush_buffer(&mut self) -> Result<Bytes> {
let mut encoded = vec![0; self.buffer.len()];
fn flush_to(&mut self, out: &mut Vec<u8>) -> 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",
Expand All @@ -103,7 +106,7 @@ impl<T: DataType> Encoder<T> for ByteStreamSplitEncoder<T> {
}

self.buffer.clear();
Ok(encoded.into())
Ok(())
}

/// return the estimated memory size of this encoder.
Expand Down Expand Up @@ -202,26 +205,29 @@ impl<T: DataType> Encoder<T> for VariableWidthByteStreamSplitEncoder<T> {
self.buffer.len()
}

fn flush_buffer(&mut self) -> Result<Bytes> {
let mut encoded = vec![0; self.buffer.len()];
fn flush_to(&mut self, out: &mut Vec<u8>) -> 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.
Expand Down
11 changes: 8 additions & 3 deletions parquet/src/encodings/encoding/dict_encoder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -132,6 +132,8 @@ impl<T: DataType> DictEncoder<T> {
/// Writes out the dictionary values with RLE encoding in a byte buffer, and return
/// the result.
pub fn write_indices(&mut self) -> Result<Bytes> {
// 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());
Expand Down Expand Up @@ -164,9 +166,6 @@ impl<T: DataType> Encoder<T> for DictEncoder<T> {
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
}
Expand All @@ -184,6 +183,12 @@ impl<T: DataType> Encoder<T> for DictEncoder<T> {
self.write_indices()
}

fn flush_to(&mut self, out: &mut Vec<u8>) -> 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`]).
Expand Down
Loading
Loading