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
26 changes: 25 additions & 1 deletion crates/iceberg/public-api.txt
Original file line number Diff line number Diff line change
Expand Up @@ -142,12 +142,16 @@ pub fn iceberg::cache::ObjectCacheProvide::manifest_list_cache(&self) -> &dyn ic
pub type iceberg::cache::ObjectCacheProvider = alloc::sync::Arc<dyn iceberg::cache::ObjectCacheProvide>
pub mod iceberg::compression
pub enum iceberg::compression::CompressionCodec
pub iceberg::compression::CompressionCodec::Brotli(u8)
pub iceberg::compression::CompressionCodec::Gzip(u8)
pub iceberg::compression::CompressionCodec::Lz4
pub iceberg::compression::CompressionCodec::Lz4Raw
pub iceberg::compression::CompressionCodec::Lzo
pub iceberg::compression::CompressionCodec::None
pub iceberg::compression::CompressionCodec::Snappy
pub iceberg::compression::CompressionCodec::Zstd(u8)
impl iceberg::compression::CompressionCodec
pub const fn iceberg::compression::CompressionCodec::brotli_default() -> Self
pub const fn iceberg::compression::CompressionCodec::gzip_default() -> Self
pub fn iceberg::compression::CompressionCodec::name(&self) -> &'static str
pub const fn iceberg::compression::CompressionCodec::zstd_default() -> Self
Expand Down Expand Up @@ -1162,12 +1166,16 @@ pub mod iceberg::partitioning
pub fn iceberg::partitioning::compute_unified_partition_type<'a>(partition_specs: impl core::iter::traits::iterator::Iterator<Item = &'a iceberg::spec::PartitionSpec>, schema: &iceberg::spec::Schema) -> iceberg::Result<iceberg::spec::StructType>
pub mod iceberg::puffin
pub enum iceberg::puffin::CompressionCodec
pub iceberg::puffin::CompressionCodec::Brotli(u8)
pub iceberg::puffin::CompressionCodec::Gzip(u8)
pub iceberg::puffin::CompressionCodec::Lz4
pub iceberg::puffin::CompressionCodec::Lz4Raw
pub iceberg::puffin::CompressionCodec::Lzo
pub iceberg::puffin::CompressionCodec::None
pub iceberg::puffin::CompressionCodec::Snappy
pub iceberg::puffin::CompressionCodec::Zstd(u8)
impl iceberg::compression::CompressionCodec
pub const fn iceberg::compression::CompressionCodec::brotli_default() -> Self
pub const fn iceberg::compression::CompressionCodec::gzip_default() -> Self
pub fn iceberg::compression::CompressionCodec::name(&self) -> &'static str
pub const fn iceberg::compression::CompressionCodec::zstd_default() -> Self
Expand Down Expand Up @@ -2764,6 +2772,11 @@ pub iceberg::spec::TableProperties::max_ref_age_ms: i64
pub iceberg::spec::TableProperties::max_snapshot_age_ms: i64
pub iceberg::spec::TableProperties::metadata_compression_codec: iceberg::compression::CompressionCodec
pub iceberg::spec::TableProperties::min_snapshots_to_keep: usize
pub iceberg::spec::TableProperties::parquet_compression_codec: iceberg::compression::CompressionCodec
pub iceberg::spec::TableProperties::parquet_dict_size_bytes: usize
pub iceberg::spec::TableProperties::parquet_page_row_limit: usize
pub iceberg::spec::TableProperties::parquet_page_size_bytes: usize
pub iceberg::spec::TableProperties::parquet_row_group_size_bytes: usize
pub iceberg::spec::TableProperties::write_data_location: core::option::Option<alloc::string::String>
pub iceberg::spec::TableProperties::write_datafusion_fanout_enabled: bool
pub iceberg::spec::TableProperties::write_folder_storage_location: core::option::Option<alloc::string::String>
Expand Down Expand Up @@ -2816,6 +2829,17 @@ pub const iceberg::spec::TableProperties::PROPERTY_PARQUET_CDC_MIN_CHUNK_SIZE: &
pub const iceberg::spec::TableProperties::PROPERTY_PARQUET_CDC_MIN_CHUNK_SIZE_DEFAULT: usize
pub const iceberg::spec::TableProperties::PROPERTY_PARQUET_CDC_NORM_LEVEL: &str
pub const iceberg::spec::TableProperties::PROPERTY_PARQUET_CDC_NORM_LEVEL_DEFAULT: i32
pub const iceberg::spec::TableProperties::PROPERTY_PARQUET_COMPRESSION_CODEC: &str
pub const iceberg::spec::TableProperties::PROPERTY_PARQUET_COMPRESSION_CODEC_DEFAULT: &str
pub const iceberg::spec::TableProperties::PROPERTY_PARQUET_COMPRESSION_LEVEL: &str
pub const iceberg::spec::TableProperties::PROPERTY_PARQUET_DICT_SIZE_BYTES: &str
pub const iceberg::spec::TableProperties::PROPERTY_PARQUET_DICT_SIZE_BYTES_DEFAULT: usize
pub const iceberg::spec::TableProperties::PROPERTY_PARQUET_PAGE_ROW_LIMIT: &str
pub const iceberg::spec::TableProperties::PROPERTY_PARQUET_PAGE_ROW_LIMIT_DEFAULT: usize
pub const iceberg::spec::TableProperties::PROPERTY_PARQUET_PAGE_SIZE_BYTES: &str
pub const iceberg::spec::TableProperties::PROPERTY_PARQUET_PAGE_SIZE_BYTES_DEFAULT: usize
pub const iceberg::spec::TableProperties::PROPERTY_PARQUET_ROW_GROUP_SIZE_BYTES: &str
pub const iceberg::spec::TableProperties::PROPERTY_PARQUET_ROW_GROUP_SIZE_BYTES_DEFAULT: usize
pub const iceberg::spec::TableProperties::PROPERTY_SNAPSHOT_COUNT: &str
pub const iceberg::spec::TableProperties::PROPERTY_UUID: &str
pub const iceberg::spec::TableProperties::PROPERTY_WRITE_DATA_LOCATION: &str
Expand Down Expand Up @@ -3306,7 +3330,7 @@ pub async fn iceberg::writer::file_writer::ParquetWriter::close(self) -> iceberg
pub async fn iceberg::writer::file_writer::ParquetWriter::write(&mut self, batch: &arrow_array::record_batch::RecordBatch) -> iceberg::Result<()>
pub struct iceberg::writer::file_writer::ParquetWriterBuilder
impl iceberg::writer::file_writer::ParquetWriterBuilder
pub fn iceberg::writer::file_writer::ParquetWriterBuilder::from_table_properties(table_props: &iceberg::spec::TableProperties, schema: iceberg::spec::SchemaRef) -> Self
pub fn iceberg::writer::file_writer::ParquetWriterBuilder::from_table_properties(table_props: &iceberg::spec::TableProperties, schema: iceberg::spec::SchemaRef) -> iceberg::Result<Self>
pub fn iceberg::writer::file_writer::ParquetWriterBuilder::new(props: parquet::file::properties::WriterProperties, schema: iceberg::spec::SchemaRef) -> Self
pub fn iceberg::writer::file_writer::ParquetWriterBuilder::new_with_match_mode(props: parquet::file::properties::WriterProperties, schema: iceberg::spec::SchemaRef, match_mode: iceberg::arrow::FieldMatchMode) -> Self
pub fn iceberg::writer::file_writer::ParquetWriterBuilder::with_match_mode(self, match_mode: iceberg::arrow::FieldMatchMode) -> Self
Expand Down
62 changes: 60 additions & 2 deletions crates/iceberg/src/compression.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,8 @@ const ZSTD_DEFAULT_LEVEL: u8 = 3;
const GZIP_DEFAULT_LEVEL: u8 = 6;
/// Maximum compression level for Gzip.
const GZIP_MAX_LEVEL: u8 = 9;
/// Default compression level for Brotli.
const BROTLI_DEFAULT_LEVEL: u8 = 1;

