From ca7d5ba577ec3b2d473d911638ecdcb837c9a9b5 Mon Sep 17 00:00:00 2001 From: Xander Date: Fri, 24 Jul 2026 13:20:56 +0100 Subject: [PATCH 1/6] feat(paquet): wire in additional parquet writer settings --- crates/iceberg/public-api.txt | 19 +- crates/iceberg/src/spec/table_properties.rs | 172 +++++++++++++ .../src/writer/file_writer/parquet_writer.rs | 228 +++++++++++++++++- .../datafusion/src/physical_plan/write.rs | 1 + 4 files changed, 407 insertions(+), 13 deletions(-) diff --git a/crates/iceberg/public-api.txt b/crates/iceberg/public-api.txt index f5a54304df..1d5ef2f8c7 100644 --- a/crates/iceberg/public-api.txt +++ b/crates/iceberg/public-api.txt @@ -2750,6 +2750,12 @@ 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: alloc::string::String +pub iceberg::spec::TableProperties::parquet_compression_level: core::option::Option +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_datafusion_fanout_enabled: bool pub iceberg::spec::TableProperties::write_format_default: alloc::string::String pub iceberg::spec::TableProperties::write_metadata_path: core::option::Option @@ -2798,6 +2804,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_METADATA_PATH: &str @@ -3270,7 +3287,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/spec/table_properties.rs b/crates/iceberg/src/spec/table_properties.rs index 379feee5c1..f7915916e6 100644 --- a/crates/iceberg/src/spec/table_properties.rs +++ b/crates/iceberg/src/spec/table_properties.rs @@ -40,6 +40,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() +} + /// Strips trailing slashes from a location, preserving a bare URI scheme root fn strip_trailing_slash(path: &str) -> &str { let mut path = path; @@ -168,6 +190,20 @@ 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 name (e.g. `zstd`, `gzip`). Validated when the + /// writer is built, not when properties are parsed. + pub parquet_compression_codec: String, + /// Parquet compression level for codecs that accept one. `None` uses the + /// codec's default level. + pub parquet_compression_level: Option, + /// 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, @@ -317,6 +353,36 @@ 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`, `brotli`, `uncompressed`). The codec name is validated when the + /// writer is built, not when properties are parsed. + 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"; @@ -413,6 +479,35 @@ impl TryFrom<&HashMap> for TableProperties { TableProperties::PROPERTY_PARQUET_CDC_NORM_LEVEL, TableProperties::PROPERTY_PARQUET_CDC_NORM_LEVEL_DEFAULT, )?, + parquet_compression_codec: parse_property( + props, + TableProperties::PROPERTY_PARQUET_COMPRESSION_CODEC, + TableProperties::PROPERTY_PARQUET_COMPRESSION_CODEC_DEFAULT.to_string(), + )?, + parquet_compression_level: parse_optional_property( + props, + TableProperties::PROPERTY_PARQUET_COMPRESSION_LEVEL, + )?, + 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(), @@ -947,4 +1042,81 @@ mod tests { let tp = TableProperties::try_from(&props).unwrap(); assert!(!tp.cdc_enabled); } + + #[test] + fn test_parquet_sizing_defaults() { + let tp = TableProperties::try_from(&HashMap::new()).unwrap(); + assert_eq!( + tp.parquet_compression_codec, + TableProperties::PROPERTY_PARQUET_COMPRESSION_CODEC_DEFAULT + ); + assert_eq!(tp.parquet_compression_level, None); + 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(); + assert_eq!(tp.parquet_compression_codec, "gzip"); + assert_eq!(tp.parquet_compression_level, Some(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) + ); + } } diff --git a/crates/iceberg/src/writer/file_writer/parquet_writer.rs b/crates/iceberg/src/writer/file_writer/parquet_writer.rs index db9f170938..a1f583d17b 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; @@ -82,22 +83,26 @@ impl ParquetWriterBuilder { /// 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, + table_props.parquet_compression_level, + )?; 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 +115,63 @@ impl ParquetWriterBuilder { } } +fn parquet_compression(codec: &str, level: Option) -> Result { + let compression = match codec.to_lowercase().as_str() { + "uncompressed" | "none" => Compression::UNCOMPRESSED, + "snappy" => Compression::SNAPPY, + "lzo" => Compression::LZO, + "lz4" => Compression::LZ4, + "lz4_raw" => Compression::LZ4_RAW, + "gzip" => { + let level = match level { + Some(l) => GzipLevel::try_new(level_as_u32("gzip", l)?) + .map_err(|e| invalid_level_error("gzip", e))?, + None => GzipLevel::default(), + }; + Compression::GZIP(level) + } + "brotli" => { + let level = match level { + Some(l) => BrotliLevel::try_new(level_as_u32("brotli", l)?) + .map_err(|e| invalid_level_error("brotli", e))?, + None => BrotliLevel::default(), + }; + Compression::BROTLI(level) + } + "zstd" => { + let level = match level { + Some(l) => ZstdLevel::try_new(l).map_err(|e| invalid_level_error("zstd", e))?, + None => ZstdLevel::default(), + }; + Compression::ZSTD(level) + } + other => { + return Err(Error::new( + ErrorKind::DataInvalid, + format!( + "Unsupported Parquet compression codec: {other}. Supported codecs: \ + uncompressed, snappy, gzip, lzo, brotli, lz4, lz4_raw, zstd" + ), + )); + } + }; + Ok(compression) +} + +fn invalid_level_error(codec: &str, source: impl Into) -> Error { + Error::new( + ErrorKind::DataInvalid, + format!("Invalid {codec} compression level"), + ) + .with_source(source) +} + +/// Resolve a `u32` compression level for the gzip/brotli codecs, whose +/// parquet-rs levels are unsigned. A negative level is invalid. +fn level_as_u32(codec: &str, level: i32) -> Result { + u32::try_from(level).map_err(|e| invalid_level_error(codec, e)) +} + impl FileWriterBuilder for ParquetWriterBuilder { type R = ParquetWriter; @@ -664,6 +726,7 @@ mod tests { use arrow_select::concat::concat_batches; use parquet::arrow::PARQUET_FIELD_ID_META_KEY; use parquet::file::statistics::ValueStatistics; + use parquet::schema::types::ColumnPath; use tempfile::TempDir; use uuid::Uuid; @@ -2341,10 +2404,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 +2425,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 +2463,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 +2477,148 @@ 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 parquet-rs's default level. + assert_eq!( + props.compression(&ColumnPath::from("id")), + Compression::ZSTD(ZstdLevel::default()) + ); + } + + #[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 tp = table_props(HashMap::from([( + TableProperties::PROPERTY_PARQUET_COMPRESSION_CODEC.to_string(), + "bogus".to_string(), + )])); + let err = ParquetWriterBuilder::from_table_properties(&tp, cdc_test_schema()).unwrap_err(); + assert_eq!(err.kind(), ErrorKind::DataInvalid); + assert!(err.to_string().contains("bogus")); + } + + #[test] + fn test_parquet_compression_codec_mapping() { + // Codecs without a level. + assert_eq!( + parquet_compression("uncompressed", None).unwrap(), + Compression::UNCOMPRESSED + ); + assert_eq!( + parquet_compression("snappy", None).unwrap(), + Compression::SNAPPY + ); + assert_eq!(parquet_compression("lz4", None).unwrap(), Compression::LZ4); + assert_eq!( + parquet_compression("lz4_raw", None).unwrap(), + Compression::LZ4_RAW + ); + assert_eq!(parquet_compression("lzo", None).unwrap(), Compression::LZO); + + // Case-insensitive codec names. + assert_eq!( + parquet_compression("ZSTD", None).unwrap(), + Compression::ZSTD(ZstdLevel::default()) + ); + + // Level-taking codecs use their default level when none is supplied. + assert_eq!( + parquet_compression("gzip", None).unwrap(), + Compression::GZIP(GzipLevel::default()) + ); + assert_eq!( + parquet_compression("brotli", None).unwrap(), + Compression::BROTLI(BrotliLevel::default()) + ); + + // Explicit levels are honored. + assert_eq!( + parquet_compression("zstd", Some(10)).unwrap(), + Compression::ZSTD(ZstdLevel::try_new(10).unwrap()) + ); + } + + #[test] + fn test_parquet_compression_invalid_inputs() { + // Unknown codec. + assert_eq!( + parquet_compression("bogus", None).unwrap_err().kind(), + ErrorKind::DataInvalid + ); + // Level out of range for zstd (valid range is 1..=22). + assert_eq!( + parquet_compression("zstd", Some(99)).unwrap_err().kind(), + ErrorKind::DataInvalid + ); + // Negative level for a u32-based codec. + let err = parquet_compression("gzip", Some(-1)).unwrap_err(); + assert_eq!(err.kind(), ErrorKind::DataInvalid); + assert!(err.to_string().contains("gzip")); + } } 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; From 67c28d845f55d82f04f216413ff6d2e03ed3ff06 Mon Sep 17 00:00:00 2001 From: Xander Date: Fri, 31 Jul 2026 00:45:41 +0100 Subject: [PATCH 2/6] fix compression level --- .../src/writer/file_writer/parquet_writer.rs | 44 +++++++++++-------- 1 file changed, 26 insertions(+), 18 deletions(-) diff --git a/crates/iceberg/src/writer/file_writer/parquet_writer.rs b/crates/iceberg/src/writer/file_writer/parquet_writer.rs index a1f583d17b..a2e442d799 100644 --- a/crates/iceberg/src/writer/file_writer/parquet_writer.rs +++ b/crates/iceberg/src/writer/file_writer/parquet_writer.rs @@ -47,6 +47,16 @@ use crate::transform::create_transform_function; use crate::writer::{CurrentFileStatus, DataFile}; use crate::{Error, ErrorKind, Result}; +/// Default compression levels, pinned to parquet-java's defaults so files are +/// comparable across implementations rather than tracking the parquet-rs +/// defaults, which may differ and could change without us noticing. +/// Default zstd level (parquet-rs uses 1). +const DEFAULT_ZSTD_COMPRESSION_LEVEL: i32 = 3; +/// Default gzip level. +const DEFAULT_GZIP_COMPRESSION_LEVEL: u32 = 6; +/// Default brotli level. +const DEFAULT_BROTLI_COMPRESSION_LEVEL: u32 = 1; + /// ParquetWriterBuilder is used to builder a [`ParquetWriter`] #[derive(Clone, Debug)] pub struct ParquetWriterBuilder { @@ -124,25 +134,23 @@ fn parquet_compression(codec: &str, level: Option) -> Result { "lz4_raw" => Compression::LZ4_RAW, "gzip" => { let level = match level { - Some(l) => GzipLevel::try_new(level_as_u32("gzip", l)?) - .map_err(|e| invalid_level_error("gzip", e))?, - None => GzipLevel::default(), + Some(l) => level_as_u32("gzip", l)?, + None => DEFAULT_GZIP_COMPRESSION_LEVEL, }; + let level = GzipLevel::try_new(level).map_err(|e| invalid_level_error("gzip", e))?; Compression::GZIP(level) } "brotli" => { let level = match level { - Some(l) => BrotliLevel::try_new(level_as_u32("brotli", l)?) - .map_err(|e| invalid_level_error("brotli", e))?, - None => BrotliLevel::default(), + Some(l) => level_as_u32("brotli", l)?, + None => DEFAULT_BROTLI_COMPRESSION_LEVEL, }; + let level = BrotliLevel::try_new(level).map_err(|e| invalid_level_error("brotli", e))?; Compression::BROTLI(level) } "zstd" => { - let level = match level { - Some(l) => ZstdLevel::try_new(l).map_err(|e| invalid_level_error("zstd", e))?, - None => ZstdLevel::default(), - }; + let level = level.unwrap_or(DEFAULT_ZSTD_COMPRESSION_LEVEL); + let level = ZstdLevel::try_new(level).map_err(|e| invalid_level_error("zstd", e))?; Compression::ZSTD(level) } other => { @@ -2503,10 +2511,10 @@ mod tests { props.dictionary_page_size_limit(), TableProperties::PROPERTY_PARQUET_DICT_SIZE_BYTES_DEFAULT ); - // Default codec is zstd at parquet-rs's default level. + // Default codec is zstd at the Java-aligned default level (3). assert_eq!( props.compression(&ColumnPath::from("id")), - Compression::ZSTD(ZstdLevel::default()) + Compression::ZSTD(ZstdLevel::try_new(DEFAULT_ZSTD_COMPRESSION_LEVEL).unwrap()) ); } @@ -2581,20 +2589,20 @@ mod tests { ); assert_eq!(parquet_compression("lzo", None).unwrap(), Compression::LZO); - // Case-insensitive codec names. + // Case-insensitive codec names. With no level, each codec uses its + // parquet-java-aligned default, pinned explicitly rather than inherited + // from the parquet-rs. assert_eq!( parquet_compression("ZSTD", None).unwrap(), - Compression::ZSTD(ZstdLevel::default()) + Compression::ZSTD(ZstdLevel::try_new(DEFAULT_ZSTD_COMPRESSION_LEVEL).unwrap()) ); - - // Level-taking codecs use their default level when none is supplied. assert_eq!( parquet_compression("gzip", None).unwrap(), - Compression::GZIP(GzipLevel::default()) + Compression::GZIP(GzipLevel::try_new(DEFAULT_GZIP_COMPRESSION_LEVEL).unwrap()) ); assert_eq!( parquet_compression("brotli", None).unwrap(), - Compression::BROTLI(BrotliLevel::default()) + Compression::BROTLI(BrotliLevel::try_new(DEFAULT_BROTLI_COMPRESSION_LEVEL).unwrap()) ); // Explicit levels are honored. From da6cb46307f16aabcbac851f465ecdced405e7e8 Mon Sep 17 00:00:00 2001 From: Xander Date: Fri, 31 Jul 2026 00:55:35 +0100 Subject: [PATCH 3/6] fmt --- crates/iceberg/src/writer/file_writer/parquet_writer.rs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/crates/iceberg/src/writer/file_writer/parquet_writer.rs b/crates/iceberg/src/writer/file_writer/parquet_writer.rs index a2e442d799..bab6c9cd38 100644 --- a/crates/iceberg/src/writer/file_writer/parquet_writer.rs +++ b/crates/iceberg/src/writer/file_writer/parquet_writer.rs @@ -145,7 +145,8 @@ fn parquet_compression(codec: &str, level: Option) -> Result { Some(l) => level_as_u32("brotli", l)?, None => DEFAULT_BROTLI_COMPRESSION_LEVEL, }; - let level = BrotliLevel::try_new(level).map_err(|e| invalid_level_error("brotli", e))?; + let level = + BrotliLevel::try_new(level).map_err(|e| invalid_level_error("brotli", e))?; Compression::BROTLI(level) } "zstd" => { From b83cf89ba8c7f399231cd11f3d0578cbd07aa64d Mon Sep 17 00:00:00 2001 From: Xander Date: Mon, 3 Aug 2026 07:51:10 -0600 Subject: [PATCH 4/6] move comments --- crates/iceberg/src/writer/file_writer/parquet_writer.rs | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/crates/iceberg/src/writer/file_writer/parquet_writer.rs b/crates/iceberg/src/writer/file_writer/parquet_writer.rs index bab6c9cd38..e1aa6e085f 100644 --- a/crates/iceberg/src/writer/file_writer/parquet_writer.rs +++ b/crates/iceberg/src/writer/file_writer/parquet_writer.rs @@ -47,9 +47,10 @@ use crate::transform::create_transform_function; use crate::writer::{CurrentFileStatus, DataFile}; use crate::{Error, ErrorKind, Result}; -/// Default compression levels, pinned to parquet-java's defaults so files are -/// comparable across implementations rather than tracking the parquet-rs -/// defaults, which may differ and could change without us noticing. +// Default compression levels are pinned to parquet-java's defaults so files are +// comparable across implementations, rather than tracking the parquet-rs +// defaults, which may differ and could change without us noticing. + /// Default zstd level (parquet-rs uses 1). const DEFAULT_ZSTD_COMPRESSION_LEVEL: i32 = 3; /// Default gzip level. @@ -92,8 +93,6 @@ impl ParquetWriterBuilder { /// Build a `ParquetWriterBuilder` from Iceberg table properties and a /// schema, translating `write.parquet.*` settings into `WriterProperties` /// instead of using parquet-rs defaults. - /// - /// parquet-rs defaults. 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, From 6b74b7254f02917650c9c0e67f708ac9ae4aede0 Mon Sep 17 00:00:00 2001 From: Xander Date: Tue, 4 Aug 2026 08:49:43 -0600 Subject: [PATCH 5/6] move to codec mod --- crates/iceberg/public-api.txt | 11 +- crates/iceberg/src/compression.rs | 62 +++++++- crates/iceberg/src/spec/table_properties.rs | 107 +++++++++++--- .../src/writer/file_writer/parquet_writer.rs | 136 +++++++----------- 4 files changed, 203 insertions(+), 113 deletions(-) diff --git a/crates/iceberg/public-api.txt b/crates/iceberg/public-api.txt index 7c99d395cd..cd0ee8a951 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 @@ -2762,8 +2770,7 @@ 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: alloc::string::String -pub iceberg::spec::TableProperties::parquet_compression_level: core::option::Option +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 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 f7915916e6..7cd0c28553 100644 --- a/crates/iceberg/src/spec/table_properties.rs +++ b/crates/iceberg/src/spec/table_properties.rs @@ -150,6 +150,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, + }) +} + /// TableProperties that contains the properties of a table. #[derive(Debug)] pub struct TableProperties { @@ -190,12 +223,10 @@ 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 name (e.g. `zstd`, `gzip`). Validated when the - /// writer is built, not when properties are parsed. - pub parquet_compression_codec: String, - /// Parquet compression level for codecs that accept one. `None` uses the - /// codec's default level. - pub parquet_compression_level: Option, + /// 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. @@ -354,8 +385,9 @@ impl TableProperties { pub const PROPERTY_PARQUET_CDC_NORM_LEVEL_DEFAULT: i32 = 0; /// Compression codec for Parquet data files (e.g. `zstd`, `gzip`, `snappy`, - /// `lz4`, `brotli`, `uncompressed`). The codec name is validated when the - /// writer is built, not when properties are parsed. + /// `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"; @@ -479,15 +511,7 @@ impl TryFrom<&HashMap> for TableProperties { TableProperties::PROPERTY_PARQUET_CDC_NORM_LEVEL, TableProperties::PROPERTY_PARQUET_CDC_NORM_LEVEL_DEFAULT, )?, - parquet_compression_codec: parse_property( - props, - TableProperties::PROPERTY_PARQUET_COMPRESSION_CODEC, - TableProperties::PROPERTY_PARQUET_COMPRESSION_CODEC_DEFAULT.to_string(), - )?, - parquet_compression_level: parse_optional_property( - props, - TableProperties::PROPERTY_PARQUET_COMPRESSION_LEVEL, - )?, + parquet_compression_codec: parse_parquet_compression(props)?, parquet_row_group_size_bytes: parse_property( props, TableProperties::PROPERTY_PARQUET_ROW_GROUP_SIZE_BYTES, @@ -1046,11 +1070,11 @@ mod tests { #[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, - TableProperties::PROPERTY_PARQUET_COMPRESSION_CODEC_DEFAULT + CompressionCodec::zstd_default() ); - assert_eq!(tp.parquet_compression_level, None); assert_eq!( tp.parquet_row_group_size_bytes, TableProperties::PROPERTY_PARQUET_ROW_GROUP_SIZE_BYTES_DEFAULT @@ -1098,8 +1122,8 @@ mod tests { ), ]); let tp = TableProperties::try_from(&props).unwrap(); - assert_eq!(tp.parquet_compression_codec, "gzip"); - assert_eq!(tp.parquet_compression_level, Some(4)); + // 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); @@ -1119,4 +1143,45 @@ mod tests { .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); + } } diff --git a/crates/iceberg/src/writer/file_writer/parquet_writer.rs b/crates/iceberg/src/writer/file_writer/parquet_writer.rs index e1aa6e085f..4e70520abe 100644 --- a/crates/iceberg/src/writer/file_writer/parquet_writer.rs +++ b/crates/iceberg/src/writer/file_writer/parquet_writer.rs @@ -37,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, @@ -47,17 +48,6 @@ use crate::transform::create_transform_function; use crate::writer::{CurrentFileStatus, DataFile}; use crate::{Error, ErrorKind, Result}; -// Default compression levels are pinned to parquet-java's defaults so files are -// comparable across implementations, rather than tracking the parquet-rs -// defaults, which may differ and could change without us noticing. - -/// Default zstd level (parquet-rs uses 1). -const DEFAULT_ZSTD_COMPRESSION_LEVEL: i32 = 3; -/// Default gzip level. -const DEFAULT_GZIP_COMPRESSION_LEVEL: u32 = 6; -/// Default brotli level. -const DEFAULT_BROTLI_COMPRESSION_LEVEL: u32 = 1; - /// ParquetWriterBuilder is used to builder a [`ParquetWriter`] #[derive(Clone, Debug)] pub struct ParquetWriterBuilder { @@ -99,10 +89,7 @@ impl ParquetWriterBuilder { max_chunk_size: table_props.cdc_max_chunk_size, norm_level: table_props.cdc_norm_level, }); - let compression = parquet_compression( - &table_props.parquet_compression_codec, - table_props.parquet_compression_level, - )?; + let compression = parquet_compression(table_props.parquet_compression_codec)?; let props = WriterProperties::builder() .set_content_defined_chunking(cdc) .set_compression(compression) @@ -124,44 +111,28 @@ impl ParquetWriterBuilder { } } -fn parquet_compression(codec: &str, level: Option) -> Result { - let compression = match codec.to_lowercase().as_str() { - "uncompressed" | "none" => Compression::UNCOMPRESSED, - "snappy" => Compression::SNAPPY, - "lzo" => Compression::LZO, - "lz4" => Compression::LZ4, - "lz4_raw" => Compression::LZ4_RAW, - "gzip" => { - let level = match level { - Some(l) => level_as_u32("gzip", l)?, - None => DEFAULT_GZIP_COMPRESSION_LEVEL, - }; - let level = GzipLevel::try_new(level).map_err(|e| invalid_level_error("gzip", e))?; +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) } - "brotli" => { - let level = match level { - Some(l) => level_as_u32("brotli", l)?, - None => DEFAULT_BROTLI_COMPRESSION_LEVEL, - }; + CompressionCodec::Brotli(level) => { let level = - BrotliLevel::try_new(level).map_err(|e| invalid_level_error("brotli", e))?; + BrotliLevel::try_new(level as u32).map_err(|e| invalid_level_error("brotli", e))?; Compression::BROTLI(level) } - "zstd" => { - let level = level.unwrap_or(DEFAULT_ZSTD_COMPRESSION_LEVEL); - let level = ZstdLevel::try_new(level).map_err(|e| invalid_level_error("zstd", e))?; - Compression::ZSTD(level) - } - other => { - return Err(Error::new( - ErrorKind::DataInvalid, - format!( - "Unsupported Parquet compression codec: {other}. Supported codecs: \ - uncompressed, snappy, gzip, lzo, brotli, lz4, lz4_raw, zstd" - ), - )); - } }; Ok(compression) } @@ -174,12 +145,6 @@ fn invalid_level_error(codec: &str, source: impl Into) -> Error { .with_source(source) } -/// Resolve a `u32` compression level for the gzip/brotli codecs, whose -/// parquet-rs levels are unsigned. A negative level is invalid. -fn level_as_u32(codec: &str, level: i32) -> Result { - u32::try_from(level).map_err(|e| invalid_level_error(codec, e)) -} - impl FileWriterBuilder for ParquetWriterBuilder { type R = ParquetWriter; @@ -733,6 +698,7 @@ 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; @@ -2514,7 +2480,7 @@ mod tests { // Default codec is zstd at the Java-aligned default level (3). assert_eq!( props.compression(&ColumnPath::from("id")), - Compression::ZSTD(ZstdLevel::try_new(DEFAULT_ZSTD_COMPRESSION_LEVEL).unwrap()) + Compression::ZSTD(ZstdLevel::try_new(3).unwrap()) ); } @@ -2562,71 +2528,65 @@ mod tests { #[test] fn test_from_table_properties_invalid_codec_errors() { - let tp = table_props(HashMap::from([( + let entries = HashMap::from([( TableProperties::PROPERTY_PARQUET_COMPRESSION_CODEC.to_string(), "bogus".to_string(), - )])); - let err = ParquetWriterBuilder::from_table_properties(&tp, cdc_test_schema()).unwrap_err(); + )]); + 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_codec_mapping() { + fn test_parquet_compression_mapping() { // Codecs without a level. assert_eq!( - parquet_compression("uncompressed", None).unwrap(), + parquet_compression(CompressionCodec::None).unwrap(), Compression::UNCOMPRESSED ); assert_eq!( - parquet_compression("snappy", None).unwrap(), + parquet_compression(CompressionCodec::Snappy).unwrap(), Compression::SNAPPY ); - assert_eq!(parquet_compression("lz4", None).unwrap(), Compression::LZ4); assert_eq!( - parquet_compression("lz4_raw", None).unwrap(), + parquet_compression(CompressionCodec::Lz4).unwrap(), + Compression::LZ4 + ); + assert_eq!( + parquet_compression(CompressionCodec::Lz4Raw).unwrap(), Compression::LZ4_RAW ); - assert_eq!(parquet_compression("lzo", None).unwrap(), Compression::LZO); + assert_eq!( + parquet_compression(CompressionCodec::Lzo).unwrap(), + Compression::LZO + ); - // Case-insensitive codec names. With no level, each codec uses its - // parquet-java-aligned default, pinned explicitly rather than inherited - // from the parquet-rs. + // Level-carrying codecs at their defaults. assert_eq!( - parquet_compression("ZSTD", None).unwrap(), - Compression::ZSTD(ZstdLevel::try_new(DEFAULT_ZSTD_COMPRESSION_LEVEL).unwrap()) + parquet_compression(CompressionCodec::zstd_default()).unwrap(), + Compression::ZSTD(ZstdLevel::try_new(3).unwrap()) ); assert_eq!( - parquet_compression("gzip", None).unwrap(), - Compression::GZIP(GzipLevel::try_new(DEFAULT_GZIP_COMPRESSION_LEVEL).unwrap()) + parquet_compression(CompressionCodec::gzip_default()).unwrap(), + Compression::GZIP(GzipLevel::try_new(6).unwrap()) ); assert_eq!( - parquet_compression("brotli", None).unwrap(), - Compression::BROTLI(BrotliLevel::try_new(DEFAULT_BROTLI_COMPRESSION_LEVEL).unwrap()) + parquet_compression(CompressionCodec::brotli_default()).unwrap(), + Compression::BROTLI(BrotliLevel::try_new(1).unwrap()) ); // Explicit levels are honored. assert_eq!( - parquet_compression("zstd", Some(10)).unwrap(), + parquet_compression(CompressionCodec::Zstd(10)).unwrap(), Compression::ZSTD(ZstdLevel::try_new(10).unwrap()) ); } #[test] - fn test_parquet_compression_invalid_inputs() { - // Unknown codec. - assert_eq!( - parquet_compression("bogus", None).unwrap_err().kind(), - ErrorKind::DataInvalid - ); - // Level out of range for zstd (valid range is 1..=22). - assert_eq!( - parquet_compression("zstd", Some(99)).unwrap_err().kind(), - ErrorKind::DataInvalid - ); - // Negative level for a u32-based codec. - let err = parquet_compression("gzip", Some(-1)).unwrap_err(); + 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("gzip")); + assert!(err.to_string().contains("zstd")); } } From a337daf2621db2b2ef3d036636d2c5cf7dc90c4e Mon Sep 17 00:00:00 2001 From: Xander Date: Wed, 5 Aug 2026 07:29:11 -0600 Subject: [PATCH 6/6] fmt --- crates/iceberg/src/spec/table_properties.rs | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/crates/iceberg/src/spec/table_properties.rs b/crates/iceberg/src/spec/table_properties.rs index ec00674092..ae03319276 100644 --- a/crates/iceberg/src/spec/table_properties.rs +++ b/crates/iceberg/src/spec/table_properties.rs @@ -173,13 +173,13 @@ fn parse_parquet_compression(properties: &HashMap) -> Result, - key: &str, - default: bool, - ) -> Result { +/// 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" +fn parse_property_bool( + properties: &HashMap, + key: &str, + default: bool, +) -> Result { properties.get(key).map_or(Ok(default), |value| { value.to_lowercase().parse::().map_err(|e| { Error::new(