Skip to content
Draft
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
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,7 @@ itertools = "0.13"
linkedbytes = "0.1.8"
metainfo = "0.7.14"
minijinja = "2.12.0"
miniz_oxide = "0.8"
mockall = "0.13.1"
mockito = "1"
motore-macros = "0.4.3"
Expand Down
1 change: 1 addition & 0 deletions crates/iceberg/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,7 @@ form_urlencoded = { workspace = true }
futures = { workspace = true }
iceberg-property-macro = { workspace = true }
itertools = { workspace = true }
miniz_oxide = { workspace = true }
moka = { version = "0.12.10", features = ["future"] }
murmur3 = { workspace = true }
once_cell = { workspace = true }
Expand Down
14 changes: 9 additions & 5 deletions crates/iceberg/public-api.txt
Original file line number Diff line number Diff line change
Expand Up @@ -2194,9 +2194,9 @@ impl iceberg::spec::ManifestListWriter
pub fn iceberg::spec::ManifestListWriter::add_manifests(&mut self, manifests: impl core::iter::traits::iterator::Iterator<Item = iceberg::spec::ManifestFile>) -> iceberg::Result<()>
pub async fn iceberg::spec::ManifestListWriter::close(self) -> iceberg::Result<iceberg::io::FileMetadata>
pub fn iceberg::spec::ManifestListWriter::next_row_id(&self) -> core::option::Option<u64>
pub fn iceberg::spec::ManifestListWriter::v1(writer: alloc::boxed::Box<dyn iceberg::io::FileWrite>, snapshot_id: i64, parent_snapshot_id: core::option::Option<i64>) -> Self
pub fn iceberg::spec::ManifestListWriter::v2(writer: alloc::boxed::Box<dyn iceberg::io::FileWrite>, snapshot_id: i64, parent_snapshot_id: core::option::Option<i64>, sequence_number: i64) -> Self
pub fn iceberg::spec::ManifestListWriter::v3(writer: alloc::boxed::Box<dyn iceberg::io::FileWrite>, snapshot_id: i64, parent_snapshot_id: core::option::Option<i64>, sequence_number: i64, first_row_id: core::option::Option<u64>) -> Self
pub fn iceberg::spec::ManifestListWriter::v1(writer: alloc::boxed::Box<dyn iceberg::io::FileWrite>, snapshot_id: i64, parent_snapshot_id: core::option::Option<i64>, compression: iceberg::compression::CompressionCodec) -> iceberg::Result<Self>
pub fn iceberg::spec::ManifestListWriter::v2(writer: alloc::boxed::Box<dyn iceberg::io::FileWrite>, snapshot_id: i64, parent_snapshot_id: core::option::Option<i64>, sequence_number: i64, compression: iceberg::compression::CompressionCodec) -> iceberg::Result<Self>
pub fn iceberg::spec::ManifestListWriter::v3(writer: alloc::boxed::Box<dyn iceberg::io::FileWrite>, snapshot_id: i64, parent_snapshot_id: core::option::Option<i64>, sequence_number: i64, first_row_id: core::option::Option<u64>, compression: iceberg::compression::CompressionCodec) -> iceberg::Result<Self>
impl core::fmt::Debug for iceberg::spec::ManifestListWriter
pub fn iceberg::spec::ManifestListWriter::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result
pub struct iceberg::spec::ManifestMetadata
Expand Down Expand Up @@ -2238,8 +2238,8 @@ pub fn iceberg::spec::ManifestWriterBuilder::build_v2_data(self) -> iceberg::spe
pub fn iceberg::spec::ManifestWriterBuilder::build_v2_deletes(self) -> iceberg::spec::ManifestWriter
pub fn iceberg::spec::ManifestWriterBuilder::build_v3_data(self) -> iceberg::spec::ManifestWriter
pub fn iceberg::spec::ManifestWriterBuilder::build_v3_deletes(self) -> iceberg::spec::ManifestWriter
pub fn iceberg::spec::ManifestWriterBuilder::new(output: iceberg::io::OutputFile, snapshot_id: core::option::Option<i64>, schema: iceberg::spec::SchemaRef, partition_spec: iceberg::spec::PartitionSpec) -> Self
pub fn iceberg::spec::ManifestWriterBuilder::new_from_encrypted(encrypted_output: iceberg::encryption::EncryptedOutputFile, snapshot_id: core::option::Option<i64>, schema: iceberg::spec::SchemaRef, partition_spec: iceberg::spec::PartitionSpec) -> iceberg::Result<Self>
pub fn iceberg::spec::ManifestWriterBuilder::new(output: iceberg::io::OutputFile, snapshot_id: core::option::Option<i64>, schema: iceberg::spec::SchemaRef, partition_spec: iceberg::spec::PartitionSpec, compression: iceberg::compression::CompressionCodec) -> iceberg::Result<Self>
pub fn iceberg::spec::ManifestWriterBuilder::new_from_encrypted(encrypted_output: iceberg::encryption::EncryptedOutputFile, snapshot_id: core::option::Option<i64>, schema: iceberg::spec::SchemaRef, partition_spec: iceberg::spec::PartitionSpec, compression: iceberg::compression::CompressionCodec) -> iceberg::Result<Self>
pub struct iceberg::spec::Map
impl iceberg::spec::Map
pub fn iceberg::spec::Map::get(&self, key: &iceberg::spec::Literal) -> core::option::Option<&core::option::Option<iceberg::spec::Literal>>
Expand Down Expand Up @@ -2862,6 +2862,9 @@ impl core::fmt::Debug for iceberg::spec::TableMetadataBuilder
pub fn iceberg::spec::TableMetadataBuilder::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result
pub struct iceberg::spec::TableProperties<'properties>
impl iceberg::spec::TableProperties<'_>
pub const iceberg::spec::TableProperties<'_>::PROPERTY_AVRO_COMPRESSION_CODEC: &'static str
pub const iceberg::spec::TableProperties<'_>::PROPERTY_AVRO_COMPRESSION_CODEC_DEFAULT: &'static str
pub const iceberg::spec::TableProperties<'_>::PROPERTY_AVRO_COMPRESSION_LEVEL: &'static str
pub const iceberg::spec::TableProperties<'_>::PROPERTY_COMMIT_MAX_RETRY_WAIT_MS: &'static str
pub const iceberg::spec::TableProperties<'_>::PROPERTY_COMMIT_MAX_RETRY_WAIT_MS_DEFAULT: u64
pub const iceberg::spec::TableProperties<'_>::PROPERTY_COMMIT_MIN_RETRY_WAIT_MS: &'static str
Expand Down Expand Up @@ -2934,6 +2937,7 @@ pub const iceberg::spec::TableProperties<'_>::PROPERTY_WRITE_TARGET_FILE_SIZE_BY
pub const iceberg::spec::TableProperties<'_>::RESERVED_PROPERTIES: [&'static str; 9]
pub fn iceberg::spec::TableProperties<'_>::data_encryption_key_size(&self) -> iceberg::Result<iceberg::encryption::AesKeySize>
impl<'properties> iceberg::spec::TableProperties<'properties>
pub fn iceberg::spec::TableProperties<'properties>::avro_compression_codec(&self) -> iceberg::Result<iceberg::compression::CompressionCodec>
pub fn iceberg::spec::TableProperties<'properties>::cdc_enabled(&self) -> iceberg::Result<bool>
pub fn iceberg::spec::TableProperties<'properties>::cdc_max_chunk_size(&self) -> iceberg::Result<usize>
pub fn iceberg::spec::TableProperties<'properties>::cdc_min_chunk_size(&self) -> iceberg::Result<usize>
Expand Down
60 changes: 48 additions & 12 deletions crates/iceberg/src/compression.rs
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,22 @@ impl CompressionCodec {
CompressionCodec::Snappy => "snappy",
}
}

