diff --git a/crates/iceberg/public-api.txt b/crates/iceberg/public-api.txt index 530aa4cbf3..acb588dd78 100644 --- a/crates/iceberg/public-api.txt +++ b/crates/iceberg/public-api.txt @@ -142,12 +142,16 @@ pub fn iceberg::cache::ObjectCacheProvide::manifest_list_cache(&self) -> &dyn ic pub type iceberg::cache::ObjectCacheProvider = alloc::sync::Arc 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 @@ -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, schema: &iceberg::spec::Schema) -> iceberg::Result 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 @@ -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 pub iceberg::spec::TableProperties::write_datafusion_fanout_enabled: bool pub iceberg::spec::TableProperties::write_folder_storage_location: core::option::Option @@ -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 @@ -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 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 diff --git a/crates/iceberg/src/compression.rs b/crates/iceberg/src/compression.rs index 929d9226e7..dbd2c97881 100644 --- a/crates/iceberg/src/compression.rs +++ b/crates/iceberg/src/compression.rs @@ -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)] @@ -42,6 +44,8 @@ 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. @@ -49,6 +53,11 @@ pub enum CompressionCodec { /// 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, } @@ -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", } } @@ -90,13 +107,24 @@ impl<'de> Deserialize<'de> for CompressionCodec { fn deserialize>(deserializer: D) -> std::result::Result { 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", ])), } } @@ -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"), } } @@ -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", @@ -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", @@ -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:?}"), diff --git a/crates/iceberg/src/spec/table_properties.rs b/crates/iceberg/src/spec/table_properties.rs index a35f1427a3..ae03319276 100644 --- a/crates/iceberg/src/spec/table_properties.rs +++ b/crates/iceberg/src/spec/table_properties.rs @@ -41,6 +41,28 @@ where }) } +/// Parse an optional property, returning `None` when the key is absent and an +/// error when the value is present but fails to parse. +fn parse_optional_property( + properties: &HashMap, + key: &str, +) -> Result> +where + ::Err: Display, +{ + properties + .get(key) + .map(|value| { + value.parse::().map_err(|e| { + Error::new( + ErrorKind::DataInvalid, + format!("Invalid value for {key}: {e}"), + ) + }) + }) + .transpose() +} + fn parse_location_property( properties: &HashMap, key: &str, @@ -117,6 +139,39 @@ pub(crate) fn parse_metadata_file_compression( } } +/// Parse the Parquet data-file compression codec (`write.parquet.compression-codec`) +/// and fold in the compression level (`write.parquet.compression-level`) for the +/// codecs that accept one (`zstd`, `gzip`, `brotli`). +fn parse_parquet_compression(properties: &HashMap) -> Result { + let value = properties + .get(TableProperties::PROPERTY_PARQUET_COMPRESSION_CODEC) + .map(|s| s.as_str()) + .unwrap_or(TableProperties::PROPERTY_PARQUET_COMPRESSION_CODEC_DEFAULT); + + let codec: CompressionCodec = + serde_json::from_value(serde_json::Value::String(value.to_lowercase())).map_err(|_| { + Error::new( + ErrorKind::DataInvalid, + format!( + "Invalid Parquet compression codec: {value}. Supported codecs: \ + uncompressed, snappy, gzip, lzo, brotli, lz4, lz4_raw, zstd" + ), + ) + })?; + + let level: Option = parse_optional_property( + properties, + TableProperties::PROPERTY_PARQUET_COMPRESSION_LEVEL, + )?; + + Ok(match (codec, level) { + (CompressionCodec::Zstd(_), Some(level)) => CompressionCodec::Zstd(level), + (CompressionCodec::Gzip(_), Some(level)) => CompressionCodec::Gzip(level), + (CompressionCodec::Brotli(_), Some(level)) => CompressionCodec::Brotli(level), + (codec, _) => codec, + }) +} + /// Parse boolean property case insensitively /// Rust standard library only accepts "true" and "false", see https://doc.rust-lang.org/std/primitive.bool.html#method.from_str /// Users might accidentally trigger fallback with valid configuration values such as "False" or "True" @@ -175,6 +230,18 @@ pub struct TableProperties { pub cdc_max_chunk_size: usize, /// Content-defined chunking normalization level (gearhash bit adjustment). pub cdc_norm_level: i32, + /// Parquet compression codec for data files, with the resolved compression + /// level folded in (from `write.parquet.compression-level`, or the codec's + /// default when unset). + pub parquet_compression_codec: CompressionCodec, + /// Approximate maximum Parquet row group size in bytes. + pub parquet_row_group_size_bytes: usize, + /// Approximate maximum Parquet data page size in bytes. + pub parquet_page_size_bytes: usize, + /// Maximum number of rows per Parquet data page. + pub parquet_page_row_limit: usize, + /// Approximate maximum Parquet dictionary page size in bytes. + pub parquet_dict_size_bytes: usize, /// The master key id used to encrypt this table's manifest list and data /// files. `None` if `encryption.key-id` is not set. pub encryption_key_id: Option, @@ -338,6 +405,37 @@ impl TableProperties { /// Default matches `parquet::file::properties::DEFAULT_CDC_NORM_LEVEL`. pub const PROPERTY_PARQUET_CDC_NORM_LEVEL_DEFAULT: i32 = 0; + /// Compression codec for Parquet data files (e.g. `zstd`, `gzip`, `snappy`, + /// `lz4`, `lz4_raw`, `brotli`, `lzo`, `uncompressed`). The codec name is + /// parsed into a [`CompressionCodec`] when properties are parsed; the level's + /// range is validated when the writer is built. + pub const PROPERTY_PARQUET_COMPRESSION_CODEC: &str = "write.parquet.compression-codec"; + /// Default Parquet compression codec. + pub const PROPERTY_PARQUET_COMPRESSION_CODEC_DEFAULT: &str = "zstd"; + /// Compression level for Parquet data files, for codecs that take one + /// (`gzip`, `zstd`, `brotli`). When unset, the codec's default level is used. + pub const PROPERTY_PARQUET_COMPRESSION_LEVEL: &str = "write.parquet.compression-level"; + + /// Approximate maximum size of a Parquet row group in bytes. + pub const PROPERTY_PARQUET_ROW_GROUP_SIZE_BYTES: &str = "write.parquet.row-group-size-bytes"; + /// Default Parquet row group size in bytes. + pub const PROPERTY_PARQUET_ROW_GROUP_SIZE_BYTES_DEFAULT: usize = 128 * 1024 * 1024; + + /// Approximate maximum size of a Parquet data page in bytes. + pub const PROPERTY_PARQUET_PAGE_SIZE_BYTES: &str = "write.parquet.page-size-bytes"; + /// Default Parquet page size in bytes. + pub const PROPERTY_PARQUET_PAGE_SIZE_BYTES_DEFAULT: usize = 1024 * 1024; + + /// Maximum number of rows per Parquet data page. + pub const PROPERTY_PARQUET_PAGE_ROW_LIMIT: &str = "write.parquet.page-row-limit"; + /// Default Parquet page row limit. + pub const PROPERTY_PARQUET_PAGE_ROW_LIMIT_DEFAULT: usize = 20000; + + /// Approximate maximum size of the Parquet dictionary page in bytes. + pub const PROPERTY_PARQUET_DICT_SIZE_BYTES: &str = "write.parquet.dict-size-bytes"; + /// Default Parquet dictionary page size in bytes. + pub const PROPERTY_PARQUET_DICT_SIZE_BYTES_DEFAULT: usize = 2 * 1024 * 1024; + /// Property key for the master key id used to encrypt the table's manifest /// list and data files as defined in https://iceberg.apache.org/docs/nightly/encryption/. pub const PROPERTY_ENCRYPTION_KEY_ID: &str = "encryption.key-id"; @@ -445,6 +543,27 @@ impl TryFrom<&HashMap> for TableProperties { TableProperties::PROPERTY_PARQUET_CDC_NORM_LEVEL, TableProperties::PROPERTY_PARQUET_CDC_NORM_LEVEL_DEFAULT, )?, + parquet_compression_codec: parse_parquet_compression(props)?, + parquet_row_group_size_bytes: parse_property( + props, + TableProperties::PROPERTY_PARQUET_ROW_GROUP_SIZE_BYTES, + TableProperties::PROPERTY_PARQUET_ROW_GROUP_SIZE_BYTES_DEFAULT, + )?, + parquet_page_size_bytes: parse_property( + props, + TableProperties::PROPERTY_PARQUET_PAGE_SIZE_BYTES, + TableProperties::PROPERTY_PARQUET_PAGE_SIZE_BYTES_DEFAULT, + )?, + parquet_page_row_limit: parse_property( + props, + TableProperties::PROPERTY_PARQUET_PAGE_ROW_LIMIT, + TableProperties::PROPERTY_PARQUET_PAGE_ROW_LIMIT_DEFAULT, + )?, + parquet_dict_size_bytes: parse_property( + props, + TableProperties::PROPERTY_PARQUET_DICT_SIZE_BYTES, + TableProperties::PROPERTY_PARQUET_DICT_SIZE_BYTES_DEFAULT, + )?, encryption_key_id: props .get(TableProperties::PROPERTY_ENCRYPTION_KEY_ID) .cloned(), @@ -978,6 +1097,124 @@ mod tests { assert!(!tp.cdc_enabled); } + #[test] + fn test_parquet_sizing_defaults() { + let tp = TableProperties::try_from(&HashMap::new()).unwrap(); + // Default codec is zstd at its default level. + assert_eq!( + tp.parquet_compression_codec, + CompressionCodec::zstd_default() + ); + assert_eq!( + tp.parquet_row_group_size_bytes, + TableProperties::PROPERTY_PARQUET_ROW_GROUP_SIZE_BYTES_DEFAULT + ); + assert_eq!( + tp.parquet_page_size_bytes, + TableProperties::PROPERTY_PARQUET_PAGE_SIZE_BYTES_DEFAULT + ); + assert_eq!( + tp.parquet_page_row_limit, + TableProperties::PROPERTY_PARQUET_PAGE_ROW_LIMIT_DEFAULT + ); + assert_eq!( + tp.parquet_dict_size_bytes, + TableProperties::PROPERTY_PARQUET_DICT_SIZE_BYTES_DEFAULT + ); + } + + #[test] + fn test_parquet_sizing_overrides() { + let props = HashMap::from([ + ( + TableProperties::PROPERTY_PARQUET_COMPRESSION_CODEC.to_string(), + "gzip".to_string(), + ), + ( + TableProperties::PROPERTY_PARQUET_COMPRESSION_LEVEL.to_string(), + "4".to_string(), + ), + ( + TableProperties::PROPERTY_PARQUET_ROW_GROUP_SIZE_BYTES.to_string(), + "1048576".to_string(), + ), + ( + TableProperties::PROPERTY_PARQUET_PAGE_SIZE_BYTES.to_string(), + "65536".to_string(), + ), + ( + TableProperties::PROPERTY_PARQUET_PAGE_ROW_LIMIT.to_string(), + "5000".to_string(), + ), + ( + TableProperties::PROPERTY_PARQUET_DICT_SIZE_BYTES.to_string(), + "131072".to_string(), + ), + ]); + let tp = TableProperties::try_from(&props).unwrap(); + // Codec name and level are folded into a single CompressionCodec. + assert_eq!(tp.parquet_compression_codec, CompressionCodec::Gzip(4)); + assert_eq!(tp.parquet_row_group_size_bytes, 1048576); + assert_eq!(tp.parquet_page_size_bytes, 65536); + assert_eq!(tp.parquet_page_row_limit, 5000); + assert_eq!(tp.parquet_dict_size_bytes, 131072); + } + + #[test] + fn test_parquet_invalid_sizing_rejected() { + let props = HashMap::from([( + TableProperties::PROPERTY_PARQUET_ROW_GROUP_SIZE_BYTES.to_string(), + "not_a_number".to_string(), + )]); + let err = TableProperties::try_from(&props).unwrap_err(); + assert_eq!(err.kind(), ErrorKind::DataInvalid); + assert!( + err.to_string() + .contains(TableProperties::PROPERTY_PARQUET_ROW_GROUP_SIZE_BYTES) + ); + } + + #[test] + fn test_parquet_all_codecs_parse() { + // Every codec name parquet-java supports must parse (parity with Java's + // `CompressionCodecName.valueOf`). + for (name, expected) in [ + ("uncompressed", CompressionCodec::None), + ("snappy", CompressionCodec::Snappy), + ("gzip", CompressionCodec::gzip_default()), + ("lzo", CompressionCodec::Lzo), + ("brotli", CompressionCodec::brotli_default()), + ("lz4", CompressionCodec::Lz4), + ("lz4_raw", CompressionCodec::Lz4Raw), + ("zstd", CompressionCodec::zstd_default()), + ] { + let props = HashMap::from([( + TableProperties::PROPERTY_PARQUET_COMPRESSION_CODEC.to_string(), + name.to_string(), + )]); + let tp = TableProperties::try_from(&props).unwrap(); + assert_eq!(tp.parquet_compression_codec, expected, "codec {name}"); + } + } + + #[test] + fn test_parquet_compression_level_ignored_for_levelless_codec() { + // A level set alongside a codec that carries none (e.g. snappy) is + // ignored rather than rejected, matching parquet-java. + let props = HashMap::from([ + ( + TableProperties::PROPERTY_PARQUET_COMPRESSION_CODEC.to_string(), + "snappy".to_string(), + ), + ( + TableProperties::PROPERTY_PARQUET_COMPRESSION_LEVEL.to_string(), + "5".to_string(), + ), + ]); + let tp = TableProperties::try_from(&props).unwrap(); + assert_eq!(tp.parquet_compression_codec, CompressionCodec::Snappy); + } + #[test] fn test_parse_boolean_property_case_insensitive() { let false_variants = ["False", "FALSE"]; diff --git a/crates/iceberg/src/writer/file_writer/parquet_writer.rs b/crates/iceberg/src/writer/file_writer/parquet_writer.rs index db9f170938..4e70520abe 100644 --- a/crates/iceberg/src/writer/file_writer/parquet_writer.rs +++ b/crates/iceberg/src/writer/file_writer/parquet_writer.rs @@ -27,6 +27,7 @@ use itertools::Itertools; use parquet::arrow::AsyncArrowWriter; use parquet::arrow::async_reader::AsyncFileReader; use parquet::arrow::async_writer::AsyncFileWriter as ArrowAsyncFileWriter; +use parquet::basic::{BrotliLevel, Compression, GzipLevel, ZstdLevel}; use parquet::file::metadata::ParquetMetaData; use parquet::file::properties::{CdcOptions, WriterProperties}; use parquet::file::statistics::Statistics; @@ -36,6 +37,7 @@ use crate::arrow::{ ArrowFileReader, DEFAULT_MAP_FIELD_NAME, FieldMatchMode, NanValueCountVisitor, get_parquet_stat_max_as_datum, get_parquet_stat_min_as_datum, }; +use crate::compression::CompressionCodec; use crate::io::{FileIO, FileWrite, OutputFile}; use crate::spec::{ DataContentType, DataFileBuilder, DataFileFormat, Datum, ListType, Literal, MapType, @@ -81,23 +83,22 @@ impl ParquetWriterBuilder { /// Build a `ParquetWriterBuilder` from Iceberg table properties and a /// schema, translating `write.parquet.*` settings into `WriterProperties` /// instead of using parquet-rs defaults. - /// - /// Currently translates the content-defined-chunking keys - /// (`write.parquet.content-defined-chunking.*`); other keys fall back to - /// parquet-rs defaults. - pub fn from_table_properties(table_props: &TableProperties, schema: SchemaRef) -> Self { + pub fn from_table_properties(table_props: &TableProperties, schema: SchemaRef) -> Result { let cdc = table_props.cdc_enabled.then_some(CdcOptions { min_chunk_size: table_props.cdc_min_chunk_size, max_chunk_size: table_props.cdc_max_chunk_size, norm_level: table_props.cdc_norm_level, }); - // TODO: translate the remaining write.parquet.* keys (e.g. compression-codec, - // row-group-size-bytes, page-size-bytes). - // This constructor is intended to be the single place that maps them. + let compression = parquet_compression(table_props.parquet_compression_codec)?; let props = WriterProperties::builder() .set_content_defined_chunking(cdc) + .set_compression(compression) + .set_max_row_group_bytes(Some(table_props.parquet_row_group_size_bytes)) + .set_data_page_size_limit(table_props.parquet_page_size_bytes) + .set_data_page_row_count_limit(table_props.parquet_page_row_limit) + .set_dictionary_page_size_limit(table_props.parquet_dict_size_bytes) .build(); - Self::new_with_match_mode(props, schema, FieldMatchMode::Id) + Ok(Self::new_with_match_mode(props, schema, FieldMatchMode::Id)) } /// Set the field match mode used to map Arrow fields to Iceberg fields. @@ -110,6 +111,40 @@ impl ParquetWriterBuilder { } } +fn parquet_compression(codec: CompressionCodec) -> Result { + let compression = match codec { + CompressionCodec::None => Compression::UNCOMPRESSED, + CompressionCodec::Snappy => Compression::SNAPPY, + CompressionCodec::Lzo => Compression::LZO, + CompressionCodec::Lz4 => Compression::LZ4, + CompressionCodec::Lz4Raw => Compression::LZ4_RAW, + CompressionCodec::Zstd(level) => { + let level = + ZstdLevel::try_new(level as i32).map_err(|e| invalid_level_error("zstd", e))?; + Compression::ZSTD(level) + } + CompressionCodec::Gzip(level) => { + let level = + GzipLevel::try_new(level as u32).map_err(|e| invalid_level_error("gzip", e))?; + Compression::GZIP(level) + } + CompressionCodec::Brotli(level) => { + let level = + BrotliLevel::try_new(level as u32).map_err(|e| invalid_level_error("brotli", e))?; + Compression::BROTLI(level) + } + }; + Ok(compression) +} + +fn invalid_level_error(codec: &str, source: impl Into) -> Error { + Error::new( + ErrorKind::DataInvalid, + format!("Invalid {codec} compression level"), + ) + .with_source(source) +} + impl FileWriterBuilder for ParquetWriterBuilder { type R = ParquetWriter; @@ -663,7 +698,9 @@ mod tests { use arrow_schema::{DataType, Field, Fields, SchemaRef as ArrowSchemaRef}; use arrow_select::concat::concat_batches; use parquet::arrow::PARQUET_FIELD_ID_META_KEY; + use parquet::basic::{BrotliLevel, Compression, GzipLevel, ZstdLevel}; use parquet::file::statistics::ValueStatistics; + use parquet::schema::types::ColumnPath; use tempfile::TempDir; use uuid::Uuid; @@ -2341,10 +2378,6 @@ mod tests { assert_eq!(upper_bounds, HashMap::from([(0, Datum::int(i32::MAX))])); } - // ----------------------------------------------------------------- - // ParquetWriterBuilder::from_table_properties - // ----------------------------------------------------------------- - fn cdc_test_schema() -> SchemaRef { Arc::new( Schema::builder() @@ -2366,7 +2399,7 @@ mod tests { #[test] fn test_from_table_properties_no_cdc_by_default() { let tp = table_props(HashMap::new()); - let builder = ParquetWriterBuilder::from_table_properties(&tp, cdc_test_schema()); + let builder = ParquetWriterBuilder::from_table_properties(&tp, cdc_test_schema()).unwrap(); assert!(builder.props.content_defined_chunking().is_none()); } @@ -2404,6 +2437,7 @@ mod tests { .new_output(format!("{}/cdc.parquet", tmp.path().to_str().unwrap())) .unwrap(); let writer = ParquetWriterBuilder::from_table_properties(&tp, cdc_test_schema()) + .unwrap() .build(output) .await .unwrap(); @@ -2417,4 +2451,142 @@ mod tests { assert_eq!(cdc.max_chunk_size, 8192); assert_eq!(cdc.norm_level, 2); } + + #[test] + fn test_from_table_properties_sizing_defaults() { + // With no properties set, the writer must use Iceberg's defaults (which + // differ from parquet-rs's own defaults), not parquet-rs's. + let tp = table_props(HashMap::new()); + let props = ParquetWriterBuilder::from_table_properties(&tp, cdc_test_schema()) + .unwrap() + .props; + + assert_eq!( + props.max_row_group_bytes(), + Some(TableProperties::PROPERTY_PARQUET_ROW_GROUP_SIZE_BYTES_DEFAULT) + ); + assert_eq!( + props.data_page_size_limit(), + TableProperties::PROPERTY_PARQUET_PAGE_SIZE_BYTES_DEFAULT + ); + assert_eq!( + props.data_page_row_count_limit(), + TableProperties::PROPERTY_PARQUET_PAGE_ROW_LIMIT_DEFAULT + ); + assert_eq!( + props.dictionary_page_size_limit(), + TableProperties::PROPERTY_PARQUET_DICT_SIZE_BYTES_DEFAULT + ); + // Default codec is zstd at the Java-aligned default level (3). + assert_eq!( + props.compression(&ColumnPath::from("id")), + Compression::ZSTD(ZstdLevel::try_new(3).unwrap()) + ); + } + + #[test] + fn test_from_table_properties_sizing_and_compression_overrides() { + let tp = table_props(HashMap::from([ + ( + TableProperties::PROPERTY_PARQUET_ROW_GROUP_SIZE_BYTES.to_string(), + "1048576".to_string(), + ), + ( + TableProperties::PROPERTY_PARQUET_PAGE_SIZE_BYTES.to_string(), + "65536".to_string(), + ), + ( + TableProperties::PROPERTY_PARQUET_PAGE_ROW_LIMIT.to_string(), + "5000".to_string(), + ), + ( + TableProperties::PROPERTY_PARQUET_DICT_SIZE_BYTES.to_string(), + "131072".to_string(), + ), + ( + TableProperties::PROPERTY_PARQUET_COMPRESSION_CODEC.to_string(), + "gzip".to_string(), + ), + ( + TableProperties::PROPERTY_PARQUET_COMPRESSION_LEVEL.to_string(), + "9".to_string(), + ), + ])); + let props = ParquetWriterBuilder::from_table_properties(&tp, cdc_test_schema()) + .unwrap() + .props; + + assert_eq!(props.max_row_group_bytes(), Some(1048576)); + assert_eq!(props.data_page_size_limit(), 65536); + assert_eq!(props.data_page_row_count_limit(), 5000); + assert_eq!(props.dictionary_page_size_limit(), 131072); + assert_eq!( + props.compression(&ColumnPath::from("id")), + Compression::GZIP(GzipLevel::try_new(9).unwrap()) + ); + } + + #[test] + fn test_from_table_properties_invalid_codec_errors() { + let entries = HashMap::from([( + TableProperties::PROPERTY_PARQUET_COMPRESSION_CODEC.to_string(), + "bogus".to_string(), + )]); + let err = TableProperties::try_from(&entries).unwrap_err(); + assert_eq!(err.kind(), ErrorKind::DataInvalid); + assert!(err.to_string().contains("bogus")); + } + + #[test] + fn test_parquet_compression_mapping() { + // Codecs without a level. + assert_eq!( + parquet_compression(CompressionCodec::None).unwrap(), + Compression::UNCOMPRESSED + ); + assert_eq!( + parquet_compression(CompressionCodec::Snappy).unwrap(), + Compression::SNAPPY + ); + assert_eq!( + parquet_compression(CompressionCodec::Lz4).unwrap(), + Compression::LZ4 + ); + assert_eq!( + parquet_compression(CompressionCodec::Lz4Raw).unwrap(), + Compression::LZ4_RAW + ); + assert_eq!( + parquet_compression(CompressionCodec::Lzo).unwrap(), + Compression::LZO + ); + + // Level-carrying codecs at their defaults. + assert_eq!( + parquet_compression(CompressionCodec::zstd_default()).unwrap(), + Compression::ZSTD(ZstdLevel::try_new(3).unwrap()) + ); + assert_eq!( + parquet_compression(CompressionCodec::gzip_default()).unwrap(), + Compression::GZIP(GzipLevel::try_new(6).unwrap()) + ); + assert_eq!( + parquet_compression(CompressionCodec::brotli_default()).unwrap(), + Compression::BROTLI(BrotliLevel::try_new(1).unwrap()) + ); + + // Explicit levels are honored. + assert_eq!( + parquet_compression(CompressionCodec::Zstd(10)).unwrap(), + Compression::ZSTD(ZstdLevel::try_new(10).unwrap()) + ); + } + + #[test] + fn test_parquet_compression_invalid_level() { + // zstd valid range is 1..=22; 99 is out of range. + let err = parquet_compression(CompressionCodec::Zstd(99)).unwrap_err(); + assert_eq!(err.kind(), ErrorKind::DataInvalid); + assert!(err.to_string().contains("zstd")); + } } diff --git a/crates/integrations/datafusion/src/physical_plan/write.rs b/crates/integrations/datafusion/src/physical_plan/write.rs index a7d771ec1b..8f48b9d4cc 100644 --- a/crates/integrations/datafusion/src/physical_plan/write.rs +++ b/crates/integrations/datafusion/src/physical_plan/write.rs @@ -226,6 +226,7 @@ impl ExecutionPlan for IcebergWriteExec { &table_props, self.table.metadata().current_schema().clone(), ) + .map_err(to_datafusion_error)? .with_match_mode(FieldMatchMode::Name); let target_file_size = table_props.write_target_file_size_bytes;