/// Data compression formats
#[derive(Debug, PartialEq, Eq, Clone, Copy, Default)]
Expand All @@ -42,13 +44,20 @@ pub enum CompressionCodec {
None,
/// LZ4 single compression frame with content size present
Lz4,
/// LZ4 raw block compression (no frame). Used by Parquet as `lz4_raw`.
Lz4Raw,
/// Zstandard single compression frame with content size present.
/// Level range is 0–22, where 0 means default compression level (not no compression).
/// Use [`CompressionCodec::zstd_default`] to construct with the default level.
Zstd(u8),
/// Gzip compression. Level range is 0–9, where 0 means no compression.
/// Use [`CompressionCodec::gzip_default`] to construct with the default level.
Gzip(u8),
/// Brotli compression. Level range is 0–11.
/// Use [`CompressionCodec::brotli_default`] to construct with the default level.
Brotli(u8),
/// LZO compression
Lzo,
/// Snappy compression
Snappy,
}
Expand All @@ -64,13 +73,21 @@ impl CompressionCodec {
CompressionCodec::Gzip(GZIP_DEFAULT_LEVEL)
}

/// Returns a Brotli codec with the default compression level.
pub const fn brotli_default() -> Self {
CompressionCodec::Brotli(BROTLI_DEFAULT_LEVEL)
}

