From cdbd49534c6d2583a06e041799c5682a3de9e866 Mon Sep 17 00:00:00 2001 From: Noritaka Sekiyama Date: Wed, 23 Sep 2026 06:57:49 +0900 Subject: [PATCH 1/4] Make UnboundPartitionField fields private and carry source-ids The four public fields are replaced by accessors, and the derived builder's visibility is limited to the crate, so an instance can only be built inside the crate or read out of a partition spec JSON. That is what lets the struct guarantee its own invariant: source_ids always holds at least one id. source_ids: Vec replaces source_id: i32. source_id() returns the single id, and an error for a multi-argument field rather than quietly handing back the first one. The serde module reads either spelling -- source-id, or source-ids for a v3 multi-argument transform -- and writes back the one that matches the field. Binding a multi-argument field now fails with a clear error. A bound PartitionField still carries a single source_id, so until that struct is converted there is nowhere to put the extra ids, and failing loudly beats dropping them on the way into TableCreation or TableUpdate::AddSpec. --- crates/catalog/rest/src/catalog.rs | 10 +- crates/iceberg/public-api.txt | 12 +- .../arrow/record_batch_partition_splitter.rs | 26 +- .../src/expr/visitors/expression_evaluator.rs | 2 +- .../visitors/inclusive_metrics_evaluator.rs | 2 +- .../src/expr/visitors/inclusive_projection.rs | 48 +-- crates/iceberg/src/partitioning.rs | 70 ++-- crates/iceberg/src/spec/partition.rs | 349 ++++++++++++++---- crates/iceberg/src/spec/snapshot_summary.rs | 4 +- crates/iceberg/src/spec/table_metadata.rs | 84 +++-- .../src/spec/table_metadata_builder.rs | 145 ++++---- 11 files changed, 502 insertions(+), 250 deletions(-) diff --git a/crates/catalog/rest/src/catalog.rs b/crates/catalog/rest/src/catalog.rs index aa36a074e7..267ea27f71 100644 --- a/crates/catalog/rest/src/catalog.rs +++ b/crates/catalog/rest/src/catalog.rs @@ -1671,7 +1671,7 @@ mod tests { use iceberg::spec::{ FormatVersion, NestedField, NullOrder, Operation, PrimitiveType, Schema, Snapshot, SnapshotLog, SortDirection, SortField, SortOrder, Summary, Transform, Type, - UnboundPartitionField, UnboundPartitionSpec, + UnboundPartitionSpec, }; use iceberg::test_utils::test_runtime; use iceberg::transaction::{ApplyTransactionAction, Transaction}; @@ -3832,13 +3832,7 @@ mod tests { .properties(HashMap::from([("owner".to_string(), "testx".to_string())])) .partition_spec( UnboundPartitionSpec::builder() - .add_partition_fields(vec![ - UnboundPartitionField::builder() - .source_id(1) - .transform(Transform::Truncate(3)) - .name("id".to_string()) - .build(), - ]) + .add_partition_field(1, "id", Transform::Truncate(3)) .unwrap() .build(), ) diff --git a/crates/iceberg/public-api.txt b/crates/iceberg/public-api.txt index f29c5efe07..1db3dc4a02 100644 --- a/crates/iceberg/public-api.txt +++ b/crates/iceberg/public-api.txt @@ -2911,10 +2911,12 @@ 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 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 @@ -2925,8 +2927,6 @@ pub fn iceberg::spec::UnboundPartitionField::from(field: iceberg::spec::Partitio 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 -impl iceberg::spec::UnboundPartitionField -pub fn iceberg::spec::UnboundPartitionField::builder() -> UnboundPartitionFieldBuilder<((), (), (), ())> impl serde_core::ser::Serialize for iceberg::spec::UnboundPartitionField pub fn iceberg::spec::UnboundPartitionField::serialize<__S>(&self, __serializer: __S) -> core::result::Result<<__S as serde_core::ser::Serializer>::Ok, <__S as serde_core::ser::Serializer>::Error> where __S: serde_core::ser::Serializer impl<'de> serde_core::de::Deserialize<'de> for iceberg::spec::UnboundPartitionField diff --git a/crates/iceberg/src/arrow/record_batch_partition_splitter.rs b/crates/iceberg/src/arrow/record_batch_partition_splitter.rs index 6b3aa73d22..fb3bdc5096 100644 --- a/crates/iceberg/src/arrow/record_batch_partition_splitter.rs +++ b/crates/iceberg/src/arrow/record_batch_partition_splitter.rs @@ -244,12 +244,13 @@ 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() .build() .unwrap(), @@ -358,12 +359,13 @@ 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() .build() .unwrap(), diff --git a/crates/iceberg/src/expr/visitors/expression_evaluator.rs b/crates/iceberg/src/expr/visitors/expression_evaluator.rs index 122f195897..449180869a 100644 --- a/crates/iceberg/src/expr/visitors/expression_evaluator.rs +++ b/crates/iceberg/src/expr/visitors/expression_evaluator.rs @@ -280,7 +280,7 @@ 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) diff --git a/crates/iceberg/src/expr/visitors/inclusive_metrics_evaluator.rs b/crates/iceberg/src/expr/visitors/inclusive_metrics_evaluator.rs index 06c92ab3e8..193ab6b367 100644 --- a/crates/iceberg/src/expr/visitors/inclusive_metrics_evaluator.rs +++ b/crates/iceberg/src/expr/visitors/inclusive_metrics_evaluator.rs @@ -1661,7 +1661,7 @@ 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) diff --git a/crates/iceberg/src/expr/visitors/inclusive_projection.rs b/crates/iceberg/src/expr/visitors/inclusive_projection.rs index d9544e4c47..f7a6622bd7 100644 --- a/crates/iceberg/src/expr/visitors/inclusive_projection.rs +++ b/crates/iceberg/src/expr/visitors/inclusive_projection.rs @@ -304,7 +304,7 @@ 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) @@ -339,12 +339,14 @@ 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() .build() .unwrap(); @@ -374,12 +376,14 @@ 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() .build() .unwrap(); @@ -409,12 +413,14 @@ 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() .build() .unwrap(); @@ -446,7 +452,7 @@ 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)) @@ -486,7 +492,7 @@ 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)) diff --git a/crates/iceberg/src/partitioning.rs b/crates/iceberg/src/partitioning.rs index ef230ef9e5..d388144989 100644 --- a/crates/iceberg/src/partitioning.rs +++ b/crates/iceberg/src/partitioning.rs @@ -261,12 +261,14 @@ 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( + crate::spec::UnboundPartitionField::builder() + .source_ids(vec![4]) + .field_id(1000) + .name("cat_old".to_string()) + .transform(Transform::Identity) + .build(), + ) .unwrap() .build() .unwrap(); @@ -274,12 +276,14 @@ 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( + crate::spec::UnboundPartitionField::builder() + .source_ids(vec![4]) + .field_id(1000) + .name("cat_new".to_string()) + .transform(Transform::Identity) + .build(), + ) .unwrap() .build() .unwrap(); @@ -297,12 +301,14 @@ 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( + crate::spec::UnboundPartitionField::builder() + .source_ids(vec![4]) + .field_id(1000) + .name("category".to_string()) + .transform(Transform::Identity) + .build(), + ) .unwrap() .build() .unwrap(); @@ -310,12 +316,14 @@ 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( + crate::spec::UnboundPartitionField::builder() + .source_ids(vec![4]) + .field_id(1000) + .name("category_v2".to_string()) + .transform(Transform::Void) + .build(), + ) .unwrap() .build() .unwrap(); @@ -371,12 +379,14 @@ 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( + crate::spec::UnboundPartitionField::builder() + .source_ids(vec![4]) + .field_id(spec_v0.fields()[0].field_id) + .name("category".to_string()) + .transform(Transform::Identity) + .build(), + ) .unwrap() .add_partition_field("ts", "ts_year", Transform::Year) .unwrap() diff --git a/crates/iceberg/src/spec/partition.rs b/crates/iceberg/src/spec/partition.rs index 3999c0d010..8b15c72f10 100644 --- a/crates/iceberg/src/spec/partition.rs +++ b/crates/iceberg/src/spec/partition.rs @@ -264,20 +264,147 @@ 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 known to be well formed once built: in +/// particular `source_ids` always holds at least one id. Construct one through +/// [`UnboundPartitionSpec::builder`], or read one out of a partition spec JSON. #[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone, TypedBuilder)] -#[serde(rename_all = "kebab-case")] +#[serde( + try_from = "self::_serde::UnboundPartitionFieldSerde", + into = "self::_serde::UnboundPartitionFieldSerde" +)] +#[builder( + builder_method(vis = "pub(crate)"), + builder_type(vis = "pub(crate)"), + build_method(vis = "pub(crate)") +)] 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, + 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(crate) fn with_field_id(self, field_id: i32) -> Self { + Self { + field_id: Some(field_id), + ..self + } + } +} + +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`. Both spellings are read; 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)) if !source_ids.is_empty() => source_ids, + (None, Some(_)) => { + return Err(invalid_data!("Empty source-ids is not allowed")); + } + (Some(source_id), Some(source_ids)) => { + // Tolerated for readers, but the two must agree + if source_ids.first() != Some(&source_id) { + return Err(invalid_data!( + "source-id {source_id} does not match the first entry of source-ids {source_ids:?}" + )); + } + source_ids + } + (None, None) => { + return Err(invalid_data!("missing field `source-id`")); + } + }; + + Ok(UnboundPartitionField { + source_ids, + field_id: value.field_id, + name: value.name, + transform: value.transform, + }) + } + } + + impl From for UnboundPartitionFieldSerde { + fn from(value: UnboundPartitionField) -> Self { + let multi_arg = value.source_ids.len() > 1; + Self { + source_id: (!multi_arg).then(|| value.source_ids[0]), + source_ids: multi_arg.then_some(value.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 +468,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, @@ -388,7 +515,7 @@ impl UnboundPartitionSpecBuilder { transformation: Transform, ) -> Result { let field = UnboundPartitionField { - source_id, + source_ids: vec![source_id], field_id: None, name: target_name.to_string(), transform: transformation, @@ -410,7 +537,7 @@ impl UnboundPartitionSpecBuilder { 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)?; + 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)?; } @@ -495,7 +622,7 @@ impl PartitionSpecBuilder { })? .id; let field = UnboundPartitionField { - source_id, + source_ids: vec![source_id], field_id: None, name: target_name.into(), transform, @@ -510,7 +637,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 +701,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 +722,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 +746,9 @@ 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 - ) + 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 +793,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,12 +908,12 @@ 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(), UnboundPartitionField::builder() - .source_id(2) + .source_ids(vec![2]) .name("name_string".to_string()) .transform(Transform::Void) .build(), @@ -802,12 +931,12 @@ 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(), UnboundPartitionField::builder() - .source_id(2) + .source_ids(vec![2]) .name("name_void".to_string()) .transform(Transform::Void) .build(), @@ -848,17 +977,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 +1004,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); @@ -1148,14 +1277,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 +1305,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 +1320,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 +1391,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 +1418,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 +1431,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 +1455,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 +1476,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 +1504,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, @@ -1415,7 +1544,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 +1557,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 +1568,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 +1589,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 +1601,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 +1625,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 +1637,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 +1662,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 +1674,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 +1699,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 +1718,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 +1760,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 +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(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 +1826,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 +1860,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 +1976,92 @@ 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"}"#, + "missing field `source-id`", + ), + ( + r#"{"source-id": 1, "source-ids": [2, 3], "name": "m", "transform": "identity"}"#, + "does not match the first entry", + ), + ] { + 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}" + ); + } } diff --git a/crates/iceberg/src/spec/snapshot_summary.rs b/crates/iceberg/src/spec/snapshot_summary.rs index a48a5abdda..db3348d825 100644 --- a/crates/iceberg/src/spec/snapshot_summary.rs +++ b/crates/iceberg/src/spec/snapshot_summary.rs @@ -828,7 +828,7 @@ 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(), @@ -978,7 +978,7 @@ 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(), diff --git a/crates/iceberg/src/spec/table_metadata.rs b/crates/iceberg/src/spec/table_metadata.rs index 24f10355bf..c698aafa17 100644 --- a/crates/iceberg/src/spec/table_metadata.rs +++ b/crates/iceberg/src/spec/table_metadata.rs @@ -1721,12 +1721,14 @@ 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() .build() .unwrap(); @@ -1869,12 +1871,14 @@ 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() .build() .unwrap(); @@ -2945,12 +2949,14 @@ 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() .build() .unwrap(); @@ -3042,12 +3048,14 @@ 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() .build() .unwrap(); @@ -3171,12 +3179,14 @@ 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() .build() .unwrap(); @@ -3257,12 +3267,14 @@ 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() .build() .unwrap(); diff --git a/crates/iceberg/src/spec/table_metadata_builder.rs b/crates/iceberg/src/spec/table_metadata_builder.rs index 3ed848d2d6..3c1aa622c9 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,19 +850,19 @@ 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); + return field.with_field_id(existing_field_id); } field }) @@ -1228,15 +1228,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()?; @@ -1666,12 +1670,14 @@ 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() .build() .unwrap() @@ -1739,20 +1745,19 @@ 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(), + // 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() .build(); @@ -1767,19 +1772,23 @@ 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() - .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() .build() .unwrap(); @@ -1831,12 +1840,14 @@ 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() .build() .unwrap(); @@ -2401,18 +2412,18 @@ 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(), + UnboundPartitionField::builder() + .source_ids(vec![3]) + .field_id(1002) + .name("z".to_string()) + .transform(Transform::Identity) + .build(), ]) .unwrap() .build(); From 8fcd7957fed90848c4b93a87513b995558705d06 Mon Sep 17 00:00:00 2001 From: Noritaka Sekiyama Date: Mon, 28 Sep 2026 09:02:07 +0900 Subject: [PATCH 2/4] Reject source-id and source-ids together in UnboundPartitionField PyIceberg already treats the two spellings as mutually exclusive, and no writer emits both, so accepting a matching pair only let the two readers disagree. Also name both spellings in the missing-id error and note why check_transform_compatibility rejects a multi-argument field. Co-authored-by: Isaac --- crates/iceberg/src/spec/partition.rs | 27 ++++++++++++++------------- 1 file changed, 14 insertions(+), 13 deletions(-) diff --git a/crates/iceberg/src/spec/partition.rs b/crates/iceberg/src/spec/partition.rs index 8b15c72f10..5e6917ecca 100644 --- a/crates/iceberg/src/spec/partition.rs +++ b/crates/iceberg/src/spec/partition.rs @@ -346,7 +346,8 @@ mod _serde { use crate::spec::Transform; /// Per the spec a single-argument field carries `source-id` and a multi-argument field - /// carries `source-ids`. Both spellings are read; the one that matches the field is written. + /// 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 { @@ -370,17 +371,15 @@ mod _serde { (None, Some(_)) => { return Err(invalid_data!("Empty source-ids is not allowed")); } - (Some(source_id), Some(source_ids)) => { - // Tolerated for readers, but the two must agree - if source_ids.first() != Some(&source_id) { - return Err(invalid_data!( - "source-id {source_id} does not match the first entry of source-ids {source_ids:?}" - )); - } - source_ids + (Some(_), Some(_)) => { + return Err(invalid_data!( + "source-id and source-ids are mutually exclusive" + )); } (None, None) => { - return Err(invalid_data!("missing field `source-id`")); + return Err(invalid_data!( + "Either `source-id` or `source-ids` must be present" + )); } }; @@ -746,6 +745,8 @@ 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<()> { + // 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") @@ -2026,11 +2027,11 @@ mod tests { ), ( r#"{"name": "m", "transform": "identity"}"#, - "missing field `source-id`", + "Either `source-id` or `source-ids` must be present", ), ( - r#"{"source-id": 1, "source-ids": [2, 3], "name": "m", "transform": "identity"}"#, - "does not match the first entry", + r#"{"source-id": 1, "source-ids": [1], "name": "m", "transform": "identity"}"#, + "mutually exclusive", ), ] { let err = serde_json::from_str::(input).unwrap_err(); From 1d48025bb64515a66d0b71a46620e9bce41aa47b Mon Sep 17 00:00:00 2001 From: Noritaka Sekiyama Date: Tue, 29 Sep 2026 12:30:59 +0900 Subject: [PATCH 3/4] Make UnboundPartitionField buildable and pass it to add_partition_field Address review: the UnboundPartitionField builder and with_field_id are public, since an unbound field is not expected to be validated yet, and UnboundPartitionSpecBuilder::add_partition_field now takes the field rather than a single source id, name and transform. With the builder public an empty source_ids can be built, so serializing one no longer panics and both spec builders reject it with a clear error. Co-authored-by: Isaac --- crates/catalog/rest/src/catalog.rs | 10 +- crates/iceberg/public-api.txt | 5 +- crates/iceberg/src/catalog/mod.rs | 28 +++- crates/iceberg/src/partitioning.rs | 22 ++- crates/iceberg/src/scan/cache.rs | 10 +- crates/iceberg/src/spec/partition.rs | 147 ++++++++++++------ crates/iceberg/src/spec/table_metadata.rs | 16 +- .../src/spec/table_metadata_builder.rs | 133 +++++++++++++--- 8 files changed, 290 insertions(+), 81 deletions(-) diff --git a/crates/catalog/rest/src/catalog.rs b/crates/catalog/rest/src/catalog.rs index 267ea27f71..caa8918043 100644 --- a/crates/catalog/rest/src/catalog.rs +++ b/crates/catalog/rest/src/catalog.rs @@ -1671,7 +1671,7 @@ mod tests { use iceberg::spec::{ FormatVersion, NestedField, NullOrder, Operation, PrimitiveType, Schema, Snapshot, SnapshotLog, SortDirection, SortField, SortOrder, Summary, Transform, Type, - UnboundPartitionSpec, + UnboundPartitionField, UnboundPartitionSpec, }; use iceberg::test_utils::test_runtime; use iceberg::transaction::{ApplyTransactionAction, Transaction}; @@ -3832,7 +3832,13 @@ mod tests { .properties(HashMap::from([("owner".to_string(), "testx".to_string())])) .partition_spec( UnboundPartitionSpec::builder() - .add_partition_field(1, "id", Transform::Truncate(3)) + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .name("id") + .transform(Transform::Truncate(3)) + .build(), + ) .unwrap() .build(), ) diff --git a/crates/iceberg/public-api.txt b/crates/iceberg/public-api.txt index 1db3dc4a02..024658fc78 100644 --- a/crates/iceberg/public-api.txt +++ b/crates/iceberg/public-api.txt @@ -2917,6 +2917,7 @@ 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 @@ -2927,6 +2928,8 @@ pub fn iceberg::spec::UnboundPartitionField::from(field: iceberg::spec::Partitio 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 +impl iceberg::spec::UnboundPartitionField +pub fn iceberg::spec::UnboundPartitionField::builder() -> UnboundPartitionFieldBuilder<((), (), (), ())> impl serde_core::ser::Serialize for iceberg::spec::UnboundPartitionField pub fn iceberg::spec::UnboundPartitionField::serialize<__S>(&self, __serializer: __S) -> core::result::Result<<__S as serde_core::ser::Serializer>::Ok, <__S as serde_core::ser::Serializer>::Error> where __S: serde_core::ser::Serializer impl<'de> serde_core::de::Deserialize<'de> for iceberg::spec::UnboundPartitionField @@ -2956,7 +2959,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/catalog/mod.rs b/crates/iceberg/src/catalog/mod.rs index cf7a94d20f..b9e8f33d62 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,29 @@ 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() - .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() - .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() .build(), }, diff --git a/crates/iceberg/src/partitioning.rs b/crates/iceberg/src/partitioning.rs index d388144989..2eae11175d 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,13 @@ 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(); } builder.build().bind(schema.clone()).unwrap() @@ -262,7 +270,7 @@ mod tests { let spec_v0 = PartitionSpec::builder(Arc::new(schema.clone())) .with_spec_id(0) .add_unbound_field( - crate::spec::UnboundPartitionField::builder() + UnboundPartitionField::builder() .source_ids(vec![4]) .field_id(1000) .name("cat_old".to_string()) @@ -277,7 +285,7 @@ mod tests { let spec_v1 = PartitionSpec::builder(Arc::new(schema.clone())) .with_spec_id(1) .add_unbound_field( - crate::spec::UnboundPartitionField::builder() + UnboundPartitionField::builder() .source_ids(vec![4]) .field_id(1000) .name("cat_new".to_string()) @@ -302,7 +310,7 @@ mod tests { let spec_v0 = PartitionSpec::builder(Arc::new(schema.clone())) .with_spec_id(0) .add_unbound_field( - crate::spec::UnboundPartitionField::builder() + UnboundPartitionField::builder() .source_ids(vec![4]) .field_id(1000) .name("category".to_string()) @@ -317,7 +325,7 @@ mod tests { let spec_v1 = PartitionSpec::builder(Arc::new(schema.clone())) .with_spec_id(1) .add_unbound_field( - crate::spec::UnboundPartitionField::builder() + UnboundPartitionField::builder() .source_ids(vec![4]) .field_id(1000) .name("category_v2".to_string()) @@ -380,7 +388,7 @@ mod tests { let spec_v1 = PartitionSpec::builder(Arc::new(schema.clone())) .with_spec_id(1) .add_unbound_field( - crate::spec::UnboundPartitionField::builder() + UnboundPartitionField::builder() .source_ids(vec![4]) .field_id(spec_v0.fields()[0].field_id) .name("category".to_string()) diff --git a/crates/iceberg/src/scan/cache.rs b/crates/iceberg/src/scan/cache.rs index b5eb3c4ce0..3013c7fdd3 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,13 @@ 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() .build(); diff --git a/crates/iceberg/src/spec/partition.rs b/crates/iceberg/src/spec/partition.rs index 5e6917ecca..391c56ba1c 100644 --- a/crates/iceberg/src/spec/partition.rs +++ b/crates/iceberg/src/spec/partition.rs @@ -265,19 +265,14 @@ impl PartitionKey { 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 known to be well formed once built: in -/// particular `source_ids` always holds at least one id. Construct one through -/// [`UnboundPartitionSpec::builder`], or read one out of a partition spec JSON. +/// Being unbound, a field built through [`UnboundPartitionField::builder`] is not validated; +/// it is checked when added to an [`UnboundPartitionSpecBuilder`] and again when bound to a +/// schema. #[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone, TypedBuilder)] #[serde( try_from = "self::_serde::UnboundPartitionFieldSerde", into = "self::_serde::UnboundPartitionFieldSerde" )] -#[builder( - builder_method(vis = "pub(crate)"), - builder_type(vis = "pub(crate)"), - build_method(vis = "pub(crate)") -)] pub struct UnboundPartitionField { /// The source column ids from the table’s schema. A single-argument transform reads one /// id; a v3 multi-argument transform reads several. @@ -287,6 +282,7 @@ pub struct UnboundPartitionField { #[builder(default, setter(strip_option(fallback = field_id_opt)))] field_id: Option, /// A partition name. + #[builder(setter(into))] name: String, /// A transform that is applied to the source column to produce a partition value. transform: Transform, @@ -308,7 +304,8 @@ impl UnboundPartitionField { } } - /// The source column ids this field reads, in order. Never empty. + /// The source column ids this field reads, in order. Never empty once the field is part + /// of a partition spec. pub fn source_ids(&self) -> &[i32] { &self.source_ids } @@ -329,7 +326,7 @@ impl UnboundPartitionField { } /// Return this field with the given partition field id assigned. - pub(crate) fn with_field_id(self, field_id: i32) -> Self { + pub fn with_field_id(self, field_id: i32) -> Self { Self { field_id: Some(field_id), ..self @@ -394,10 +391,13 @@ mod _serde { impl From for UnboundPartitionFieldSerde { fn from(value: UnboundPartitionField) -> Self { - let multi_arg = value.source_ids.len() > 1; + 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: (!multi_arg).then(|| value.source_ids[0]), - source_ids: multi_arg.then_some(value.source_ids), + source_id, + source_ids, field_id: value.field_id, name: value.name, transform: value.transform, @@ -507,19 +507,15 @@ 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_ids: vec![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_source_ids_set(&field)?; + 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. @@ -529,21 +525,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_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) - } - /// Build the unbound partition spec. pub fn build(self) -> UnboundPartitionSpec { UnboundPartitionSpec { @@ -636,6 +622,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_source_ids_set(&field)?; 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)?; @@ -793,6 +780,17 @@ trait CorePartitionSpecValidator { Ok(()) } + /// Ensure that the partition field reads at least one source column. + fn check_source_ids_set(&self, field: &UnboundPartitionField) -> Result<()> { + if field.source_ids.is_empty() { + return Err(invalid_data!( + "Partition field '{}' has no source id", + field.name + )); + } + Ok(()) + } + /// For a single source-column transformations must be unique. fn check_for_redundant_partitions( &self, @@ -1014,7 +1012,13 @@ 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() .build(); @@ -1261,9 +1265,21 @@ 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() - .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_err(); } @@ -1522,12 +1538,20 @@ 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() .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_err(); assert!(err.message().contains("redundant partition")); @@ -2065,4 +2089,39 @@ mod tests { "unexpected error: {err}" ); } + + #[test] + fn test_unbound_partition_field_without_source_ids() { + let field = UnboundPartitionField::builder() + .source_ids(vec![]) + .name("m") + .transform(Transform::Identity) + .build(); + + // Unbound, so building and serializing it does not validate + let serialized = serde_json::to_value(&field).unwrap(); + assert_eq!(Some(&serde_json::json!([])), serialized.get("source-ids")); + + let err = UnboundPartitionSpec::builder() + .add_partition_field(field.clone()) + .unwrap_err(); + assert!( + err.to_string().contains("has no source id"), + "unexpected error: {err}" + ); + + let schema = Schema::builder() + .with_fields(vec![ + NestedField::required(1, "a", Type::Primitive(PrimitiveType::Int)).into(), + ]) + .build() + .unwrap(); + let err = PartitionSpec::builder(schema) + .add_unbound_field(field) + .unwrap_err(); + assert!( + err.to_string().contains("has no source id"), + "unexpected error: {err}" + ); + } } diff --git a/crates/iceberg/src/spec/table_metadata.rs b/crates/iceberg/src/spec/table_metadata.rs index c698aafa17..dbe4d451d1 100644 --- a/crates/iceberg/src/spec/table_metadata.rs +++ b/crates/iceberg/src/spec/table_metadata.rs @@ -4493,7 +4493,13 @@ 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() .build(), SortOrder::unsorted_order(), @@ -4508,7 +4514,13 @@ 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() .build(), ) diff --git a/crates/iceberg/src/spec/table_metadata_builder.rs b/crates/iceberg/src/spec/table_metadata_builder.rs index 3c1aa622c9..009f6e78c9 100644 --- a/crates/iceberg/src/spec/table_metadata_builder.rs +++ b/crates/iceberg/src/spec/table_metadata_builder.rs @@ -862,9 +862,10 @@ impl TableMetadataBuilder { && let Some(&existing_field_id) = equivalent_field_ids.get(&(field.source_ids().to_vec(), field.transform())) { - return field.with_field_id(existing_field_id); + field.with_field_id(existing_field_id) + } else { + field } - field }) .collect(); @@ -1436,7 +1437,13 @@ 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() .build() } @@ -1826,7 +1833,13 @@ 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() .build(); @@ -2748,7 +2761,13 @@ 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() .build(); @@ -2807,7 +2826,13 @@ 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() .build(); @@ -2856,7 +2881,13 @@ 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() .build(); @@ -2879,7 +2910,13 @@ 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() .build(); @@ -2909,7 +2946,13 @@ 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() .build(); @@ -3005,7 +3048,13 @@ 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() .build(); @@ -3054,7 +3103,13 @@ 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() .build(); @@ -3104,7 +3159,13 @@ 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() .build(); @@ -3128,7 +3189,13 @@ 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() .build(); @@ -3529,7 +3596,13 @@ 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() .build(); @@ -3548,7 +3621,13 @@ 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() .build(); let builder = metadata.into_builder(Some("s3://bucket/table/metadata/v1.json".to_string())); @@ -3558,11 +3637,29 @@ 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 + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![1]) + .name("id") + .transform(Transform::Identity) + .build(), + ) // Should reuse 1000 .unwrap() - .add_partition_field(2, "data_bucket", Transform::Bucket(10)) // Should reuse 1001 + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![2]) + .name("data_bucket") + .transform(Transform::Bucket(10)) + .build(), + ) // Should reuse 1001 .unwrap() - .add_partition_field(3, "year", Transform::Year) // Should get new 1002 + .add_partition_field( + UnboundPartitionField::builder() + .source_ids(vec![3]) + .name("year") + .transform(Transform::Year) + .build(), + ) // Should get new 1002 .unwrap() .build(); let builder = metadata.into_builder(Some("s3://bucket/table/metadata/v2.json".to_string())); From 4e72beebaf9d88fad412a1e0a23f0c3f699429da Mon Sep 17 00:00:00 2001 From: Noritaka Sekiyama Date: Tue, 29 Sep 2026 20:33:43 +0900 Subject: [PATCH 4/4] Validate UnboundPartitionField when it is built Address review: the builder's build() now returns Result and rejects an empty source_ids, following FileScanTask. Deserialization goes through the same builder, so the check lives in one place and the spec builders no longer need their own. Co-authored-by: Isaac --- crates/catalog/rest/src/catalog.rs | 3 +- crates/iceberg/public-api.txt | 2 + .../arrow/record_batch_partition_splitter.rs | 6 +- crates/iceberg/src/catalog/mod.rs | 9 +- .../src/expr/visitors/expression_evaluator.rs | 3 +- .../visitors/inclusive_metrics_evaluator.rs | 3 +- .../src/expr/visitors/inclusive_projection.rs | 18 ++- crates/iceberg/src/partitioning.rs | 18 ++- crates/iceberg/src/scan/cache.rs | 3 +- crates/iceberg/src/spec/partition.rs | 108 ++++++++---------- crates/iceberg/src/spec/snapshot_summary.rs | 6 +- crates/iceberg/src/spec/table_metadata.rs | 24 ++-- .../src/spec/table_metadata_builder.rs | 70 ++++++++---- 13 files changed, 158 insertions(+), 115 deletions(-) diff --git a/crates/catalog/rest/src/catalog.rs b/crates/catalog/rest/src/catalog.rs index caa8918043..b7d3b3163c 100644 --- a/crates/catalog/rest/src/catalog.rs +++ b/crates/catalog/rest/src/catalog.rs @@ -3837,7 +3837,8 @@ mod tests { .source_ids(vec![1]) .name("id") .transform(Transform::Truncate(3)) - .build(), + .build() + .unwrap(), ) .unwrap() .build(), diff --git a/crates/iceberg/public-api.txt b/crates/iceberg/public-api.txt index 024658fc78..19a5d49b17 100644 --- a/crates/iceberg/public-api.txt +++ b/crates/iceberg/public-api.txt @@ -2925,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 diff --git a/crates/iceberg/src/arrow/record_batch_partition_splitter.rs b/crates/iceberg/src/arrow/record_batch_partition_splitter.rs index fb3bdc5096..0b679040c6 100644 --- a/crates/iceberg/src/arrow/record_batch_partition_splitter.rs +++ b/crates/iceberg/src/arrow/record_batch_partition_splitter.rs @@ -249,7 +249,8 @@ mod tests { .source_ids(vec![1]) .name("id_bucket".to_string()) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .build() @@ -364,7 +365,8 @@ mod tests { .source_ids(vec![1]) .name("id_bucket".to_string()) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .build() diff --git a/crates/iceberg/src/catalog/mod.rs b/crates/iceberg/src/catalog/mod.rs index b9e8f33d62..88681040cd 100644 --- a/crates/iceberg/src/catalog/mod.rs +++ b/crates/iceberg/src/catalog/mod.rs @@ -1664,7 +1664,8 @@ mod tests { .source_ids(vec![4]) .name("ts_day") .transform(Transform::Day) - .build(), + .build() + .unwrap(), ) .unwrap() .add_partition_field( @@ -1672,7 +1673,8 @@ mod tests { .source_ids(vec![1]) .name("id_bucket") .transform(Transform::Bucket(16)) - .build(), + .build() + .unwrap(), ) .unwrap() .add_partition_field( @@ -1680,7 +1682,8 @@ mod tests { .source_ids(vec![2]) .name("id_truncate") .transform(Transform::Truncate(4)) - .build(), + .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 449180869a..31ce94a99b 100644 --- a/crates/iceberg/src/expr/visitors/expression_evaluator.rs +++ b/crates/iceberg/src/expr/visitors/expression_evaluator.rs @@ -284,7 +284,8 @@ mod tests { .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 193ab6b367..8ccd1364d9 100644 --- a/crates/iceberg/src/expr/visitors/inclusive_metrics_evaluator.rs +++ b/crates/iceberg/src/expr/visitors/inclusive_metrics_evaluator.rs @@ -1665,7 +1665,8 @@ mod test { .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 f7a6622bd7..88508af844 100644 --- a/crates/iceberg/src/expr/visitors/inclusive_projection.rs +++ b/crates/iceberg/src/expr/visitors/inclusive_projection.rs @@ -308,7 +308,8 @@ mod tests { .name("a".to_string()) .field_id(1) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .build() @@ -345,7 +346,8 @@ mod tests { .field_id(1000) .name("year".to_string()) .transform(Transform::Year) - .build(), + .build() + .unwrap(), ]) .unwrap() .build() @@ -382,7 +384,8 @@ mod tests { .field_id(1000) .name("month".to_string()) .transform(Transform::Month) - .build(), + .build() + .unwrap(), ]) .unwrap() .build() @@ -419,7 +422,8 @@ mod tests { .field_id(1000) .name("day".to_string()) .transform(Transform::Day) - .build(), + .build() + .unwrap(), ]) .unwrap() .build() @@ -456,7 +460,8 @@ mod tests { .name("name_truncate".to_string()) .field_id(3) .transform(Transform::Truncate(4)) - .build(), + .build() + .unwrap(), ) .unwrap() .build() @@ -496,7 +501,8 @@ mod tests { .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 2eae11175d..6ac62c68c4 100644 --- a/crates/iceberg/src/partitioning.rs +++ b/crates/iceberg/src/partitioning.rs @@ -206,7 +206,8 @@ mod tests { .source_ids(vec![source_id]) .name(name) .transform(transform) - .build(), + .build() + .unwrap(), ) .unwrap(); } @@ -275,7 +276,8 @@ mod tests { .field_id(1000) .name("cat_old".to_string()) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .build() @@ -290,7 +292,8 @@ mod tests { .field_id(1000) .name("cat_new".to_string()) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .build() @@ -315,7 +318,8 @@ mod tests { .field_id(1000) .name("category".to_string()) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .build() @@ -330,7 +334,8 @@ mod tests { .field_id(1000) .name("category_v2".to_string()) .transform(Transform::Void) - .build(), + .build() + .unwrap(), ) .unwrap() .build() @@ -393,7 +398,8 @@ mod tests { .field_id(spec_v0.fields()[0].field_id) .name("category".to_string()) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .add_partition_field("ts", "ts_year", Transform::Year) diff --git a/crates/iceberg/src/scan/cache.rs b/crates/iceberg/src/scan/cache.rs index 3013c7fdd3..b87786e6cf 100644 --- a/crates/iceberg/src/scan/cache.rs +++ b/crates/iceberg/src/scan/cache.rs @@ -293,7 +293,8 @@ mod tests { .source_ids(vec![2]) .name("part") .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .build(); diff --git a/crates/iceberg/src/spec/partition.rs b/crates/iceberg/src/spec/partition.rs index 391c56ba1c..17aa3de968 100644 --- a/crates/iceberg/src/spec/partition.rs +++ b/crates/iceberg/src/spec/partition.rs @@ -265,14 +265,15 @@ impl PartitionKey { pub type UnboundPartitionSpecRef = Arc; /// Unbound partition field can be built without a schema and later bound to a schema. /// -/// Being unbound, a field built through [`UnboundPartitionField::builder`] is not validated; -/// it is checked when added to an [`UnboundPartitionSpecBuilder`] and again when 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( try_from = "self::_serde::UnboundPartitionFieldSerde", into = "self::_serde::UnboundPartitionFieldSerde" )] +#[builder(build_method(into = Result))] pub struct UnboundPartitionField { /// The source column ids from the table’s schema. A single-argument transform reads one /// id; a v3 multi-argument transform reads several. @@ -304,8 +305,7 @@ impl UnboundPartitionField { } } - /// The source column ids this field reads, in order. Never empty once the field is part - /// of a partition spec. + /// The source column ids this field reads, in order. Never empty. pub fn source_ids(&self) -> &[i32] { &self.source_ids } @@ -332,6 +332,20 @@ impl UnboundPartitionField { ..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 { @@ -364,10 +378,7 @@ mod _serde { 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)) if !source_ids.is_empty() => source_ids, - (None, Some(_)) => { - return Err(invalid_data!("Empty source-ids is not allowed")); - } + (None, Some(source_ids)) => source_ids, (Some(_), Some(_)) => { return Err(invalid_data!( "source-id and source-ids are mutually exclusive" @@ -380,12 +391,12 @@ mod _serde { } }; - Ok(UnboundPartitionField { - source_ids, - field_id: value.field_id, - name: value.name, - transform: value.transform, - }) + UnboundPartitionField::builder() + .source_ids(source_ids) + .field_id_opt(value.field_id) + .name(value.name) + .transform(value.transform) + .build() } } @@ -509,7 +520,6 @@ impl UnboundPartitionSpecBuilder { /// Add a new partition field to the partition spec from an unbound partition field. pub fn add_partition_field(mut self, field: UnboundPartitionField) -> Result { self.check_name_set_and_unique(&field.name)?; - self.check_source_ids_set(&field)?; 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)?; @@ -622,7 +632,6 @@ 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_source_ids_set(&field)?; 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)?; @@ -780,17 +789,6 @@ trait CorePartitionSpecValidator { Ok(()) } - /// Ensure that the partition field reads at least one source column. - fn check_source_ids_set(&self, field: &UnboundPartitionField) -> Result<()> { - if field.source_ids.is_empty() { - return Err(invalid_data!( - "Partition field '{}' has no source id", - field.name - )); - } - Ok(()) - } - /// For a single source-column transformations must be unique. fn check_for_redundant_partitions( &self, @@ -910,12 +908,14 @@ mod tests { .source_ids(vec![1]) .name("id".to_string()) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), UnboundPartitionField::builder() .source_ids(vec![2]) .name("name_string".to_string()) .transform(Transform::Void) - .build(), + .build() + .unwrap(), ]) .unwrap() .with_spec_id(1) @@ -933,12 +933,14 @@ mod tests { .source_ids(vec![1]) .name("id_void".to_string()) .transform(Transform::Void) - .build(), + .build() + .unwrap(), UnboundPartitionField::builder() .source_ids(vec![2]) .name("name_void".to_string()) .transform(Transform::Void) - .build(), + .build() + .unwrap(), ]) .unwrap() .build() @@ -1017,7 +1019,8 @@ mod tests { .source_ids(vec![4]) .name("ts_day") .transform(Transform::Day) - .build(), + .build() + .unwrap(), ) .unwrap() .build(); @@ -1270,7 +1273,8 @@ mod tests { .source_ids(vec![1]) .name("ts_day") .transform(Transform::Day) - .build(), + .build() + .unwrap(), ) .unwrap() .add_partition_field( @@ -1278,7 +1282,8 @@ mod tests { .source_ids(vec![2]) .name("ts_day") .transform(Transform::Day) - .build(), + .build() + .unwrap(), ) .unwrap_err(); } @@ -1543,7 +1548,8 @@ mod tests { .source_ids(vec![1]) .name("id_bucket[16]") .transform(Transform::Bucket(16)) - .build(), + .build() + .unwrap(), ) .unwrap() .add_partition_field( @@ -1551,7 +1557,8 @@ mod tests { .source_ids(vec![1]) .name("id_bucket_with_other_name") .transform(Transform::Bucket(16)) - .build(), + .build() + .unwrap(), ) .unwrap_err(); assert!(err.message().contains("redundant partition")); @@ -2091,36 +2098,15 @@ mod tests { } #[test] - fn test_unbound_partition_field_without_source_ids() { - let field = UnboundPartitionField::builder() + fn test_unbound_partition_field_builder_rejects_empty_source_ids() { + let err = UnboundPartitionField::builder() .source_ids(vec![]) .name("m") .transform(Transform::Identity) - .build(); - - // Unbound, so building and serializing it does not validate - let serialized = serde_json::to_value(&field).unwrap(); - assert_eq!(Some(&serde_json::json!([])), serialized.get("source-ids")); - - let err = UnboundPartitionSpec::builder() - .add_partition_field(field.clone()) - .unwrap_err(); - assert!( - err.to_string().contains("has no source id"), - "unexpected error: {err}" - ); - - let schema = Schema::builder() - .with_fields(vec![ - NestedField::required(1, "a", Type::Primitive(PrimitiveType::Int)).into(), - ]) .build() - .unwrap(); - let err = PartitionSpec::builder(schema) - .add_unbound_field(field) .unwrap_err(); assert!( - err.to_string().contains("has no source id"), + 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 db3348d825..f8e69c4c86 100644 --- a/crates/iceberg/src/spec/snapshot_summary.rs +++ b/crates/iceberg/src/spec/snapshot_summary.rs @@ -831,7 +831,8 @@ mod tests { .source_ids(vec![2]) .name("year".to_string()) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ]) .unwrap() .with_spec_id(1) @@ -981,7 +982,8 @@ mod tests { .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 dbe4d451d1..83468602cd 100644 --- a/crates/iceberg/src/spec/table_metadata.rs +++ b/crates/iceberg/src/spec/table_metadata.rs @@ -1727,7 +1727,8 @@ mod tests { .field_id(1000) .name("ts_day".to_string()) .transform(Transform::Day) - .build(), + .build() + .unwrap(), ) .unwrap() .build() @@ -1877,7 +1878,8 @@ mod tests { .field_id(1000) .name("ts_day".to_string()) .transform(Transform::Day) - .build(), + .build() + .unwrap(), ) .unwrap() .build() @@ -2955,7 +2957,8 @@ mod tests { .field_id(1000) .name("x".to_string()) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .build() @@ -3054,7 +3057,8 @@ mod tests { .field_id(1000) .name("x".to_string()) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .build() @@ -3185,7 +3189,8 @@ mod tests { .field_id(1000) .name("x".to_string()) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .build() @@ -3273,7 +3278,8 @@ mod tests { .field_id(1000) .name("x".to_string()) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .build() @@ -4498,7 +4504,8 @@ mod tests { .source_ids(vec![2]) .name("y") .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .build(), @@ -4519,7 +4526,8 @@ mod tests { .source_ids(vec![3]) .name("z") .transform(Transform::Identity) - .build(), + .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 009f6e78c9..9a8206cb27 100644 --- a/crates/iceberg/src/spec/table_metadata_builder.rs +++ b/crates/iceberg/src/spec/table_metadata_builder.rs @@ -1442,7 +1442,8 @@ mod tests { .source_ids(vec![2]) .name("y") .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .build() @@ -1684,6 +1685,7 @@ mod tests { .name("y".to_string()) .transform(Transform::Identity) .build() + .unwrap() ) .unwrap() .build() @@ -1758,13 +1760,15 @@ mod tests { .field_id(1000) .name("y".to_string()) .transform(Transform::Identity) - .build(), + .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(), + .build() + .unwrap(), ]) .unwrap() .build(); @@ -1785,7 +1789,8 @@ mod tests { .field_id(1000) .name("y".to_string()) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .add_unbound_field( @@ -1794,7 +1799,8 @@ mod tests { .field_id(1001) .name("z".to_string()) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .build() @@ -1838,7 +1844,8 @@ mod tests { .source_ids(vec![1]) .name("y_bucket[2]") .transform(Transform::Bucket(2)) - .build(), + .build() + .unwrap(), ) .unwrap() .build(); @@ -1859,7 +1866,8 @@ mod tests { .field_id(1001) .name("y_bucket[2]".to_string()) .transform(Transform::Bucket(2)) - .build(), + .build() + .unwrap(), ) .unwrap() .build() @@ -2430,13 +2438,15 @@ mod tests { .field_id(1000) .name("y".to_string()) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), UnboundPartitionField::builder() .source_ids(vec![3]) .field_id(1002) .name("z".to_string()) .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ]) .unwrap() .build(); @@ -2766,7 +2776,8 @@ mod tests { .source_ids(vec![1]) .name("bucket_data") .transform(Transform::Bucket(16)) - .build(), + .build() + .unwrap(), ) .unwrap() .build(); @@ -2831,7 +2842,8 @@ mod tests { .source_ids(vec![1]) .name("partition_col") .transform(Transform::Bucket(16)) - .build(), + .build() + .unwrap(), ) .unwrap() .build(); @@ -2886,7 +2898,8 @@ mod tests { .source_ids(vec![1]) .name("data_bucket") .transform(Transform::Bucket(16)) - .build(), + .build() + .unwrap(), ) .unwrap() .build(); @@ -2915,7 +2928,8 @@ mod tests { .source_ids(vec![1]) .name("existing_field") .transform(Transform::Bucket(8)) - .build(), + .build() + .unwrap(), ) .unwrap() .build(); @@ -2951,7 +2965,8 @@ mod tests { .source_ids(vec![1]) .name("bucket_data") .transform(Transform::Bucket(16)) - .build(), + .build() + .unwrap(), ) .unwrap() .build(); @@ -3053,7 +3068,8 @@ mod tests { .source_ids(vec![2]) .name("partition_data") .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .build(); @@ -3108,7 +3124,8 @@ mod tests { .source_ids(vec![1]) .name("bucket_data") .transform(Transform::Bucket(16)) - .build(), + .build() + .unwrap(), ) .unwrap() .build(); @@ -3164,7 +3181,8 @@ mod tests { .source_ids(vec![1]) .name("data_bucket") .transform(Transform::Bucket(16)) - .build(), + .build() + .unwrap(), ) .unwrap() .build(); @@ -3194,7 +3212,8 @@ mod tests { .source_ids(vec![2]) .name("new_partition_field") .transform(Transform::Bucket(8)) - .build(), + .build() + .unwrap(), ) .unwrap() .build(); @@ -3601,7 +3620,8 @@ mod tests { .source_ids(vec![1]) .name("id") .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) .unwrap() .build(); @@ -3626,7 +3646,8 @@ mod tests { .source_ids(vec![2]) .name("data_bucket") .transform(Transform::Bucket(10)) - .build(), + .build() + .unwrap(), ) .unwrap() .build(); @@ -3642,7 +3663,8 @@ mod tests { .source_ids(vec![1]) .name("id") .transform(Transform::Identity) - .build(), + .build() + .unwrap(), ) // Should reuse 1000 .unwrap() .add_partition_field( @@ -3650,7 +3672,8 @@ mod tests { .source_ids(vec![2]) .name("data_bucket") .transform(Transform::Bucket(10)) - .build(), + .build() + .unwrap(), ) // Should reuse 1001 .unwrap() .add_partition_field( @@ -3658,7 +3681,8 @@ mod tests { .source_ids(vec![3]) .name("year") .transform(Transform::Year) - .build(), + .build() + .unwrap(), ) // Should get new 1002 .unwrap() .build();