diff --git a/crates/catalog/rest/src/catalog.rs b/crates/catalog/rest/src/catalog.rs index aa36a074e7..b7d3b3163c 100644 --- a/crates/catalog/rest/src/catalog.rs +++ b/crates/catalog/rest/src/catalog.rs @@ -3832,13 +3832,14 @@ mod tests { .properties(HashMap::from([("owner".to_string(), "testx".to_string())])) .partition_spec( UnboundPartitionSpec::builder() - .add_partition_fields(vec![ + .add_partition_field( UnboundPartitionField::builder() - .source_id(1) + .source_ids(vec![1]) + .name("id") .transform(Transform::Truncate(3)) - .name("id".to_string()) - .build(), - ]) + .build() + .unwrap(), + ) .unwrap() .build(), ) diff --git a/crates/iceberg/public-api.txt b/crates/iceberg/public-api.txt index f29c5efe07..19a5d49b17 100644 --- a/crates/iceberg/public-api.txt +++ b/crates/iceberg/public-api.txt @@ -2911,10 +2911,13 @@ pub fn iceberg::spec::TableProperties<'properties>::write_target_file_size_bytes impl<'properties> core::fmt::Debug for iceberg::spec::TableProperties<'properties> pub fn iceberg::spec::TableProperties<'properties>::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result pub struct iceberg::spec::UnboundPartitionField -pub iceberg::spec::UnboundPartitionField::field_id: core::option::Option -pub iceberg::spec::UnboundPartitionField::name: alloc::string::String -pub iceberg::spec::UnboundPartitionField::source_id: i32 -pub iceberg::spec::UnboundPartitionField::transform: iceberg::spec::Transform +impl iceberg::spec::UnboundPartitionField +pub fn iceberg::spec::UnboundPartitionField::field_id(&self) -> core::option::Option +pub fn iceberg::spec::UnboundPartitionField::name(&self) -> &str +pub fn iceberg::spec::UnboundPartitionField::source_id(&self) -> iceberg::Result +pub fn iceberg::spec::UnboundPartitionField::source_ids(&self) -> &[i32] +pub fn iceberg::spec::UnboundPartitionField::transform(&self) -> iceberg::spec::Transform +pub fn iceberg::spec::UnboundPartitionField::with_field_id(self, field_id: i32) -> Self impl core::clone::Clone for iceberg::spec::UnboundPartitionField pub fn iceberg::spec::UnboundPartitionField::clone(&self) -> iceberg::spec::UnboundPartitionField impl core::cmp::Eq for iceberg::spec::UnboundPartitionField @@ -2922,6 +2925,8 @@ impl core::cmp::PartialEq for iceberg::spec::UnboundPartitionField pub fn iceberg::spec::UnboundPartitionField::eq(&self, other: &iceberg::spec::UnboundPartitionField) -> bool impl core::convert::From for iceberg::spec::UnboundPartitionField pub fn iceberg::spec::UnboundPartitionField::from(field: iceberg::spec::PartitionField) -> Self +impl core::convert::From for iceberg::Result +pub fn iceberg::Result::from(field: iceberg::spec::UnboundPartitionField) -> Self impl core::fmt::Debug for iceberg::spec::UnboundPartitionField pub fn iceberg::spec::UnboundPartitionField::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result impl core::marker::StructuralPartialEq for iceberg::spec::UnboundPartitionField @@ -2956,7 +2961,7 @@ impl<'de> serde_core::de::Deserialize<'de> for iceberg::spec::UnboundPartitionSp pub fn iceberg::spec::UnboundPartitionSpec::deserialize<__D>(__deserializer: __D) -> core::result::Result::Error> where __D: serde_core::de::Deserializer<'de> pub struct iceberg::spec::UnboundPartitionSpecBuilder impl iceberg::spec::UnboundPartitionSpecBuilder -pub fn iceberg::spec::UnboundPartitionSpecBuilder::add_partition_field(self, source_id: i32, target_name: impl alloc::string::ToString, transformation: iceberg::spec::Transform) -> iceberg::Result +pub fn iceberg::spec::UnboundPartitionSpecBuilder::add_partition_field(self, field: iceberg::spec::UnboundPartitionField) -> iceberg::Result pub fn iceberg::spec::UnboundPartitionSpecBuilder::add_partition_fields(self, fields: impl core::iter::traits::collect::IntoIterator) -> iceberg::Result pub fn iceberg::spec::UnboundPartitionSpecBuilder::build(self) -> iceberg::spec::UnboundPartitionSpec pub fn iceberg::spec::UnboundPartitionSpecBuilder::new() -> Self diff --git a/crates/iceberg/src/arrow/record_batch_partition_splitter.rs b/crates/iceberg/src/arrow/record_batch_partition_splitter.rs index 6b3aa73d22..0b679040c6 100644 --- a/crates/iceberg/src/arrow/record_batch_partition_splitter.rs +++ b/crates/iceberg/src/arrow/record_batch_partition_splitter.rs @@ -244,12 +244,14 @@ mod tests { let partition_spec = Arc::new( PartitionSpecBuilder::new(schema.clone()) .with_spec_id(1) - .add_unbound_field(UnboundPartitionField { - source_id: 1, - field_id: None, - name: "id_bucket".to_string(), - transform: Transform::Identity, - }) + .add_unbound_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .name("id_bucket".to_string()) + .transform(Transform::Identity) + .build() + .unwrap(), + ) .unwrap() .build() .unwrap(), @@ -358,12 +360,14 @@ mod tests { let partition_spec = Arc::new( PartitionSpecBuilder::new(schema.clone()) .with_spec_id(1) - .add_unbound_field(UnboundPartitionField { - source_id: 1, - field_id: None, - name: "id_bucket".to_string(), - transform: Transform::Identity, - }) + .add_unbound_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .name("id_bucket".to_string()) + .transform(Transform::Identity) + .build() + .unwrap(), + ) .unwrap() .build() .unwrap(), diff --git a/crates/iceberg/src/catalog/mod.rs b/crates/iceberg/src/catalog/mod.rs index cf7a94d20f..88681040cd 100644 --- a/crates/iceberg/src/catalog/mod.rs +++ b/crates/iceberg/src/catalog/mod.rs @@ -1103,8 +1103,8 @@ mod tests { PartitionStatisticsFile, PrimitiveType, Schema, Snapshot, SnapshotReference, SnapshotRetention, SortDirection, SortField, SortOrder, SqlViewRepresentation, StatisticsFile, Summary, TableMetadata, TableMetadataBuilder, Transform, Type, - UnboundPartitionSpec, ViewFormatVersion, ViewRepresentation, ViewRepresentations, - ViewVersion, + UnboundPartitionField, UnboundPartitionSpec, ViewFormatVersion, ViewRepresentation, + ViewRepresentations, ViewVersion, }; use crate::table::Table; use crate::test_utils::test_runtime; @@ -1659,11 +1659,32 @@ mod tests { "#, TableUpdate::AddSpec { spec: UnboundPartitionSpec::builder() - .add_partition_field(4, "ts_day".to_string(), Transform::Day) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![4]) + .name("ts_day") + .transform(Transform::Day) + .build() + .unwrap(), + ) .unwrap() - .add_partition_field(1, "id_bucket".to_string(), Transform::Bucket(16)) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .name("id_bucket") + .transform(Transform::Bucket(16)) + .build() + .unwrap(), + ) .unwrap() - .add_partition_field(2, "id_truncate".to_string(), Transform::Truncate(4)) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![2]) + .name("id_truncate") + .transform(Transform::Truncate(4)) + .build() + .unwrap(), + ) .unwrap() .build(), }, diff --git a/crates/iceberg/src/expr/visitors/expression_evaluator.rs b/crates/iceberg/src/expr/visitors/expression_evaluator.rs index 122f195897..31ce94a99b 100644 --- a/crates/iceberg/src/expr/visitors/expression_evaluator.rs +++ b/crates/iceberg/src/expr/visitors/expression_evaluator.rs @@ -280,11 +280,12 @@ mod tests { .with_spec_id(1) .add_unbound_field( UnboundPartitionField::builder() - .source_id(1) + .source_ids(vec![1]) .name("a".to_string()) .field_id(1) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .build() diff --git a/crates/iceberg/src/expr/visitors/inclusive_metrics_evaluator.rs b/crates/iceberg/src/expr/visitors/inclusive_metrics_evaluator.rs index 06c92ab3e8..8ccd1364d9 100644 --- a/crates/iceberg/src/expr/visitors/inclusive_metrics_evaluator.rs +++ b/crates/iceberg/src/expr/visitors/inclusive_metrics_evaluator.rs @@ -1661,11 +1661,12 @@ mod test { .with_spec_id(1) .add_unbound_fields(vec![ UnboundPartitionField::builder() - .source_id(1) + .source_ids(vec![1]) .name("a".to_string()) .field_id(1) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ]) .unwrap() .build() diff --git a/crates/iceberg/src/expr/visitors/inclusive_projection.rs b/crates/iceberg/src/expr/visitors/inclusive_projection.rs index d9544e4c47..88508af844 100644 --- a/crates/iceberg/src/expr/visitors/inclusive_projection.rs +++ b/crates/iceberg/src/expr/visitors/inclusive_projection.rs @@ -304,11 +304,12 @@ mod tests { .with_spec_id(1) .add_unbound_field( UnboundPartitionField::builder() - .source_id(1) + .source_ids(vec![1]) .name("a".to_string()) .field_id(1) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .build() @@ -339,12 +340,15 @@ mod tests { let partition_spec = PartitionSpec::builder(arc_schema.clone()) .with_spec_id(1) - .add_unbound_fields(vec![UnboundPartitionField { - source_id: 2, - name: "year".to_string(), - field_id: Some(1000), - transform: Transform::Year, - }]) + .add_unbound_fields(vec![ + UnboundPartitionField::builder() + .source_ids(vec![2]) + .field_id(1000) + .name("year".to_string()) + .transform(Transform::Year) + .build() + .unwrap(), + ]) .unwrap() .build() .unwrap(); @@ -374,12 +378,15 @@ mod tests { let partition_spec = PartitionSpec::builder(arc_schema.clone()) .with_spec_id(1) - .add_unbound_fields(vec![UnboundPartitionField { - source_id: 2, - name: "month".to_string(), - field_id: Some(1000), - transform: Transform::Month, - }]) + .add_unbound_fields(vec![ + UnboundPartitionField::builder() + .source_ids(vec![2]) + .field_id(1000) + .name("month".to_string()) + .transform(Transform::Month) + .build() + .unwrap(), + ]) .unwrap() .build() .unwrap(); @@ -409,12 +416,15 @@ mod tests { let partition_spec = PartitionSpec::builder(arc_schema.clone()) .with_spec_id(1) - .add_unbound_fields(vec![UnboundPartitionField { - source_id: 2, - name: "day".to_string(), - field_id: Some(1000), - transform: Transform::Day, - }]) + .add_unbound_fields(vec![ + UnboundPartitionField::builder() + .source_ids(vec![2]) + .field_id(1000) + .name("day".to_string()) + .transform(Transform::Day) + .build() + .unwrap(), + ]) .unwrap() .build() .unwrap(); @@ -446,11 +456,12 @@ mod tests { .with_spec_id(1) .add_unbound_field( UnboundPartitionField::builder() - .source_id(3) + .source_ids(vec![3]) .name("name_truncate".to_string()) .field_id(3) .transform(Transform::Truncate(4)) - .build(), + .build() + .unwrap(), ) .unwrap() .build() @@ -486,11 +497,12 @@ mod tests { .with_spec_id(1) .add_unbound_field( UnboundPartitionField::builder() - .source_id(1) + .source_ids(vec![1]) .name("a_bucket[7]".to_string()) .field_id(1) .transform(Transform::Bucket(7)) - .build(), + .build() + .unwrap(), ) .unwrap() .build() diff --git a/crates/iceberg/src/partitioning.rs b/crates/iceberg/src/partitioning.rs index ef230ef9e5..6ac62c68c4 100644 --- a/crates/iceberg/src/partitioning.rs +++ b/crates/iceberg/src/partitioning.rs @@ -177,7 +177,9 @@ mod tests { use std::sync::Arc; use super::*; - use crate::spec::{NestedField, PrimitiveType, Transform, Type, UnboundPartitionSpec}; + use crate::spec::{ + NestedField, PrimitiveType, Transform, Type, UnboundPartitionField, UnboundPartitionSpec, + }; fn test_schema() -> Schema { Schema::builder() @@ -199,7 +201,14 @@ mod tests { let mut builder = UnboundPartitionSpec::builder().with_spec_id(spec_id); for (source_id, name, transform) in fields { builder = builder - .add_partition_field(source_id, name, transform) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![source_id]) + .name(name) + .transform(transform) + .build() + .unwrap(), + ) .unwrap(); } builder.build().bind(schema.clone()).unwrap() @@ -261,12 +270,15 @@ mod tests { // Spec 0: old name let spec_v0 = PartitionSpec::builder(Arc::new(schema.clone())) .with_spec_id(0) - .add_unbound_field(crate::spec::UnboundPartitionField { - source_id: 4, - field_id: Some(1000), - name: "cat_old".to_string(), - transform: Transform::Identity, - }) + .add_unbound_field( + UnboundPartitionField::builder() + .source_ids(vec![4]) + .field_id(1000) + .name("cat_old".to_string()) + .transform(Transform::Identity) + .build() + .unwrap(), + ) .unwrap() .build() .unwrap(); @@ -274,12 +286,15 @@ mod tests { // Spec 1: newer name, same field_id let spec_v1 = PartitionSpec::builder(Arc::new(schema.clone())) .with_spec_id(1) - .add_unbound_field(crate::spec::UnboundPartitionField { - source_id: 4, - field_id: Some(1000), - name: "cat_new".to_string(), - transform: Transform::Identity, - }) + .add_unbound_field( + UnboundPartitionField::builder() + .source_ids(vec![4]) + .field_id(1000) + .name("cat_new".to_string()) + .transform(Transform::Identity) + .build() + .unwrap(), + ) .unwrap() .build() .unwrap(); @@ -297,12 +312,15 @@ mod tests { // Spec 0 (older): category partitioned by identity let spec_v0 = PartitionSpec::builder(Arc::new(schema.clone())) .with_spec_id(0) - .add_unbound_field(crate::spec::UnboundPartitionField { - source_id: 4, - field_id: Some(1000), - name: "category".to_string(), - transform: Transform::Identity, - }) + .add_unbound_field( + UnboundPartitionField::builder() + .source_ids(vec![4]) + .field_id(1000) + .name("category".to_string()) + .transform(Transform::Identity) + .build() + .unwrap(), + ) .unwrap() .build() .unwrap(); @@ -310,12 +328,15 @@ mod tests { // Spec 1 (newer): same field_id voided (partition dropped) let spec_v1 = PartitionSpec::builder(Arc::new(schema.clone())) .with_spec_id(1) - .add_unbound_field(crate::spec::UnboundPartitionField { - source_id: 4, - field_id: Some(1000), - name: "category_v2".to_string(), - transform: Transform::Void, - }) + .add_unbound_field( + UnboundPartitionField::builder() + .source_ids(vec![4]) + .field_id(1000) + .name("category_v2".to_string()) + .transform(Transform::Void) + .build() + .unwrap(), + ) .unwrap() .build() .unwrap(); @@ -371,12 +392,15 @@ mod tests { // Spec 1: partition by category + ts_year let spec_v1 = PartitionSpec::builder(Arc::new(schema.clone())) .with_spec_id(1) - .add_unbound_field(crate::spec::UnboundPartitionField { - source_id: 4, - field_id: Some(spec_v0.fields()[0].field_id), - name: "category".to_string(), - transform: Transform::Identity, - }) + .add_unbound_field( + UnboundPartitionField::builder() + .source_ids(vec![4]) + .field_id(spec_v0.fields()[0].field_id) + .name("category".to_string()) + .transform(Transform::Identity) + .build() + .unwrap(), + ) .unwrap() .add_partition_field("ts", "ts_year", Transform::Year) .unwrap() diff --git a/crates/iceberg/src/scan/cache.rs b/crates/iceberg/src/scan/cache.rs index b5eb3c4ce0..b87786e6cf 100644 --- a/crates/iceberg/src/scan/cache.rs +++ b/crates/iceberg/src/scan/cache.rs @@ -270,7 +270,7 @@ mod tests { use crate::expr::{Bind, BoundPredicate, Reference}; use crate::spec::{ Datum, FormatVersion, NestedField, PrimitiveType, Schema, SortOrder, TableMetadataBuilder, - Transform, Type, UnboundPartitionSpec, + Transform, Type, UnboundPartitionField, UnboundPartitionSpec, }; /// A historical spec whose source column was dropped from the current schema resolves to an @@ -288,7 +288,14 @@ mod tests { // Spec 0 partitions on `part` (id 2). let spec = UnboundPartitionSpec::builder() .with_spec_id(0) - .add_partition_field(2, "part", Transform::Identity) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![2]) + .name("part") + .transform(Transform::Identity) + .build() + .unwrap(), + ) .unwrap() .build(); diff --git a/crates/iceberg/src/spec/partition.rs b/crates/iceberg/src/spec/partition.rs index 3999c0d010..17aa3de968 100644 --- a/crates/iceberg/src/spec/partition.rs +++ b/crates/iceberg/src/spec/partition.rs @@ -264,20 +264,157 @@ impl PartitionKey { /// Reference to [`UnboundPartitionSpec`]. pub type UnboundPartitionSpecRef = Arc; /// Unbound partition field can be built without a schema and later bound to a schema. +/// +/// The fields are private so that an instance is always well formed: in particular +/// `source_ids` holds at least one id, which [`UnboundPartitionField::builder`] checks when the +/// field is built. #[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone, TypedBuilder)] -#[serde(rename_all = "kebab-case")] +#[serde( + try_from = "self::_serde::UnboundPartitionFieldSerde", + into = "self::_serde::UnboundPartitionFieldSerde" +)] +#[builder(build_method(into = Result))] pub struct UnboundPartitionField { - /// A source column id from the table’s schema - pub source_id: i32, + /// The source column ids from the table’s schema. A single-argument transform reads one + /// id; a v3 multi-argument transform reads several. + source_ids: Vec, /// A partition field id that is used to identify a partition field and is unique within a partition spec. /// In v2 table metadata, it is unique across all partition specs. #[builder(default, setter(strip_option(fallback = field_id_opt)))] - #[serde(skip_serializing_if = "Option::is_none")] - pub field_id: Option, + field_id: Option, /// A partition name. - pub name: String, + #[builder(setter(into))] + name: String, /// A transform that is applied to the source column to produce a partition value. - pub transform: Transform, + transform: Transform, +} + +impl UnboundPartitionField { + /// The single source column id this field reads. + /// + /// Returns an error for a multi-argument field, which reads several columns and therefore + /// has no single source id. Use [`Self::source_ids`] to handle both shapes. + pub fn source_id(&self) -> Result { + match self.source_ids.as_slice() { + [source_id] => Ok(*source_id), + source_ids => Err(invalid_data!( + "Partition field '{}' reads {} source columns and has no single source id", + self.name, + source_ids.len() + )), + } + } + + /// The source column ids this field reads, in order. Never empty. + pub fn source_ids(&self) -> &[i32] { + &self.source_ids + } + + /// The partition field id, when one was assigned. + pub fn field_id(&self) -> Option { + self.field_id + } + + /// The partition name. + pub fn name(&self) -> &str { + &self.name + } + + /// The transform applied to the source columns to produce a partition value. + pub fn transform(&self) -> Transform { + self.transform + } + + /// Return this field with the given partition field id assigned. + pub fn with_field_id(self, field_id: i32) -> Self { + Self { + field_id: Some(field_id), + ..self + } + } + + fn validate(&self) -> Result<()> { + if self.source_ids.is_empty() { + return Err(invalid_data!("Empty source-ids is not allowed")); + } + Ok(()) + } +} + +impl From for Result { + fn from(field: UnboundPartitionField) -> Self { + field.validate()?; + Ok(field) + } +} + +mod _serde { + use serde::{Deserialize, Serialize}; + + use super::UnboundPartitionField; + use crate::Error; + use crate::error::invalid_data; + use crate::spec::Transform; + + /// Per the spec a single-argument field carries `source-id` and a multi-argument field + /// carries `source-ids`. Either spelling is read, but not both at once; the one that + /// matches the field is written. + #[derive(Serialize, Deserialize)] + #[serde(rename_all = "kebab-case")] + pub(super) struct UnboundPartitionFieldSerde { + #[serde(default, skip_serializing_if = "Option::is_none")] + source_id: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + source_ids: Option>, + #[serde(default, skip_serializing_if = "Option::is_none")] + field_id: Option, + name: String, + transform: Transform, + } + + impl TryFrom for UnboundPartitionField { + type Error = Error; + + fn try_from(value: UnboundPartitionFieldSerde) -> Result { + let source_ids = match (value.source_id, value.source_ids) { + (Some(source_id), None) => vec![source_id], + (None, Some(source_ids)) => source_ids, + (Some(_), Some(_)) => { + return Err(invalid_data!( + "source-id and source-ids are mutually exclusive" + )); + } + (None, None) => { + return Err(invalid_data!( + "Either `source-id` or `source-ids` must be present" + )); + } + }; + + UnboundPartitionField::builder() + .source_ids(source_ids) + .field_id_opt(value.field_id) + .name(value.name) + .transform(value.transform) + .build() + } + } + + impl From for UnboundPartitionFieldSerde { + fn from(value: UnboundPartitionField) -> Self { + let (source_id, source_ids) = match value.source_ids.as_slice() { + [source_id] => (Some(*source_id), None), + _ => (None, Some(value.source_ids)), + }; + Self { + source_id, + source_ids, + field_id: value.field_id, + name: value.name, + transform: value.transform, + } + } + } } /// Unbound partition spec can be built without a schema and later bound to a schema. @@ -341,7 +478,7 @@ fn has_sequential_ids(field_ids: impl Iterator) -> bool { impl From for UnboundPartitionField { fn from(field: PartitionField) -> Self { UnboundPartitionField { - source_id: field.source_id, + source_ids: vec![field.source_id], field_id: Some(field.field_id), name: field.name, transform: field.transform, @@ -381,19 +518,14 @@ impl UnboundPartitionSpecBuilder { } /// Add a new partition field to the partition spec from an unbound partition field. - pub fn add_partition_field( - self, - source_id: i32, - target_name: impl ToString, - transformation: Transform, - ) -> Result { - let field = UnboundPartitionField { - source_id, - field_id: None, - name: target_name.to_string(), - transform: transformation, - }; - self.add_partition_field_internal(field) + pub fn add_partition_field(mut self, field: UnboundPartitionField) -> Result { + self.check_name_set_and_unique(&field.name)?; + self.check_for_redundant_partitions(&field.source_ids, &field.transform)?; + if let Some(partition_field_id) = field.field_id { + self.check_partition_id_unique(partition_field_id)?; + } + self.fields.push(field); + Ok(self) } /// Add multiple partition fields to the partition spec. @@ -403,21 +535,11 @@ impl UnboundPartitionSpecBuilder { ) -> Result { let mut builder = self; for field in fields { - builder = builder.add_partition_field_internal(field)?; + builder = builder.add_partition_field(field)?; } Ok(builder) } - fn add_partition_field_internal(mut self, field: UnboundPartitionField) -> Result { - self.check_name_set_and_unique(&field.name)?; - self.check_for_redundant_partitions(field.source_id, &field.transform)?; - if let Some(partition_field_id) = field.field_id { - self.check_partition_id_unique(partition_field_id)?; - } - self.fields.push(field); - Ok(self) - } - /// Build the unbound partition spec. pub fn build(self) -> UnboundPartitionSpec { UnboundPartitionSpec { @@ -495,7 +617,7 @@ impl PartitionSpecBuilder { })? .id; let field = UnboundPartitionField { - source_id, + source_ids: vec![source_id], field_id: None, name: target_name.into(), transform, @@ -510,7 +632,7 @@ impl PartitionSpecBuilder { /// Otherwise, a new `field_id` is assigned. pub fn add_unbound_field(mut self, field: UnboundPartitionField) -> Result { self.check_name_set_and_unique(&field.name)?; - self.check_for_redundant_partitions(field.source_id, &field.transform)?; + self.check_for_redundant_partitions(&field.source_ids, &field.transform)?; Self::check_name_does_not_collide_with_schema(&field, &self.schema)?; Self::check_transform_compatibility(&field, &self.schema)?; if let Some(partition_field_id) = field.field_id { @@ -574,7 +696,7 @@ impl PartitionSpecBuilder { }; bound_fields.push(PartitionField { - source_id: field.source_id, + source_id: field.source_id()?, field_id: partition_field_id, name: field.name, transform: field.transform, @@ -595,14 +717,14 @@ impl PartitionSpecBuilder { match schema.field_by_name(field.name.as_str()) { Some(schema_collision) => { if field.transform == Transform::Identity { - if schema_collision.id == field.source_id { + if schema_collision.id == field.source_id()? { Ok(()) } else { Err(invalid_data!( "Cannot create identity partition sourced from different field in schema. Field name '{}' has id `{}` in schema but partition source id is `{}`", field.name, schema_collision.id, - field.source_id + field.source_id()? )) } } else { @@ -619,11 +741,11 @@ impl PartitionSpecBuilder { /// Ensure that the transformation of the field is compatible with type of the field /// in the schema. Implicitly also checks if the source field exists in the schema. fn check_transform_compatibility(field: &UnboundPartitionField, schema: &Schema) -> Result<()> { - let schema_field = schema.field_by_id(field.source_id).ok_or_else(|| { - invalid_data!( - "Cannot find partition source field with id `{}` in schema", - field.source_id - ) + // A multi-argument field is rejected here on purpose: a bound `PartitionField` has + // no place for more than one source id yet. + let source_id = field.source_id()?; + let schema_field = schema.field_by_id(source_id).ok_or_else(|| { + invalid_data!("Cannot find partition source field with id `{source_id}` in schema") })?; if field.transform != Transform::Void { @@ -668,15 +790,19 @@ trait CorePartitionSpecValidator { } /// For a single source-column transformations must be unique. - fn check_for_redundant_partitions(&self, source_id: i32, transform: &Transform) -> Result<()> { + fn check_for_redundant_partitions( + &self, + source_ids: &[i32], + transform: &Transform, + ) -> Result<()> { let collision = self.fields().iter().find(|f| { - f.source_id == source_id && f.transform.dedup_name() == transform.dedup_name() + f.source_ids == source_ids && f.transform.dedup_name() == transform.dedup_name() }); if let Some(collision) = collision { Err(invalid_data!( - "Cannot add redundant partition with source id `{}` and transform `{}`. A partition with the same source id and transform already exists with name `{}`", - source_id, + "Cannot add redundant partition with source ids `{:?}` and transform `{}`. A partition with the same source ids and transform already exists with name `{}`", + source_ids, transform.dedup_name(), collision.name )) @@ -779,15 +905,17 @@ mod tests { let partition_spec = PartitionSpec::builder(schema.clone()) .add_unbound_fields(vec![ UnboundPartitionField::builder() - .source_id(1) + .source_ids(vec![1]) .name("id".to_string()) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), UnboundPartitionField::builder() - .source_id(2) + .source_ids(vec![2]) .name("name_string".to_string()) .transform(Transform::Void) - .build(), + .build() + .unwrap(), ]) .unwrap() .with_spec_id(1) @@ -802,15 +930,17 @@ mod tests { .with_spec_id(1) .add_unbound_fields(vec![ UnboundPartitionField::builder() - .source_id(1) + .source_ids(vec![1]) .name("id_void".to_string()) .transform(Transform::Void) - .build(), + .build() + .unwrap(), UnboundPartitionField::builder() - .source_id(2) + .source_ids(vec![2]) .name("name_void".to_string()) .transform(Transform::Void) - .build(), + .build() + .unwrap(), ]) .unwrap() .build() @@ -848,17 +978,17 @@ mod tests { let partition_spec: UnboundPartitionSpec = serde_json::from_str(spec).unwrap(); assert_eq!(Some(1), partition_spec.spec_id); - assert_eq!(4, partition_spec.fields[0].source_id); + assert_eq!(4, partition_spec.fields[0].source_id().unwrap()); assert_eq!(Some(1000), partition_spec.fields[0].field_id); assert_eq!("ts_day", partition_spec.fields[0].name); assert_eq!(Transform::Day, partition_spec.fields[0].transform); - assert_eq!(1, partition_spec.fields[1].source_id); + assert_eq!(1, partition_spec.fields[1].source_id().unwrap()); assert_eq!(Some(1001), partition_spec.fields[1].field_id); assert_eq!("id_bucket", partition_spec.fields[1].name); assert_eq!(Transform::Bucket(16), partition_spec.fields[1].transform); - assert_eq!(2, partition_spec.fields[2].source_id); + assert_eq!(2, partition_spec.fields[2].source_id().unwrap()); assert_eq!(Some(1002), partition_spec.fields[2].field_id); assert_eq!("id_truncate", partition_spec.fields[2].name); assert_eq!(Transform::Truncate(4), partition_spec.fields[2].transform); @@ -875,7 +1005,7 @@ mod tests { let partition_spec: UnboundPartitionSpec = serde_json::from_str(spec).unwrap(); assert_eq!(None, partition_spec.spec_id); - assert_eq!(4, partition_spec.fields[0].source_id); + assert_eq!(4, partition_spec.fields[0].source_id().unwrap()); assert_eq!(None, partition_spec.fields[0].field_id); assert_eq!("ts_day", partition_spec.fields[0].name); assert_eq!(Transform::Day, partition_spec.fields[0].transform); @@ -884,7 +1014,14 @@ mod tests { #[test] fn test_unbound_partition_spec_serialization_skips_none_fields() { let spec = UnboundPartitionSpec::builder() - .add_partition_field(4, "ts_day".to_string(), Transform::Day) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![4]) + .name("ts_day") + .transform(Transform::Day) + .build() + .unwrap(), + ) .unwrap() .build(); @@ -1131,9 +1268,23 @@ mod tests { #[test] fn test_builder_disallow_duplicate_names() { UnboundPartitionSpec::builder() - .add_partition_field(1, "ts_day".to_string(), Transform::Day) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .name("ts_day") + .transform(Transform::Day) + .build() + .unwrap(), + ) .unwrap() - .add_partition_field(2, "ts_day".to_string(), Transform::Day) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![2]) + .name("ts_day") + .transform(Transform::Day) + .build() + .unwrap(), + ) .unwrap_err(); } @@ -1148,14 +1299,14 @@ mod tests { .unwrap(); PartitionSpec::builder(schema.clone()) .add_unbound_field(UnboundPartitionField { - source_id: 1, + source_ids: vec![1], field_id: Some(1000), name: "id".to_string(), transform: Transform::Identity, }) .unwrap() .add_unbound_field(UnboundPartitionField { - source_id: 2, + source_ids: vec![2], field_id: Some(1000), name: "id_bucket".to_string(), transform: Transform::Bucket(16), @@ -1176,14 +1327,14 @@ mod tests { let spec = PartitionSpec::builder(schema.clone()) .with_spec_id(1) .add_unbound_field(UnboundPartitionField { - source_id: 1, + source_ids: vec![1], name: "id".to_string(), transform: Transform::Identity, field_id: Some(1012), }) .unwrap() .add_unbound_field(UnboundPartitionField { - source_id: 2, + source_ids: vec![2], name: "name_void".to_string(), transform: Transform::Void, field_id: None, @@ -1191,7 +1342,7 @@ mod tests { .unwrap() // Should keep its ID even if its lower .add_unbound_field(UnboundPartitionField { - source_id: 3, + source_ids: vec![3], name: "year".to_string(), transform: Transform::Year, field_id: Some(1), @@ -1262,7 +1413,7 @@ mod tests { let err = PartitionSpec::builder(schema) .with_spec_id(1) .add_unbound_field(UnboundPartitionField { - source_id: 1, + source_ids: vec![1], field_id: None, name: "id".to_string(), transform: Transform::Bucket(16), @@ -1289,7 +1440,7 @@ mod tests { PartitionSpec::builder(schema.clone()) .with_spec_id(1) .add_unbound_field(UnboundPartitionField { - source_id: 1, + source_ids: vec![1], field_id: None, name: "id".to_string(), transform: Transform::Identity, @@ -1302,7 +1453,7 @@ mod tests { PartitionSpec::builder(schema) .with_spec_id(1) .add_unbound_field(UnboundPartitionField { - source_id: 2, + source_ids: vec![2], field_id: None, name: "id".to_string(), transform: Transform::Identity, @@ -1326,13 +1477,13 @@ mod tests { .with_spec_id(1) .add_unbound_fields(vec![ UnboundPartitionField { - source_id: 1, + source_ids: vec![1], field_id: None, name: "id_bucket".to_string(), transform: Transform::Bucket(16), }, UnboundPartitionField { - source_id: 2, + source_ids: vec![2], field_id: None, name: "name".to_string(), transform: Transform::Identity, @@ -1347,13 +1498,13 @@ mod tests { .with_spec_id(1) .add_unbound_fields(vec![ UnboundPartitionField { - source_id: 1, + source_ids: vec![1], field_id: None, name: "id_bucket".to_string(), transform: Transform::Bucket(16), }, UnboundPartitionField { - source_id: 4, + source_ids: vec![4], field_id: None, name: "name".to_string(), transform: Transform::Identity, @@ -1375,7 +1526,7 @@ mod tests { let err = PartitionSpec::builder(schema) .with_spec_id(1) .add_unbound_fields(vec![UnboundPartitionField { - source_id: 2, + source_ids: vec![2], field_id: None, name: "v_part".to_string(), transform: Transform::Identity, @@ -1392,12 +1543,22 @@ mod tests { fn test_builder_disallows_redundant() { let err = UnboundPartitionSpec::builder() .with_spec_id(1) - .add_partition_field(1, "id_bucket[16]".to_string(), Transform::Bucket(16)) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .name("id_bucket[16]") + .transform(Transform::Bucket(16)) + .build() + .unwrap(), + ) .unwrap() .add_partition_field( - 1, - "id_bucket_with_other_name".to_string(), - Transform::Bucket(16), + UnboundPartitionField::builder() + .source_ids(vec![1]) + .name("id_bucket_with_other_name") + .transform(Transform::Bucket(16)) + .build() + .unwrap(), ) .unwrap_err(); assert!(err.message().contains("redundant partition")); @@ -1415,7 +1576,7 @@ mod tests { PartitionSpec::builder(schema) .with_spec_id(1) .add_unbound_field(UnboundPartitionField { - source_id: 1, + source_ids: vec![1], field_id: None, name: "id_year".to_string(), transform: Transform::Year, @@ -1428,7 +1589,7 @@ mod tests { let spec = UnboundPartitionSpec::builder() .with_spec_id(1) .add_partition_fields(vec![UnboundPartitionField { - source_id: 1, + source_ids: vec![1], field_id: None, name: "id_bucket[16]".to_string(), transform: Transform::Bucket(16), @@ -1439,7 +1600,7 @@ mod tests { assert_eq!(spec, UnboundPartitionSpec { spec_id: Some(1), fields: vec![UnboundPartitionField { - source_id: 1, + source_ids: vec![1], field_id: None, name: "id_bucket[16]".to_string(), transform: Transform::Bucket(16), @@ -1460,7 +1621,7 @@ mod tests { let partition_spec_1 = PartitionSpec::builder(schema.clone()) .with_spec_id(1) .add_unbound_field(UnboundPartitionField { - source_id: 1, + source_ids: vec![1], field_id: None, name: "id_bucket".to_string(), transform: Transform::Bucket(16), @@ -1472,7 +1633,7 @@ mod tests { let partition_spec_2 = PartitionSpec::builder(schema) .with_spec_id(1) .add_unbound_field(UnboundPartitionField { - source_id: 1, + source_ids: vec![1], field_id: None, name: "id_bucket".to_string(), transform: Transform::Bucket(16), @@ -1496,7 +1657,7 @@ mod tests { let partition_spec_1 = PartitionSpec::builder(schema.clone()) .with_spec_id(1) .add_unbound_field(UnboundPartitionField { - source_id: 1, + source_ids: vec![1], field_id: None, name: "id_bucket".to_string(), transform: Transform::Bucket(16), @@ -1508,7 +1669,7 @@ mod tests { let partition_spec_2 = PartitionSpec::builder(schema) .with_spec_id(1) .add_unbound_field(UnboundPartitionField { - source_id: 1, + source_ids: vec![1], field_id: None, name: "id_bucket".to_string(), transform: Transform::Bucket(32), @@ -1533,7 +1694,7 @@ mod tests { let partition_spec_1 = PartitionSpec::builder(schema.clone()) .with_spec_id(1) .add_unbound_field(UnboundPartitionField { - source_id: 1, + source_ids: vec![1], field_id: None, name: "id_bucket".to_string(), transform: Transform::Bucket(16), @@ -1545,7 +1706,7 @@ mod tests { let partition_spec_2 = PartitionSpec::builder(schema) .with_spec_id(1) .add_unbound_field(UnboundPartitionField { - source_id: 2, + source_ids: vec![2], field_id: None, name: "id_bucket".to_string(), transform: Transform::Bucket(16), @@ -1570,14 +1731,14 @@ mod tests { let partition_spec_1 = PartitionSpec::builder(schema.clone()) .with_spec_id(1) .add_unbound_field(UnboundPartitionField { - source_id: 1, + source_ids: vec![1], field_id: None, name: "id_bucket".to_string(), transform: Transform::Bucket(16), }) .unwrap() .add_unbound_field(UnboundPartitionField { - source_id: 2, + source_ids: vec![2], field_id: None, name: "name".to_string(), transform: Transform::Identity, @@ -1589,14 +1750,14 @@ mod tests { let partition_spec_2 = PartitionSpec::builder(schema) .with_spec_id(1) .add_unbound_field(UnboundPartitionField { - source_id: 2, + source_ids: vec![2], field_id: None, name: "name".to_string(), transform: Transform::Identity, }) .unwrap() .add_unbound_field(UnboundPartitionField { - source_id: 1, + source_ids: vec![1], field_id: None, name: "id_bucket".to_string(), transform: Transform::Bucket(16), @@ -1631,14 +1792,14 @@ mod tests { let spec = PartitionSpec::builder(schema) .with_spec_id(1) .add_unbound_field(UnboundPartitionField { - source_id: 1, + source_ids: vec![1], field_id: Some(1001), name: "id".to_string(), transform: Transform::Identity, }) .unwrap() .add_unbound_field(UnboundPartitionField { - source_id: 2, + source_ids: vec![2], field_id: Some(1000), name: "name".to_string(), transform: Transform::Identity, @@ -1663,14 +1824,14 @@ mod tests { let spec = PartitionSpec::builder(schema) .with_spec_id(1) .add_unbound_field(UnboundPartitionField { - source_id: 1, + source_ids: vec![1], field_id: Some(1000), name: "id".to_string(), transform: Transform::Identity, }) .unwrap() .add_unbound_field(UnboundPartitionField { - source_id: 2, + source_ids: vec![2], field_id: Some(1001), name: "name".to_string(), transform: Transform::Identity, @@ -1697,14 +1858,14 @@ mod tests { let spec = PartitionSpec::builder(schema) .with_spec_id(1) .add_unbound_field(UnboundPartitionField { - source_id: 1, + source_ids: vec![1], field_id: Some(999), name: "id".to_string(), transform: Transform::Identity, }) .unwrap() .add_unbound_field(UnboundPartitionField { - source_id: 2, + source_ids: vec![2], field_id: Some(1000), name: "name".to_string(), transform: Transform::Identity, @@ -1731,14 +1892,14 @@ mod tests { let spec = PartitionSpec::builder(schema) .with_spec_id(1) .add_unbound_field(UnboundPartitionField { - source_id: 1, + source_ids: vec![1], field_id: Some(1000), name: "id".to_string(), transform: Transform::Identity, }) .unwrap() .add_unbound_field(UnboundPartitionField { - source_id: 2, + source_ids: vec![2], field_id: Some(1002), name: "name".to_string(), transform: Transform::Identity, @@ -1847,4 +2008,106 @@ mod tests { "data=a%2Fb%2Fc%2Fd/data_truc_10=a%2Fb%2Fc%2Fd" ); } + + #[test] + fn test_unbound_partition_field_reads_source_ids_only() { + let field: UnboundPartitionField = serde_json::from_str( + r#"{"source-ids": [1, 2], "name": "m", "transform": "bucket[4]"}"#, + ) + .unwrap(); + + assert_eq!([1, 2], field.source_ids()); + // a multi-argument field has no single source id + assert!( + field + .source_id() + .unwrap_err() + .to_string() + .contains("has no single source id") + ); + + let serialized = serde_json::to_value(&field).unwrap(); + assert_eq!( + Some(&serde_json::json!([1, 2])), + serialized.get("source-ids") + ); + assert!(serialized.get("source-id").is_none()); + } + + #[test] + fn test_unbound_partition_field_single_source_id_round_trip() { + let field: UnboundPartitionField = + serde_json::from_str(r#"{"source-id": 1, "name": "m", "transform": "identity"}"#) + .unwrap(); + + assert_eq!(1, field.source_id().unwrap()); + assert_eq!([1], field.source_ids()); + + // a single-argument field keeps writing source-id + let serialized = serde_json::to_value(&field).unwrap(); + assert_eq!(Some(&serde_json::json!(1)), serialized.get("source-id")); + assert!(serialized.get("source-ids").is_none()); + } + + #[test] + fn test_unbound_partition_field_rejects_malformed_source_ids() { + for (input, expected) in [ + ( + r#"{"source-ids": [], "name": "m", "transform": "identity"}"#, + "Empty source-ids is not allowed", + ), + ( + r#"{"name": "m", "transform": "identity"}"#, + "Either `source-id` or `source-ids` must be present", + ), + ( + r#"{"source-id": 1, "source-ids": [1], "name": "m", "transform": "identity"}"#, + "mutually exclusive", + ), + ] { + let err = serde_json::from_str::(input).unwrap_err(); + assert!( + err.to_string().contains(expected), + "unexpected error for {input}: {err}" + ); + } + } + + #[test] + fn test_binding_a_multi_argument_field_fails_loudly() { + let schema = Schema::builder() + .with_fields(vec![ + NestedField::required(1, "a", Type::Primitive(PrimitiveType::Int)).into(), + NestedField::required(2, "b", Type::Primitive(PrimitiveType::Int)).into(), + ]) + .build() + .unwrap(); + + let spec: UnboundPartitionSpec = serde_json::from_str( + r#"{"spec-id": 1, "fields": [{"source-ids": [1, 2], "field-id": 1000, "name": "m", "transform": "bucket[4]"}]}"#, + ) + .unwrap(); + + // A bound PartitionField still carries a single source id, so binding has to refuse + // rather than quietly keep the first one. + let err = spec.bind(schema).unwrap_err(); + assert!( + err.to_string().contains("has no single source id"), + "unexpected error: {err}" + ); + } + + #[test] + fn test_unbound_partition_field_builder_rejects_empty_source_ids() { + let err = UnboundPartitionField::builder() + .source_ids(vec![]) + .name("m") + .transform(Transform::Identity) + .build() + .unwrap_err(); + assert!( + err.to_string().contains("Empty source-ids is not allowed"), + "unexpected error: {err}" + ); + } } diff --git a/crates/iceberg/src/spec/snapshot_summary.rs b/crates/iceberg/src/spec/snapshot_summary.rs index a48a5abdda..f8e69c4c86 100644 --- a/crates/iceberg/src/spec/snapshot_summary.rs +++ b/crates/iceberg/src/spec/snapshot_summary.rs @@ -828,10 +828,11 @@ mod tests { PartitionSpec::builder(schema.clone()) .add_unbound_fields(vec![ UnboundPartitionField::builder() - .source_id(2) + .source_ids(vec![2]) .name("year".to_string()) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ]) .unwrap() .with_spec_id(1) @@ -978,10 +979,11 @@ mod tests { PartitionSpec::builder(schema.clone()) .add_unbound_fields(vec![ UnboundPartitionField::builder() - .source_id(2) + .source_ids(vec![2]) .name("year".to_string()) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ]) .unwrap() .with_spec_id(1) diff --git a/crates/iceberg/src/spec/table_metadata.rs b/crates/iceberg/src/spec/table_metadata.rs index 24f10355bf..83468602cd 100644 --- a/crates/iceberg/src/spec/table_metadata.rs +++ b/crates/iceberg/src/spec/table_metadata.rs @@ -1721,12 +1721,15 @@ mod tests { let partition_spec = PartitionSpec::builder(schema.clone()) .with_spec_id(0) - .add_unbound_field(UnboundPartitionField { - name: "ts_day".to_string(), - transform: Transform::Day, - source_id: 4, - field_id: Some(1000), - }) + .add_unbound_field( + UnboundPartitionField::builder() + .source_ids(vec![4]) + .field_id(1000) + .name("ts_day".to_string()) + .transform(Transform::Day) + .build() + .unwrap(), + ) .unwrap() .build() .unwrap(); @@ -1869,12 +1872,15 @@ mod tests { let partition_spec = PartitionSpec::builder(schema.clone()) .with_spec_id(0) - .add_unbound_field(UnboundPartitionField { - name: "ts_day".to_string(), - transform: Transform::Day, - source_id: 4, - field_id: Some(1000), - }) + .add_unbound_field( + UnboundPartitionField::builder() + .source_ids(vec![4]) + .field_id(1000) + .name("ts_day".to_string()) + .transform(Transform::Day) + .build() + .unwrap(), + ) .unwrap() .build() .unwrap(); @@ -2945,12 +2951,15 @@ mod tests { let partition_spec = PartitionSpec::builder(schema.clone()) .with_spec_id(0) - .add_unbound_field(UnboundPartitionField { - name: "x".to_string(), - transform: Transform::Identity, - source_id: 1, - field_id: Some(1000), - }) + .add_unbound_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .field_id(1000) + .name("x".to_string()) + .transform(Transform::Identity) + .build() + .unwrap(), + ) .unwrap() .build() .unwrap(); @@ -3042,12 +3051,15 @@ mod tests { let partition_spec = PartitionSpec::builder(schema2.clone()) .with_spec_id(0) - .add_unbound_field(UnboundPartitionField { - name: "x".to_string(), - transform: Transform::Identity, - source_id: 1, - field_id: Some(1000), - }) + .add_unbound_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .field_id(1000) + .name("x".to_string()) + .transform(Transform::Identity) + .build() + .unwrap(), + ) .unwrap() .build() .unwrap(); @@ -3171,12 +3183,15 @@ mod tests { let partition_spec = PartitionSpec::builder(schema.clone()) .with_spec_id(0) - .add_unbound_field(UnboundPartitionField { - name: "x".to_string(), - transform: Transform::Identity, - source_id: 1, - field_id: Some(1000), - }) + .add_unbound_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .field_id(1000) + .name("x".to_string()) + .transform(Transform::Identity) + .build() + .unwrap(), + ) .unwrap() .build() .unwrap(); @@ -3257,12 +3272,15 @@ mod tests { let partition_spec = PartitionSpec::builder(schema.clone()) .with_spec_id(0) - .add_unbound_field(UnboundPartitionField { - name: "x".to_string(), - transform: Transform::Identity, - source_id: 1, - field_id: Some(1000), - }) + .add_unbound_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .field_id(1000) + .name("x".to_string()) + .transform(Transform::Identity) + .build() + .unwrap(), + ) .unwrap() .build() .unwrap(); @@ -4481,7 +4499,14 @@ mod tests { schema.clone(), UnboundPartitionSpec::builder() .with_spec_id(0) - .add_partition_field(2, "y", Transform::Identity) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![2]) + .name("y") + .transform(Transform::Identity) + .build() + .unwrap(), + ) .unwrap() .build(), SortOrder::unsorted_order(), @@ -4496,7 +4521,14 @@ mod tests { .into_builder(None) .add_partition_spec( UnboundPartitionSpec::builder() - .add_partition_field(3, "z", Transform::Identity) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![3]) + .name("z") + .transform(Transform::Identity) + .build() + .unwrap(), + ) .unwrap() .build(), ) diff --git a/crates/iceberg/src/spec/table_metadata_builder.rs b/crates/iceberg/src/spec/table_metadata_builder.rs index 3ed848d2d6..9a8206cb27 100644 --- a/crates/iceberg/src/spec/table_metadata_builder.rs +++ b/crates/iceberg/src/spec/table_metadata_builder.rs @@ -741,7 +741,7 @@ impl TableMetadataBuilder { for partition_field in unbound_spec.fields() { let exists_in_any_schema = self .metadata - .name_exists_in_any_schema(&partition_field.name); + .name_exists_in_any_schema(partition_field.name()); // Skip if partition field name doesn't conflict with any schema field if !exists_in_any_schema { @@ -749,15 +749,15 @@ impl TableMetadataBuilder { } // If name exists in schemas, validate against current schema rules - if let Some(schema_field) = current_schema.field_by_name(&partition_field.name) { + if let Some(schema_field) = current_schema.field_by_name(partition_field.name()) { let is_identity_transform = - partition_field.transform == crate::spec::Transform::Identity; - let has_matching_source_id = schema_field.id == partition_field.source_id; + partition_field.transform() == crate::spec::Transform::Identity; + let has_matching_source_id = schema_field.id == partition_field.source_id()?; if !is_identity_transform { return Err(invalid_data!( "Cannot create partition with name '{}' that conflicts with schema field and is not an identity transform.", - partition_field.name + partition_field.name() )); } @@ -765,9 +765,9 @@ impl TableMetadataBuilder { return Err(invalid_data!( "Cannot create identity partition sourced from different field in schema. \ Field name '{}' has id `{}` in schema but partition source id is `{}`", - partition_field.name, + partition_field.name(), schema_field.id, - partition_field.source_id + partition_field.source_id()? )); } } @@ -850,21 +850,22 @@ impl TableMetadataBuilder { .partition_specs .values() .flat_map(|spec| spec.fields()) - .map(|field| ((field.source_id, &field.transform), field.field_id)) + .map(|field| ((vec![field.source_id], field.transform), field.field_id)) .collect(); // Create new fields with reused field IDs where possible let fields = unbound_spec .fields .into_iter() - .map(|mut field| { - if field.field_id.is_none() + .map(|field| { + if field.field_id().is_none() && let Some(&existing_field_id) = - equivalent_field_ids.get(&(field.source_id, &field.transform)) + equivalent_field_ids.get(&(field.source_ids().to_vec(), field.transform())) { - field.field_id = Some(existing_field_id); + field.with_field_id(existing_field_id) + } else { + field } - field }) .collect(); @@ -1228,15 +1229,19 @@ impl TableMetadataBuilder { // Re-build partition spec with new ids let mut fresh_spec = PartitionSpecBuilder::new(fresh_schema.clone()); for field in spec.fields() { - let source_field_name = previous_id_to_name.get(&field.source_id).ok_or_else(|| { + let source_id = field.source_id()?; + let source_field_name = previous_id_to_name.get(&source_id).ok_or_else(|| { invalid_data!( "Cannot find source column with id {} for partition column {} in schema.", - field.source_id, - field.name + source_id, + field.name() ) })?; - fresh_spec = - fresh_spec.add_partition_field(source_field_name, &field.name, field.transform)?; + fresh_spec = fresh_spec.add_partition_field( + source_field_name, + field.name(), + field.transform(), + )?; } let fresh_spec = fresh_spec.build()?; @@ -1432,7 +1437,14 @@ mod tests { fn partition_spec() -> UnboundPartitionSpec { UnboundPartitionSpec::builder() .with_spec_id(0) - .add_partition_field(2, "y", Transform::Identity) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![2]) + .name("y") + .transform(Transform::Identity) + .build() + .unwrap(), + ) .unwrap() .build() } @@ -1666,12 +1678,15 @@ mod tests { // partition_spec() has None set for field-id spec: PartitionSpec::builder(schema()) .with_spec_id(0) - .add_unbound_field(UnboundPartitionField { - name: "y".to_string(), - transform: Transform::Identity, - source_id: 2, - field_id: Some(1000) - }) + .add_unbound_field( + UnboundPartitionField::builder() + .source_ids(vec![2]) + .field_id(1000) + .name("y".to_string()) + .transform(Transform::Identity) + .build() + .unwrap() + ) .unwrap() .build() .unwrap() @@ -1739,20 +1754,21 @@ mod tests { let added_spec = UnboundPartitionSpec::builder() .with_spec_id(10) .add_partition_fields(vec![ - UnboundPartitionField { - // The previous field - has field_id set - name: "y".to_string(), - transform: Transform::Identity, - source_id: 2, - field_id: Some(1000), - }, - UnboundPartitionField { - // A new field without field id - should still be without field id in changes - name: "z".to_string(), - transform: Transform::Identity, - source_id: 3, - field_id: None, - }, + // The previous field - has field_id set + UnboundPartitionField::builder() + .source_ids(vec![2]) + .field_id(1000) + .name("y".to_string()) + .transform(Transform::Identity) + .build() + .unwrap(), + // A new field without field id - should still be without field id in changes + UnboundPartitionField::builder() + .source_ids(vec![3]) + .name("z".to_string()) + .transform(Transform::Identity) + .build() + .unwrap(), ]) .unwrap() .build(); @@ -1767,19 +1783,25 @@ mod tests { let expected_change = added_spec.with_spec_id(1); let expected_spec = PartitionSpec::builder(schema()) .with_spec_id(1) - .add_unbound_field(UnboundPartitionField { - name: "y".to_string(), - transform: Transform::Identity, - source_id: 2, - field_id: Some(1000), - }) + .add_unbound_field( + UnboundPartitionField::builder() + .source_ids(vec![2]) + .field_id(1000) + .name("y".to_string()) + .transform(Transform::Identity) + .build() + .unwrap(), + ) .unwrap() - .add_unbound_field(UnboundPartitionField { - name: "z".to_string(), - transform: Transform::Identity, - source_id: 3, - field_id: Some(1001), - }) + .add_unbound_field( + UnboundPartitionField::builder() + .source_ids(vec![3]) + .field_id(1001) + .name("z".to_string()) + .transform(Transform::Identity) + .build() + .unwrap(), + ) .unwrap() .build() .unwrap(); @@ -1817,7 +1839,14 @@ mod tests { let schema = builder.get_current_schema().unwrap().clone(); let added_spec = UnboundPartitionSpec::builder() .with_spec_id(10) - .add_partition_field(1, "y_bucket[2]", Transform::Bucket(2)) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .name("y_bucket[2]") + .transform(Transform::Bucket(2)) + .build() + .unwrap(), + ) .unwrap() .build(); @@ -1831,12 +1860,15 @@ mod tests { let expected_spec = PartitionSpec::builder(schema) .with_spec_id(1) - .add_unbound_field(UnboundPartitionField { - name: "y_bucket[2]".to_string(), - transform: Transform::Bucket(2), - source_id: 1, - field_id: Some(1001), - }) + .add_unbound_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .field_id(1001) + .name("y_bucket[2]".to_string()) + .transform(Transform::Bucket(2)) + .build() + .unwrap(), + ) .unwrap() .build() .unwrap(); @@ -2401,18 +2433,20 @@ mod tests { let added_spec = UnboundPartitionSpec::builder() .with_spec_id(10) .add_partition_fields(vec![ - UnboundPartitionField { - name: "y".to_string(), - transform: Transform::Identity, - source_id: 2, - field_id: Some(1000), - }, - UnboundPartitionField { - name: "z".to_string(), - transform: Transform::Identity, - source_id: 3, - field_id: Some(1002), - }, + UnboundPartitionField::builder() + .source_ids(vec![2]) + .field_id(1000) + .name("y".to_string()) + .transform(Transform::Identity) + .build() + .unwrap(), + UnboundPartitionField::builder() + .source_ids(vec![3]) + .field_id(1002) + .name("z".to_string()) + .transform(Transform::Identity) + .build() + .unwrap(), ]) .unwrap() .build(); @@ -2737,7 +2771,14 @@ mod tests { let partition_spec_with_bucket = UnboundPartitionSpec::builder() .with_spec_id(0) - .add_partition_field(1, "bucket_data", Transform::Bucket(16)) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .name("bucket_data") + .transform(Transform::Bucket(16)) + .build() + .unwrap(), + ) .unwrap() .build(); @@ -2796,7 +2837,14 @@ mod tests { let partition_spec = UnboundPartitionSpec::builder() .with_spec_id(0) - .add_partition_field(1, "partition_col", Transform::Bucket(16)) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .name("partition_col") + .transform(Transform::Bucket(16)) + .build() + .unwrap(), + ) .unwrap() .build(); @@ -2845,7 +2893,14 @@ mod tests { let partition_spec = UnboundPartitionSpec::builder() .with_spec_id(0) - .add_partition_field(1, "data_bucket", Transform::Bucket(16)) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .name("data_bucket") + .transform(Transform::Bucket(16)) + .build() + .unwrap(), + ) .unwrap() .build(); @@ -2868,7 +2923,14 @@ mod tests { let conflicting_partition_spec = UnboundPartitionSpec::builder() .with_spec_id(1) - .add_partition_field(1, "existing_field", Transform::Bucket(8)) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .name("existing_field") + .transform(Transform::Bucket(8)) + .build() + .unwrap(), + ) .unwrap() .build(); @@ -2898,7 +2960,14 @@ mod tests { let partition_spec = UnboundPartitionSpec::builder() .with_spec_id(0) - .add_partition_field(1, "bucket_data", Transform::Bucket(16)) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .name("bucket_data") + .transform(Transform::Bucket(16)) + .build() + .unwrap(), + ) .unwrap() .build(); @@ -2994,7 +3063,14 @@ mod tests { let partition_spec = UnboundPartitionSpec::builder() .with_spec_id(0) - .add_partition_field(2, "partition_data", Transform::Identity) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![2]) + .name("partition_data") + .transform(Transform::Identity) + .build() + .unwrap(), + ) .unwrap() .build(); @@ -3043,7 +3119,14 @@ mod tests { let partition_spec = UnboundPartitionSpec::builder() .with_spec_id(0) - .add_partition_field(1, "bucket_data", Transform::Bucket(16)) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .name("bucket_data") + .transform(Transform::Bucket(16)) + .build() + .unwrap(), + ) .unwrap() .build(); @@ -3093,7 +3176,14 @@ mod tests { let partition_spec = UnboundPartitionSpec::builder() .with_spec_id(0) - .add_partition_field(1, "data_bucket", Transform::Bucket(16)) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .name("data_bucket") + .transform(Transform::Bucket(16)) + .build() + .unwrap(), + ) .unwrap() .build(); @@ -3117,7 +3207,14 @@ mod tests { // Try to add a partition spec with a field name that does NOT conflict with existing schema fields let non_conflicting_partition_spec = UnboundPartitionSpec::builder() .with_spec_id(1) - .add_partition_field(2, "new_partition_field", Transform::Bucket(8)) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![2]) + .name("new_partition_field") + .transform(Transform::Bucket(8)) + .build() + .unwrap(), + ) .unwrap() .build(); @@ -3518,7 +3615,14 @@ mod tests { // Create initial table with spec 0: identity(id) -> field_id = 1000 let initial_spec = UnboundPartitionSpec::builder() - .add_partition_field(1, "id", Transform::Identity) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .name("id") + .transform(Transform::Identity) + .build() + .unwrap(), + ) .unwrap() .build(); @@ -3537,7 +3641,14 @@ mod tests { // Add spec 1: bucket(data) -> field_id = 1001 let spec1 = UnboundPartitionSpec::builder() - .add_partition_field(2, "data_bucket", Transform::Bucket(10)) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![2]) + .name("data_bucket") + .transform(Transform::Bucket(10)) + .build() + .unwrap(), + ) .unwrap() .build(); let builder = metadata.into_builder(Some("s3://bucket/table/metadata/v1.json".to_string())); @@ -3547,11 +3658,32 @@ mod tests { // Add spec 2: identity(id) + bucket(data) + year(timestamp) // Should reuse field_id 1000 for identity(id) and 1001 for bucket(data) let spec2 = UnboundPartitionSpec::builder() - .add_partition_field(1, "id", Transform::Identity) // Should reuse 1000 - .unwrap() - .add_partition_field(2, "data_bucket", Transform::Bucket(10)) // Should reuse 1001 - .unwrap() - .add_partition_field(3, "year", Transform::Year) // Should get new 1002 + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .name("id") + .transform(Transform::Identity) + .build() + .unwrap(), + ) // Should reuse 1000 + .unwrap() + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![2]) + .name("data_bucket") + .transform(Transform::Bucket(10)) + .build() + .unwrap(), + ) // Should reuse 1001 + .unwrap() + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![3]) + .name("year") + .transform(Transform::Year) + .build() + .unwrap(), + ) // Should get new 1002 .unwrap() .build(); let builder = metadata.into_builder(Some("s3://bucket/table/metadata/v2.json".to_string()));