/// Parses a codec name (case-insensitive), returning `None` for unrecognized names.
/// Codecs that carry a level get their default level.
pub(crate) fn from_name(name: &str) -> Option<Self> {
Some(match name.to_lowercase().as_str() {
"none" | "uncompressed" => CompressionCodec::None,
"lz4" => CompressionCodec::Lz4,
"lz4_raw" => CompressionCodec::Lz4Raw,
"zstd" => CompressionCodec::zstd_default(),
"gzip" => CompressionCodec::gzip_default(),
"brotli" => CompressionCodec::brotli_default(),
"lzo" => CompressionCodec::Lzo,
"snappy" => CompressionCodec::Snappy,
_ => return None,
})
}
}

// Note: serialize/deserialize do not round-trip the compression level. Iceberg configuration
Expand All @@ -106,16 +122,8 @@ impl Serialize for CompressionCodec {
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" | "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, &[
Self::from_name(&s).ok_or_else(|| {
serde::de::Error::unknown_variant(&s, &[
"none",
"uncompressed",
"lz4",
Expand All @@ -125,8 +133,8 @@ impl<'de> Deserialize<'de> for CompressionCodec {
"brotli",
"lzo",
"snappy",
])),
}
])
})
}
}

Expand Down Expand Up @@ -219,6 +227,17 @@ impl CompressionCodec {
}
}