/// Returns the codec name as used in serialization and error messages.
pub fn name(&self) -> &'static str {
match self {
CompressionCodec::None => "none",
CompressionCodec::Lz4 => "lz4",
CompressionCodec::Lz4Raw => "lz4_raw",
CompressionCodec::Zstd(_) => "zstd",
CompressionCodec::Gzip(_) => "gzip",
CompressionCodec::Brotli(_) => "brotli",
CompressionCodec::Lzo => "lzo",
CompressionCodec::Snappy => "snappy",
}
}
Expand All @@ -90,13 +107,24 @@ impl<'de> Deserialize<'de> for CompressionCodec {
fn deserialize<D: Deserializer<'de>>(deserializer: D) -> std::result::Result<Self, D::Error> {
let s = String::deserialize(deserializer)?;
match s.to_lowercase().as_str() {
"none" => Ok(CompressionCodec::None),
"none" | "uncompressed" => Ok(CompressionCodec::None),
"lz4" => Ok(CompressionCodec::Lz4),
"lz4_raw" => Ok(CompressionCodec::Lz4Raw),
"zstd" => Ok(CompressionCodec::zstd_default()),
"gzip" => Ok(CompressionCodec::gzip_default()),
"brotli" => Ok(CompressionCodec::brotli_default()),
"lzo" => Ok(CompressionCodec::Lzo),
"snappy" => Ok(CompressionCodec::Snappy),
other => Err(serde::de::Error::unknown_variant(other, &[
"none", "lz4", "zstd", "gzip", "snappy",
"none",
"uncompressed",
"lz4",
"lz4_raw",
"zstd",
"gzip",
"brotli",
"lzo",
"snappy",
])),
}
}
Expand All @@ -107,8 +135,11 @@ impl fmt::Display for CompressionCodec {
match self {
CompressionCodec::None => write!(f, "None"),
CompressionCodec::Lz4 => write!(f, "Lz4"),
CompressionCodec::Lz4Raw => write!(f, "Lz4Raw"),
CompressionCodec::Zstd(level) => write!(f, "Zstd(level={level})"),
CompressionCodec::Gzip(level) => write!(f, "Gzip(level={level})"),
CompressionCodec::Brotli(level) => write!(f, "Brotli(level={level})"),
CompressionCodec::Lzo => write!(f, "Lzo"),
CompressionCodec::Snappy => write!(f, "Snappy"),
}
}
Expand All @@ -129,6 +160,18 @@ impl CompressionCodec {
decoder.read_to_end(&mut decompressed)?;
Ok(decompressed)
}
CompressionCodec::Lz4Raw => Err(Error::new(
ErrorKind::FeatureUnsupported,
"LZ4_RAW decompression is not supported currently",
)),
CompressionCodec::Brotli(_) => Err(Error::new(
ErrorKind::FeatureUnsupported,
"Brotli decompression is not supported currently",
)),
CompressionCodec::Lzo => Err(Error::new(
ErrorKind::FeatureUnsupported,
"LZO decompression is not supported currently",
)),
CompressionCodec::Snappy => Err(Error::new(
ErrorKind::FeatureUnsupported,
"Snappy decompression is not supported currently",
Expand Down Expand Up @@ -157,6 +200,18 @@ impl CompressionCodec {
encoder.write_all(&bytes)?;
Ok(encoder.finish()?)
}
CompressionCodec::Lz4Raw => Err(Error::new(
ErrorKind::FeatureUnsupported,
"LZ4_RAW compression is not supported currently",
)),
CompressionCodec::Brotli(_) => Err(Error::new(
ErrorKind::FeatureUnsupported,
"Brotli compression is not supported currently",
)),
CompressionCodec::Lzo => Err(Error::new(
ErrorKind::FeatureUnsupported,
"LZO compression is not supported currently",
)),
CompressionCodec::Snappy => Err(Error::new(
ErrorKind::FeatureUnsupported,
"Snappy compression is not supported currently",
Expand All @@ -179,7 +234,10 @@ impl CompressionCodec {
CompressionCodec::None => Ok(""),
CompressionCodec::Gzip(_) => Ok(".gz"),
codec @ (CompressionCodec::Lz4
| CompressionCodec::Lz4Raw
| CompressionCodec::Zstd(_)
| CompressionCodec::Brotli(_)
| CompressionCodec::Lzo
| CompressionCodec::Snappy) => Err(Error::new(
ErrorKind::FeatureUnsupported,
format!("suffix not defined for {codec:?}"),
Expand Down
Loading
Loading