Skip to content
Merged
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
2 changes: 1 addition & 1 deletion parquet/benches/metadata.rs
Original file line number Diff line number Diff line change
Expand Up @@ -90,7 +90,7 @@ fn encoded_meta(is_nullable: bool, has_lists: bool) -> Vec<u8> {
.map(|j| {
ColumnChunkMetaData::builder(column_desc_ptrs[j].clone())
.set_encodings(vec![Encoding::PLAIN, Encoding::RLE_DICTIONARY])
.set_compression(parquet::basic::Compression::UNCOMPRESSED)
.set_compression_codec(parquet::basic::CompressionCodec::UNCOMPRESSED)
.set_num_values(rng.random_range(1..1000000))
.set_total_compressed_size(rng.random_range(50000..5000000))
.set_data_page_offset(rng.random_range(4..2000000000))
Expand Down
155 changes: 109 additions & 46 deletions parquet/src/basic.rs
Original file line number Diff line number Diff line change
Expand Up @@ -798,6 +798,33 @@ fn i32_to_encoding(val: i32) -> Encoding {
// ----------------------------------------------------------------------
// Mirrors thrift enum `CompressionCodec`

thrift_enum!(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I like that this matches the name CompressionCodec used in the thrift definition: https://github.com/apache/parquet-format/blob/96edf77704b60b6f3ca2232c218c64eff6c874d3/src/main/thrift/parquet.thrift#L642

I think it makes the code easier to understand

/// Supported compression algorithms.
///
/// Codecs added in format version X.Y can be read by readers based on X.Y and later.
/// Codec support may vary between readers based on the format version and
/// libraries available at runtime.
///
/// See [Compression.md] for a detailed specification of these algorithms.
///
/// [Compression.md]: https://github.com/apache/parquet-format/blob/master/Compression.md
enum CompressionCodec {
UNCOMPRESSED = 0;
SNAPPY = 1;
GZIP = 2;
LZO = 3;
BROTLI = 4; // Added in 2.4
LZ4 = 5; // DEPRECATED (Added in 2.4)
ZSTD = 6; // Added in 2.4
LZ4_RAW = 7; // Added in 2.9
}
);

// NOTE: This enum likely belongs in file::properties now, but moving it there would be a
// breaking API change, that's probably not worth the pain. If a new codec is added to the
// Parquet specification, or any other breaking changes are made to this enum, this can be
// revisited.

/// Supported block compression algorithms.
///
/// Block compression can yield non-trivial improvements to storage efficiency at the expense
Expand Down Expand Up @@ -834,51 +861,33 @@ pub enum Compression {
LZ4_RAW,
}

impl<'a, R: ThriftCompactInputProtocol<'a>> ReadThrift<'a, R> for Compression {
fn read_thrift(prot: &mut R) -> Result<Self> {
let val = prot.read_i32()?;
Ok(match val {
0 => Self::UNCOMPRESSED,
1 => Self::SNAPPY,
2 => Self::GZIP(Default::default()),
3 => Self::LZO,
4 => Self::BROTLI(Default::default()),
5 => Self::LZ4,
6 => Self::ZSTD(Default::default()),
7 => Self::LZ4_RAW,
_ => return Err(general_err!("Unexpected CompressionCodec {}", val)),
})
}
}

// TODO(ets): explore replacing this with a thrift_enum!(ThriftCompression) for the serialization
// and then provide `From` impls to convert back and forth. This is necessary due to the addition
// of compression level to some variants.
impl WriteThrift for Compression {
const ELEMENT_TYPE: ElementType = ElementType::I32;

fn write_thrift<W: Write>(&self, writer: &mut ThriftCompactOutputProtocol<W>) -> Result<()> {
let id: i32 = match *self {
Self::UNCOMPRESSED => 0,
Self::SNAPPY => 1,
Self::GZIP(_) => 2,
Self::LZO => 3,
Self::BROTLI(_) => 4,
Self::LZ4 => 5,
Self::ZSTD(_) => 6,
Self::LZ4_RAW => 7,
};
writer.write_i32(id)
impl From<CompressionCodec> for Compression {
fn from(value: CompressionCodec) -> Self {
match value {
CompressionCodec::UNCOMPRESSED => Compression::UNCOMPRESSED,
CompressionCodec::SNAPPY => Compression::SNAPPY,
CompressionCodec::GZIP => Compression::GZIP(Default::default()),
CompressionCodec::LZO => Compression::LZO,
CompressionCodec::BROTLI => Compression::BROTLI(Default::default()),
CompressionCodec::LZ4 => Compression::LZ4,
CompressionCodec::ZSTD => Compression::ZSTD(Default::default()),
CompressionCodec::LZ4_RAW => Compression::LZ4_RAW,
}
}
}

write_thrift_field!(Compression, FieldType::I32);

impl Compression {
/// Returns the codec type of this compression setting as a string, without the compression
/// level.
pub(crate) fn codec_to_string(self) -> String {
format!("{self:?}").split('(').next().unwrap().to_owned()
impl From<Compression> for CompressionCodec {
fn from(value: Compression) -> Self {
match value {
Compression::UNCOMPRESSED => CompressionCodec::UNCOMPRESSED,
Compression::SNAPPY => CompressionCodec::SNAPPY,
Compression::GZIP(_) => CompressionCodec::GZIP,
Compression::LZO => CompressionCodec::LZO,
Compression::BROTLI(_) => CompressionCodec::BROTLI,
Compression::LZ4 => CompressionCodec::LZ4,
Compression::ZSTD(_) => CompressionCodec::ZSTD,
Compression::LZ4_RAW => CompressionCodec::LZ4_RAW,
}
}
}

Expand Down Expand Up @@ -2140,11 +2149,65 @@ mod tests {
}

#[test]
fn test_compression_codec_to_string() {
assert_eq!(Compression::UNCOMPRESSED.codec_to_string(), "UNCOMPRESSED");
fn test_compression_conversion() {
assert_eq!(
CompressionCodec::from(Compression::UNCOMPRESSED),
CompressionCodec::UNCOMPRESSED
);
assert_eq!(
CompressionCodec::from(Compression::SNAPPY),
CompressionCodec::SNAPPY
);
assert_eq!(
CompressionCodec::from(Compression::GZIP(Default::default())),
CompressionCodec::GZIP
);
assert_eq!(
CompressionCodec::from(Compression::LZO),
CompressionCodec::LZO
);
assert_eq!(
CompressionCodec::from(Compression::BROTLI(Default::default())),
CompressionCodec::BROTLI
);
assert_eq!(
CompressionCodec::from(Compression::LZ4),
CompressionCodec::LZ4
);
assert_eq!(
CompressionCodec::from(Compression::ZSTD(Default::default())),
CompressionCodec::ZSTD
);
assert_eq!(
CompressionCodec::from(Compression::LZ4_RAW),
CompressionCodec::LZ4_RAW
);

assert_eq!(
Compression::from(CompressionCodec::UNCOMPRESSED),
Compression::UNCOMPRESSED
);
assert_eq!(
Compression::from(CompressionCodec::SNAPPY),
Compression::SNAPPY
);
assert_eq!(
Compression::from(CompressionCodec::GZIP),
Compression::GZIP(Default::default())
);
assert_eq!(Compression::from(CompressionCodec::LZO), Compression::LZO);
assert_eq!(
Compression::from(CompressionCodec::BROTLI),
Compression::BROTLI(Default::default())
);
assert_eq!(Compression::from(CompressionCodec::LZ4), Compression::LZ4);
assert_eq!(
Compression::from(CompressionCodec::ZSTD),
Compression::ZSTD(Default::default())
);
assert_eq!(
Compression::ZSTD(ZstdLevel::default()).codec_to_string(),
"ZSTD"
Compression::from(CompressionCodec::LZ4_RAW),
Compression::LZ4_RAW
);
}

Expand Down
22 changes: 11 additions & 11 deletions parquet/src/bin/parquet-layout.rs
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,7 @@ use parquet::file::metadata::ParquetMetaDataReader;
use serde::Serialize;
use thrift::protocol::TCompactInputProtocol;

use parquet::basic::Compression;
use parquet::basic::CompressionCodec;
use parquet::errors::Result;
use parquet::file::reader::ChunkReader;
#[allow(deprecated)]
Expand Down Expand Up @@ -118,7 +118,7 @@ fn do_layout<C: ChunkReader>(reader: &C) -> Result<ParquetFile> {
.iter()
.zip(schema.columns())
.map(|(column, column_schema)| {
let compression = compression(column.compression());
let compression = compression(column.compression_codec());
let mut pages = vec![];

let mut start = column
Expand Down Expand Up @@ -225,16 +225,16 @@ fn read_page_header<C: ChunkReader>(reader: &C, offset: u64) -> Result<(usize, P
}

/// Returns a string representation for a given compression
fn compression(compression: Compression) -> Option<&'static str> {
fn compression(compression: CompressionCodec) -> Option<&'static str> {
match compression {
Compression::UNCOMPRESSED => None,
Compression::SNAPPY => Some("snappy"),
Compression::GZIP(_) => Some("gzip"),
Compression::LZO => Some("lzo"),
Compression::BROTLI(_) => Some("brotli"),
Compression::LZ4 => Some("lz4"),
Compression::ZSTD(_) => Some("zstd"),
Compression::LZ4_RAW => Some("lz4_raw"),
CompressionCodec::UNCOMPRESSED => None,
CompressionCodec::SNAPPY => Some("snappy"),
CompressionCodec::GZIP => Some("gzip"),
CompressionCodec::LZO => Some("lzo"),
CompressionCodec::BROTLI => Some("brotli"),
CompressionCodec::LZ4 => Some("lz4"),
CompressionCodec::ZSTD => Some("zstd"),
CompressionCodec::LZ4_RAW => Some("lz4_raw"),
}
}

Expand Down
4 changes: 2 additions & 2 deletions parquet/src/file/metadata/memory.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@
//! Memory calculations for [`ParquetMetadata::memory_size`]
//!
//! [`ParquetMetadata::memory_size`]: crate::file::metadata::ParquetMetaData::memory_size
use crate::basic::{BoundaryOrder, ColumnOrder, Compression, Encoding, PageType};
use crate::basic::{BoundaryOrder, ColumnOrder, CompressionCodec, Encoding, PageType};
use crate::data_type::private::ParquetValueType;
use crate::file::metadata::{
ColumnChunkMetaData, FileMetaData, KeyValue, PageEncodingStats, ParquetPageEncodingStats,
Expand Down Expand Up @@ -206,7 +206,7 @@ impl HeapSize for SortingColumn {
0 // no heap allocations
}
}
impl HeapSize for Compression {
impl HeapSize for CompressionCodec {
fn heap_size(&self) -> usize {
0 // no heap allocations
}
Expand Down
Loading
Loading