/// Replace the compression level, if this variant carries one.
/// Variants without a level are returned unchanged.
pub(crate) fn with_level(self, level: u8) -> Self {
match self {
CompressionCodec::Zstd(_) => CompressionCodec::Zstd(level),
CompressionCodec::Gzip(_) => CompressionCodec::Gzip(level),
CompressionCodec::Brotli(_) => CompressionCodec::Brotli(level),
other => other,
}
}

pub(crate) fn is_none(&self) -> bool {
matches!(self, CompressionCodec::None)
}
Expand Down Expand Up @@ -315,6 +334,23 @@ mod tests {
assert!(zstd_err.to_string().contains("suffix not defined for Zstd"));
}

#[test]
fn test_with_level() {
for (codec, expected) in [
(CompressionCodec::zstd_default(), CompressionCodec::Zstd(5)),
(CompressionCodec::gzip_default(), CompressionCodec::Gzip(5)),
(
CompressionCodec::brotli_default(),
CompressionCodec::Brotli(5),
),
(CompressionCodec::None, CompressionCodec::None),
(CompressionCodec::Lz4, CompressionCodec::Lz4),
(CompressionCodec::Snappy, CompressionCodec::Snappy),
] {
assert_eq!(codec.with_level(5), expected);
}
}

#[test]
fn test_display() {
assert_eq!(CompressionCodec::None.to_string(), "None");
Expand Down
15 changes: 13 additions & 2 deletions crates/iceberg/src/io/object_cache.rs
Original file line number Diff line number Diff line change
Expand Up @@ -225,6 +225,7 @@ mod tests {

use super::*;
use crate::TableIdent;
use crate::compression::CompressionCodec;
use crate::io::{FileIO, OutputFile};
use crate::spec::{
DataContentType, DataFileBuilder, DataFileFormat, FormatVersion, Literal,
Expand Down Expand Up @@ -307,7 +308,9 @@ mod tests {
Some(current_snapshot.snapshot_id()),
current_schema.clone(),
current_partition_spec.as_ref().clone(),
CompressionCodec::None,
)
.unwrap()
.build_v2_data();
writer
.add_entry(
Expand Down Expand Up @@ -344,7 +347,9 @@ mod tests {
current_snapshot.snapshot_id(),
current_snapshot.parent_snapshot_id(),
current_snapshot.sequence_number(),
);
CompressionCodec::None,
)
.unwrap();
manifest_list_write
.add_manifests(vec![data_file_manifest].into_iter())
.unwrap();
Expand All @@ -362,7 +367,9 @@ mod tests {
Some(current_snapshot.snapshot_id()),
current_schema.clone(),
current_partition_spec.as_ref().clone(),
CompressionCodec::None,
)
.unwrap()
.build_v1();
writer
.add_entry(
Expand Down Expand Up @@ -398,7 +405,9 @@ mod tests {
manifest_list_writer,
current_snapshot.snapshot_id(),
current_snapshot.parent_snapshot_id(),
);
CompressionCodec::None,
)
.unwrap();
manifest_list_write
.add_manifests(vec![data_file_manifest].into_iter())
.unwrap();
Expand Down Expand Up @@ -618,7 +627,9 @@ mod tests {
Some(1),
schema,
partition_spec,
CompressionCodec::None,
)
.unwrap()
.build_v3_data();
writer
.add_entry(
Expand Down
Loading
Loading