From a67f9c40a1e3d2067c1f1366ce64832ead09f6c8 Mon Sep 17 00:00:00 2001 From: Noritaka Sekiyama Date: Wed, 30 Sep 2026 10:29:56 +0900 Subject: [PATCH] Make PartitionField fields private and carry source-ids The second of the per-struct refactors after UnboundPartitionField. The four fields are private behind accessors, source_ids replaces source_id, and the builder validates at build time. A bound field can now hold several source ids, so binding a multi-argument field succeeds instead of failing. Such a transform cannot be evaluated yet, so it binds and reads as Transform::Unknown, as the spec requires of v3 readers; transform compatibility is checked for every source id. Co-authored-by: Isaac --- crates/iceberg/public-api.txt | 13 +- .../src/arrow/partition_value_calculator.rs | 6 +- .../src/arrow/record_batch_transformer.rs | 14 +- .../src/expr/visitors/inclusive_projection.rs | 5 +- .../src/expr/visitors/strict_projection.rs | 5 +- crates/iceberg/src/partitioning.rs | 38 +- crates/iceberg/src/scan/cache.rs | 3 +- crates/iceberg/src/spec/partition.rs | 379 +++++++++++++++--- crates/iceberg/src/spec/table_metadata.rs | 29 +- .../src/spec/table_metadata_builder.rs | 11 +- .../src/writer/file_writer/parquet_writer.rs | 18 +- 11 files changed, 404 insertions(+), 117 deletions(-) diff --git a/crates/iceberg/public-api.txt b/crates/iceberg/public-api.txt index 58c7db7dd6..692c7f2ab2 100644 --- a/crates/iceberg/public-api.txt +++ b/crates/iceberg/public-api.txt @@ -2316,24 +2316,25 @@ pub fn iceberg::spec::NestedField::serialize<__S>(&self, __serializer: __S) -> c impl<'de> serde_core::de::Deserialize<'de> for iceberg::spec::NestedField pub fn iceberg::spec::NestedField::deserialize<__D>(__deserializer: __D) -> core::result::Result::Error> where __D: serde_core::de::Deserializer<'de> pub struct iceberg::spec::PartitionField -pub iceberg::spec::PartitionField::field_id: i32 -pub iceberg::spec::PartitionField::name: alloc::string::String -pub iceberg::spec::PartitionField::source_id: i32 -pub iceberg::spec::PartitionField::transform: iceberg::spec::Transform impl iceberg::spec::PartitionField +pub fn iceberg::spec::PartitionField::field_id(&self) -> i32 pub fn iceberg::spec::PartitionField::into_unbound(self) -> iceberg::spec::UnboundPartitionField +pub fn iceberg::spec::PartitionField::name(&self) -> &str +pub fn iceberg::spec::PartitionField::source_id(&self) -> iceberg::Result +pub fn iceberg::spec::PartitionField::source_ids(&self) -> &[i32] +pub fn iceberg::spec::PartitionField::transform(&self) -> iceberg::spec::Transform impl core::clone::Clone for iceberg::spec::PartitionField pub fn iceberg::spec::PartitionField::clone(&self) -> iceberg::spec::PartitionField impl core::cmp::Eq for iceberg::spec::PartitionField impl core::cmp::PartialEq for iceberg::spec::PartitionField pub fn iceberg::spec::PartitionField::eq(&self, other: &iceberg::spec::PartitionField) -> bool +impl core::convert::From for iceberg::Result +pub fn iceberg::Result::from(field: iceberg::spec::PartitionField) -> Self impl core::convert::From for iceberg::spec::UnboundPartitionField pub fn iceberg::spec::UnboundPartitionField::from(field: iceberg::spec::PartitionField) -> Self impl core::fmt::Debug for iceberg::spec::PartitionField pub fn iceberg::spec::PartitionField::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result impl core::marker::StructuralPartialEq for iceberg::spec::PartitionField -impl iceberg::spec::PartitionField -pub fn iceberg::spec::PartitionField::builder() -> PartitionFieldBuilder<((), (), (), ())> impl serde_core::ser::Serialize for iceberg::spec::PartitionField pub fn iceberg::spec::PartitionField::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::PartitionField diff --git a/crates/iceberg/src/arrow/partition_value_calculator.rs b/crates/iceberg/src/arrow/partition_value_calculator.rs index c8c917779a..6591b9b9e3 100644 --- a/crates/iceberg/src/arrow/partition_value_calculator.rs +++ b/crates/iceberg/src/arrow/partition_value_calculator.rs @@ -73,15 +73,15 @@ impl PartitionValueCalculator { let transform_functions: Vec = partition_spec .fields() .iter() - .map(|pf| create_transform_function(&pf.transform)) + .map(|pf| create_transform_function(&pf.transform())) .collect::>>()?; // Extract source field IDs for projection let source_field_ids: Vec = partition_spec .fields() .iter() - .map(|pf| pf.source_id) - .collect(); + .map(|pf| pf.source_id()) + .collect::>()?; // Create projector for extracting source columns let projector = RecordBatchProjector::from_iceberg_schema( diff --git a/crates/iceberg/src/arrow/record_batch_transformer.rs b/crates/iceberg/src/arrow/record_batch_transformer.rs index 19976a81c5..86fa5c1513 100644 --- a/crates/iceberg/src/arrow/record_batch_transformer.rs +++ b/crates/iceberg/src/arrow/record_batch_transformer.rs @@ -82,10 +82,12 @@ fn constants_map( for (pos, field) in partition_spec.fields().iter().enumerate() { // Only identity transforms should use constant values from partition metadata - if matches!(field.transform, Transform::Identity) { + if matches!(field.transform(), Transform::Identity) { + // An identity transform reads exactly one source column. + let source_id = field.source_id()?; // The source column may have been dropped from the schema after the spec was // created. It cannot be projected in that case, so no constant is needed. - let Some(iceberg_field) = schema.field_by_id(field.source_id) else { + let Some(iceberg_field) = schema.field_by_id(source_id) else { continue; }; @@ -97,7 +99,7 @@ fn constants_map( ErrorKind::Unexpected, format!( "Partition field {} has non-primitive type {:?}", - field.source_id, iceberg_field.field_type + source_id, iceberg_field.field_type ), )); } @@ -115,14 +117,14 @@ fn constants_map( Some(Literal::Primitive(value)) => { // Create a Datum from the primitive type and value let datum = Datum::new(prim_type.clone(), value.clone()); - constants.insert(field.source_id, datum); + constants.insert(source_id, datum); } Some(literal) => { return Err(Error::new( ErrorKind::Unexpected, format!( "Partition field {} has non-primitive value: {:?}", - field.source_id, literal + source_id, literal ), )); } @@ -1116,7 +1118,7 @@ pub(crate) fn build_partition_constant( // Find matching field in this file's partition spec by field_id let value = spec_fields .iter() - .position(|f| f.field_id == unified_field.id) + .position(|f| f.field_id() == unified_field.id) .and_then(|pos| { // Get the value from partition_data at this position match &partition_data[pos] { diff --git a/crates/iceberg/src/expr/visitors/inclusive_projection.rs b/crates/iceberg/src/expr/visitors/inclusive_projection.rs index 88508af844..aba22f35b2 100644 --- a/crates/iceberg/src/expr/visitors/inclusive_projection.rs +++ b/crates/iceberg/src/expr/visitors/inclusive_projection.rs @@ -41,7 +41,8 @@ impl InclusiveProjection { if let std::collections::hash_map::Entry::Vacant(e) = self.cached_parts.entry(field_id) { let mut parts: Vec = vec![]; for partition_spec_field in self.partition_spec.fields() { - if partition_spec_field.source_id == field_id { + // A multi-argument field has an unknown transform, which cannot project + if partition_spec_field.source_ids() == [field_id] { parts.push(partition_spec_field.clone()) } } @@ -68,7 +69,7 @@ impl InclusiveProjection { .iter() .try_fold(Predicate::AlwaysTrue, |res, part| { Ok( - if let Some(pred_for_part) = part.transform.project(&part.name, predicate)? { + if let Some(pred_for_part) = part.transform().project(part.name(), predicate)? { if res == Predicate::AlwaysTrue { pred_for_part } else { diff --git a/crates/iceberg/src/expr/visitors/strict_projection.rs b/crates/iceberg/src/expr/visitors/strict_projection.rs index ebc6212c76..3a8e48d75e 100644 --- a/crates/iceberg/src/expr/visitors/strict_projection.rs +++ b/crates/iceberg/src/expr/visitors/strict_projection.rs @@ -45,7 +45,8 @@ impl StrictProjection { if let std::collections::hash_map::Entry::Vacant(e) = self.cached_parts.entry(field_id) { let mut parts: Vec = vec![]; for partition_spec_field in self.partition_spec.fields() { - if partition_spec_field.source_id == field_id { + // A multi-argument field has an unknown transform, which cannot project + if partition_spec_field.source_ids() == [field_id] { parts.push(partition_spec_field.clone()) } } @@ -81,7 +82,7 @@ impl StrictProjection { // the day, but does match the original predicate. Ok( if let Some(pred_for_part) = - part.transform.strict_project(&part.name, predicate)? + part.transform().strict_project(part.name(), predicate)? { if res == Predicate::AlwaysFalse { pred_for_part diff --git a/crates/iceberg/src/partitioning.rs b/crates/iceberg/src/partitioning.rs index 6ac62c68c4..7da8d9be70 100644 --- a/crates/iceberg/src/partitioning.rs +++ b/crates/iceberg/src/partitioning.rs @@ -62,17 +62,17 @@ pub fn compute_unified_partition_type<'a>( for spec in &specs { for field in spec.fields() { - let field_id = field.field_id; + let field_id = field.field_id(); // Reject unknown transforms up front: we cannot determine their result type, // so we cannot build a partition column for them. This check must precede the // active_field_ids filter below, otherwise an unknown transform could be // silently skipped. - if matches!(field.transform, Transform::Unknown) { + if matches!(field.transform(), Transform::Unknown) { return Err(invalid_data!( "Partition field '{}' uses an unknown transform whose result type \ cannot be determined", - field.name + field.name() )); } @@ -80,17 +80,18 @@ pub fn compute_unified_partition_type<'a>( continue; } - let source_field = match schema.field_by_id(field.source_id) { + // Unknown transforms were rejected above, so the field reads one source column. + let source_field = match schema.field_by_id(field.source_id()?) { Some(f) => f, None => continue, }; match field_map.get(&field_id) { None => { - let res_type = field.transform.result_type(&source_field.field_type)?; + let res_type = field.transform().result_type(&source_field.field_type)?; field_map.insert(field_id, field); type_map.insert(field_id, res_type); - name_map.insert(field_id, field.name.clone()); + name_map.insert(field_id, field.name().to_string()); } Some(existing) => { // V1 tables do not guarantee field ids are unique across specs, so two @@ -99,8 +100,8 @@ pub fn compute_unified_partition_type<'a>( return Err(invalid_data!( "Conflicting partition fields for field id {field_id}: \ '{}' and '{}'", - field.name, - existing.name + field.name(), + existing.name() )); } @@ -108,7 +109,7 @@ pub fn compute_unified_partition_type<'a>( // newer spec voided the field but an older spec has a real transform, // keep the older spec's type. if is_void_transform(existing) && !is_void_transform(field) { - let res_type = field.transform.result_type(&source_field.field_type)?; + let res_type = field.transform().result_type(&source_field.field_type)?; field_map.insert(field_id, field); type_map.insert(field_id, res_type); } @@ -143,16 +144,16 @@ pub fn compute_unified_partition_type<'a>( } fn is_void_transform(field: &PartitionField) -> bool { - matches!(field.transform, Transform::Void) + matches!(field.transform(), Transform::Void) } /// Two partition fields with the same field id are compatible if they share the same /// source id and have compatible transforms. Matches Java's /// `Partitioning.equivalentIgnoringNames`. fn equivalent_ignoring_names(field: &PartitionField, other: &PartitionField) -> bool { - field.field_id == other.field_id - && field.source_id == other.source_id - && compatible_transforms(&field.transform, &other.transform) + field.field_id() == other.field_id() + && field.source_ids() == other.source_ids() + && compatible_transforms(&field.transform(), &other.transform()) } /// Transforms are compatible if they are equal, or if either is Void (a dropped field). @@ -167,8 +168,13 @@ fn all_active_field_ids<'a>( ) -> HashSet { partition_specs .flat_map(|spec| spec.fields().iter()) - .filter(|field| schema.field_by_id(field.source_id).is_some()) - .map(|field| field.field_id) + .filter(|field| { + field + .source_ids() + .iter() + .all(|source_id| schema.field_by_id(*source_id).is_some()) + }) + .map(|field| field.field_id()) .collect() } @@ -395,7 +401,7 @@ mod tests { .add_unbound_field( UnboundPartitionField::builder() .source_ids(vec![4]) - .field_id(spec_v0.fields()[0].field_id) + .field_id(spec_v0.fields()[0].field_id()) .name("category".to_string()) .transform(Transform::Identity) .build() diff --git a/crates/iceberg/src/scan/cache.rs b/crates/iceberg/src/scan/cache.rs index b87786e6cf..91d0d7b92f 100644 --- a/crates/iceberg/src/scan/cache.rs +++ b/crates/iceberg/src/scan/cache.rs @@ -83,7 +83,8 @@ impl PartitionFilterCache { let has_dropped_source_column = partition_spec .fields() .iter() - .any(|field| schema.field_by_id(field.source_id).is_none()); + .flat_map(|field| field.source_ids()) + .any(|source_id| schema.field_by_id(*source_id).is_none()); let partition_filter = if has_dropped_source_column { BoundPredicate::AlwaysTrue diff --git a/crates/iceberg/src/spec/partition.rs b/crates/iceberg/src/spec/partition.rs index 17aa3de968..408b981276 100644 --- a/crates/iceberg/src/spec/partition.rs +++ b/crates/iceberg/src/spec/partition.rs @@ -34,25 +34,105 @@ pub(crate) const UNPARTITIONED_LAST_ASSIGNED_ID: i32 = 999; pub(crate) const DEFAULT_PARTITION_SPEC_ID: i32 = 0; /// Partition fields capture the transform from table data to partition values. +/// +/// The fields are private so that an instance is always well formed: `source_ids` holds at +/// least one id, and a field reading several source columns always carries +/// [`Transform::Unknown`], because multi-argument transforms cannot be evaluated yet and the +/// spec requires v3 readers to read such tables while ignoring them. #[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone, TypedBuilder)] -#[serde(rename_all = "kebab-case")] +#[serde( + try_from = "self::_serde::PartitionFieldSerde", + into = "self::_serde::PartitionFieldSerde" +)] +#[builder( + builder_method(vis = "pub(crate)"), + builder_type(vis = "pub(crate)"), + build_method(vis = "pub(crate)", into = Result) +)] pub struct PartitionField { - /// 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. - pub field_id: i32, + field_id: i32, /// A partition name. - pub name: String, + #[builder(setter(into))] + name: String, /// A transform that is applied to the source column to produce a partition value. - pub transform: Transform, + transform: Transform, } impl PartitionField { + /// 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. + pub fn field_id(&self) -> i32 { + 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 + } + /// To unbound partition field pub fn into_unbound(self) -> UnboundPartitionField { self.into() } + + fn validate(&self) -> Result<()> { + if self.source_ids.is_empty() { + return Err(invalid_data!("Empty source-ids is not allowed")); + } + if self.source_ids.len() > 1 && self.transform != Transform::Unknown { + return Err(invalid_data!( + "Partition field '{}' reads several source columns, so its transform must be unknown", + self.name + )); + } + Ok(()) + } +} + +impl From for Result { + fn from(field: PartitionField) -> Self { + field.validate()?; + Ok(field) + } +} + +/// The transform a bound field carries: a multi-argument transform cannot be evaluated yet, so +/// it is read as [`Transform::Unknown`]. +fn bound_transform(source_ids: &[i32], transform: Transform) -> Transform { + if source_ids.len() > 1 { + Transform::Unknown + } else { + transform + } } /// Reference to [`PartitionSpec`]. @@ -109,7 +189,11 @@ impl PartitionSpec { pub fn partition_type(&self, schema: &Schema) -> Result { let mut struct_fields = Vec::with_capacity(self.fields.len()); for partition_field in &self.fields { - let res_type = match schema.field_by_id(partition_field.source_id) { + let source_field = match partition_field.source_ids.as_slice() { + [source_id] => schema.field_by_id(*source_id), + _ => None, + }; + let res_type = match source_field { Some(field) => partition_field.transform.result_type(&field.field_type)?, // Historical specs may reference dropped source columns. Retain every // field's position and any result type that is independent of its source. @@ -169,7 +253,7 @@ impl PartitionSpec { } for (this_field, other_field) in self.fields.iter().zip(other.fields.iter()) { - if this_field.source_id != other_field.source_id + if this_field.source_ids != other_field.source_ids || this_field.name != other_field.name || this_field.transform != other_field.transform { @@ -351,11 +435,77 @@ impl From for Result { mod _serde { use serde::{Deserialize, Serialize}; - use super::UnboundPartitionField; + use super::{PartitionField, UnboundPartitionField, bound_transform}; use crate::Error; use crate::error::invalid_data; use crate::spec::Transform; + /// Resolve the `source-id` / `source-ids` pair read from a field's JSON. + fn read_source_ids( + source_id: Option, + source_ids: Option>, + ) -> Result, Error> { + match (source_id, source_ids) { + (Some(source_id), None) => Ok(vec![source_id]), + (None, Some(source_ids)) => Ok(source_ids), + (Some(_), Some(_)) => Err(invalid_data!( + "source-id and source-ids are mutually exclusive" + )), + (None, None) => Err(invalid_data!( + "Either `source-id` or `source-ids` must be present" + )), + } + } + + /// Write `source-id` for a single-argument field and `source-ids` otherwise. + fn write_source_ids(source_ids: Vec) -> (Option, Option>) { + match source_ids.as_slice() { + [source_id] => (Some(*source_id), None), + _ => (None, Some(source_ids)), + } + } + + /// Same spelling rules as [`UnboundPartitionFieldSerde`], with a required `field-id`. + #[derive(Serialize, Deserialize)] + #[serde(rename_all = "kebab-case")] + pub(super) struct PartitionFieldSerde { + #[serde(default, skip_serializing_if = "Option::is_none")] + source_id: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + source_ids: Option>, + field_id: i32, + name: String, + transform: Transform, + } + + impl TryFrom for PartitionField { + type Error = Error; + + fn try_from(value: PartitionFieldSerde) -> Result { + let source_ids = read_source_ids(value.source_id, value.source_ids)?; + let transform = bound_transform(&source_ids, value.transform); + PartitionField::builder() + .source_ids(source_ids) + .field_id(value.field_id) + .name(value.name) + .transform(transform) + .build() + } + } + + impl From for PartitionFieldSerde { + fn from(value: PartitionField) -> Self { + let (source_id, source_ids) = write_source_ids(value.source_ids); + Self { + source_id, + source_ids, + field_id: value.field_id, + name: value.name, + transform: value.transform, + } + } + } + /// Per the spec a single-argument field carries `source-id` and a multi-argument field /// carries `source-ids`. Either spelling is read, but not both at once; the one that /// matches the field is written. @@ -376,21 +526,7 @@ mod _serde { type Error = Error; fn try_from(value: UnboundPartitionFieldSerde) -> Result { - let source_ids = match (value.source_id, value.source_ids) { - (Some(source_id), None) => vec![source_id], - (None, Some(source_ids)) => source_ids, - (Some(_), Some(_)) => { - return Err(invalid_data!( - "source-id and source-ids are mutually exclusive" - )); - } - (None, None) => { - return Err(invalid_data!( - "Either `source-id` or `source-ids` must be present" - )); - } - }; - + let source_ids = read_source_ids(value.source_id, value.source_ids)?; UnboundPartitionField::builder() .source_ids(source_ids) .field_id_opt(value.field_id) @@ -402,10 +538,7 @@ mod _serde { impl From for UnboundPartitionFieldSerde { fn from(value: UnboundPartitionField) -> Self { - let (source_id, source_ids) = match value.source_ids.as_slice() { - [source_id] => (Some(*source_id), None), - _ => (None, Some(value.source_ids)), - }; + let (source_id, source_ids) = write_source_ids(value.source_ids); Self { source_id, source_ids, @@ -478,7 +611,7 @@ fn has_sequential_ids(field_ids: impl Iterator) -> bool { impl From for UnboundPartitionField { fn from(field: PartitionField) -> Self { UnboundPartitionField { - source_ids: vec![field.source_id], + source_ids: field.source_ids, field_id: Some(field.field_id), name: field.name, transform: field.transform, @@ -695,11 +828,12 @@ impl PartitionSpecBuilder { last_assigned_field_id }; + let transform = bound_transform(&field.source_ids, field.transform); bound_fields.push(PartitionField { - source_id: field.source_id()?, + source_ids: field.source_ids, field_id: partition_field_id, name: field.name, - transform: field.transform, + transform, }) } @@ -741,12 +875,18 @@ 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") - })?; + let mut schema_fields = Vec::with_capacity(field.source_ids.len()); + for &source_id in &field.source_ids { + schema_fields.push(schema.field_by_id(source_id).ok_or_else(|| { + invalid_data!("Cannot find partition source field with id `{source_id}` in schema") + })?); + } + + // A multi-argument field is bound with an unknown transform, which is not evaluated, + // so beyond its source columns existing there is nothing to check. + let [schema_field] = schema_fields.as_slice() else { + return Ok(()); + }; if field.transform != Transform::Void { if !schema_field.field_type.is_primitive() { @@ -868,17 +1008,17 @@ mod tests { "#; let partition_spec: PartitionSpec = serde_json::from_str(spec).unwrap(); - assert_eq!(4, partition_spec.fields[0].source_id); + assert_eq!(4, partition_spec.fields[0].source_id().unwrap()); assert_eq!(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!(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!(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); @@ -1249,12 +1389,15 @@ mod tests { ] { let spec = PartitionSpec { spec_id: 0, - fields: vec![PartitionField { - source_id: 1, - field_id: 1000, - name: "partition".to_string(), - transform, - }], + fields: vec![ + PartitionField::builder() + .source_ids(vec![1]) + .field_id(1000) + .name("partition") + .transform(transform) + .build() + .unwrap(), + ], }; assert_eq!( spec.partition_type(&schema).unwrap(), @@ -1380,12 +1523,15 @@ mod tests { assert_eq!(spec, PartitionSpec { spec_id: 1, - fields: vec![PartitionField { - source_id: 1, - field_id: 1000, - name: "id_bucket[16]".to_string(), - transform: Transform::Bucket(16), - }], + fields: vec![ + PartitionField::builder() + .source_ids(vec![1]) + .field_id(1000) + .name("id_bucket[16]") + .transform(Transform::Bucket(16)) + .build() + .unwrap() + ], }); assert_eq!( spec.partition_type(&schema).unwrap(), @@ -2074,7 +2220,7 @@ mod tests { } #[test] - fn test_binding_a_multi_argument_field_fails_loudly() { + fn test_binding_a_multi_argument_field() { let schema = Schema::builder() .with_fields(vec![ NestedField::required(1, "a", Type::Primitive(PrimitiveType::Int)).into(), @@ -2088,15 +2234,136 @@ mod tests { ) .unwrap(); - // A bound PartitionField still carries a single source id, so binding has to refuse - // rather than quietly keep the first one. + // Every id is kept, and the transform, which cannot be evaluated, binds as unknown + let bound = spec.bind(schema.clone()).unwrap(); + let field = &bound.fields()[0]; + assert_eq!([1, 2], field.source_ids()); + assert_eq!(Transform::Unknown, field.transform()); + assert_eq!( + bound.partition_type(&schema).unwrap(), + StructType::new(vec![ + NestedField::optional(1000, "m", Type::Primitive(PrimitiveType::String)).into() + ]) + ); + + // and the ids survive the way back to an unbound field + assert_eq!([1, 2], bound.into_unbound().fields()[0].source_ids()); + + // A source column that is not in the schema is still rejected + let spec: UnboundPartitionSpec = serde_json::from_str( + r#"{"spec-id": 1, "fields": [{"source-ids": [1, 3], "field-id": 1000, "name": "m", "transform": "bucket[4]"}]}"#, + ) + .unwrap(); let err = spec.bind(schema).unwrap_err(); assert!( - err.to_string().contains("has no single source id"), + err.to_string() + .contains("Cannot find partition source field with id `3`"), + "unexpected error: {err}" + ); + } + + #[test] + fn test_partition_field_reads_source_ids() { + let field: PartitionField = serde_json::from_str( + r#"{"source-ids": [1, 2], "field-id": 1000, "name": "m", "transform": "bucket[4]"}"#, + ) + .unwrap(); + assert_eq!([1, 2], field.source_ids()); + assert_eq!(Transform::Unknown, field.transform()); + assert!( + field + .source_id() + .unwrap_err() + .to_string() + .contains("has no single source id") + ); + + // A multi-argument field writes source-ids only + 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()); + + // A single-argument field keeps writing source-id, including one read from source-ids + for input in [ + r#"{"source-id": 1, "field-id": 1000, "name": "m", "transform": "bucket[4]"}"#, + r#"{"source-ids": [1], "field-id": 1000, "name": "m", "transform": "bucket[4]"}"#, + ] { + let field: PartitionField = serde_json::from_str(input).unwrap(); + assert_eq!(1, field.source_id().unwrap()); + assert_eq!(Transform::Bucket(4), field.transform()); + 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_partition_field_rejects_malformed_source_ids() { + for (input, expected) in [ + ( + r#"{"source-ids": [], "field-id": 1000, "name": "m", "transform": "identity"}"#, + "Empty source-ids is not allowed", + ), + ( + r#"{"field-id": 1000, "name": "m", "transform": "identity"}"#, + "Either `source-id` or `source-ids` must be present", + ), + ( + r#"{"source-id": 1, "source-ids": [1], "field-id": 1000, "name": "m", "transform": "identity"}"#, + "mutually exclusive", + ), + ] { + let err = serde_json::from_str::(input).unwrap_err(); + assert!( + err.to_string().contains(expected), + "unexpected error for {input}: {err}" + ); + } + } + + #[test] + fn test_partition_field_builder_validates() { + let err = PartitionField::builder() + .source_ids(vec![]) + .field_id(1000) + .name("m") + .transform(Transform::Identity) + .build() + .unwrap_err(); + assert!( + err.to_string().contains("Empty source-ids is not allowed"), + "unexpected error: {err}" + ); + + let err = PartitionField::builder() + .source_ids(vec![1, 2]) + .field_id(1000) + .name("m") + .transform(Transform::Bucket(4)) + .build() + .unwrap_err(); + assert!( + err.to_string().contains("its transform must be unknown"), "unexpected error: {err}" ); } + #[test] + fn test_is_compatible_with_compares_every_source_id() { + let spec = |source_ids: &str| -> PartitionSpec { + serde_json::from_str(&format!( + r#"{{"spec-id": 1, "fields": [{{"source-ids": {source_ids}, "field-id": 1000, "name": "m", "transform": "bucket[4]"}}]}}"# + )) + .unwrap() + }; + + assert!(spec("[1, 2]").is_compatible_with(&spec("[1, 2]"))); + assert!(!spec("[1, 2]").is_compatible_with(&spec("[1, 3]"))); + } + #[test] fn test_unbound_partition_field_builder_rejects_empty_source_ids() { let err = UnboundPartitionField::builder() diff --git a/crates/iceberg/src/spec/table_metadata.rs b/crates/iceberg/src/spec/table_metadata.rs index 83468602cd..baa9b195c2 100644 --- a/crates/iceberg/src/spec/table_metadata.rs +++ b/crates/iceberg/src/spec/table_metadata.rs @@ -158,7 +158,7 @@ impl TableMetadata { pub(crate) fn partition_name_exists(&self, name: &str) -> bool { self.partition_specs .values() - .any(|spec| spec.fields().iter().any(|pf| pf.name == name)) + .any(|spec| spec.fields().iter().any(|pf| pf.name() == name)) } /// Check if a field name exists in any schema. @@ -563,15 +563,18 @@ impl TableMetadata { for field in self.default_spec.fields() { // Historical specs may reference dropped columns, but active default-spec // fields need a source for new writes. Void fields do not read their source. - if field.transform != Transform::Void - && self.current_schema().field_by_id(field.source_id).is_none() - { - return Err(invalid_data!( - "Default partition spec {} references missing source field {} in current schema {}", - self.default_spec.spec_id(), - field.source_id, - self.current_schema_id - )); + if field.transform() == Transform::Void { + continue; + } + for &source_id in field.source_ids() { + if self.current_schema().field_by_id(source_id).is_none() { + return Err(invalid_data!( + "Default partition spec {} references missing source field {} in current schema {}", + self.default_spec.spec_id(), + source_id, + self.current_schema_id + )); + } } } @@ -3421,9 +3424,9 @@ mod tests { let default_spec = &desered_type.default_spec; assert_eq!(default_spec.spec_id(), 2); assert_eq!(default_spec.fields().len(), 1); - assert_eq!(default_spec.fields()[0].name, "y"); - assert_eq!(default_spec.fields()[0].transform, Transform::Identity); - assert_eq!(default_spec.fields()[0].source_id, 2); + assert_eq!(default_spec.fields()[0].name(), "y"); + assert_eq!(default_spec.fields()[0].transform(), Transform::Identity); + assert_eq!(default_spec.fields()[0].source_id().unwrap(), 2); } #[test] diff --git a/crates/iceberg/src/spec/table_metadata_builder.rs b/crates/iceberg/src/spec/table_metadata_builder.rs index 9a8206cb27..7d25e4b5b7 100644 --- a/crates/iceberg/src/spec/table_metadata_builder.rs +++ b/crates/iceberg/src/spec/table_metadata_builder.rs @@ -850,7 +850,12 @@ impl TableMetadataBuilder { .partition_specs .values() .flat_map(|spec| spec.fields()) - .map(|field| ((vec![field.source_id], field.transform), field.field_id)) + .map(|field| { + ( + (field.source_ids().to_vec(), field.transform()), + field.field_id(), + ) + }) .collect(); // Create new fields with reused field IDs where possible @@ -2799,7 +2804,7 @@ mod tests { .default_partition_spec() .fields() .iter() - .map(|f| f.name.clone()) + .map(|f| f.name().to_string()) .collect(); assert!(partition_field_names.contains(&"bucket_data".to_string())); @@ -3691,7 +3696,7 @@ mod tests { // Verify field ID reuse: spec 2 should reuse IDs from specs 0 and 1, assign new ID for new field let spec2 = result.metadata.partition_spec_by_id(2).unwrap(); - let field_ids: Vec = spec2.fields().iter().map(|f| f.field_id).collect(); + let field_ids: Vec = spec2.fields().iter().map(|f| f.field_id()).collect(); assert_eq!(field_ids, vec![1000, 1001, 1002]); // Reused 1000, 1001; new 1002 } } diff --git a/crates/iceberg/src/writer/file_writer/parquet_writer.rs b/crates/iceberg/src/writer/file_writer/parquet_writer.rs index 7cdc488889..72e32642b7 100644 --- a/crates/iceberg/src/writer/file_writer/parquet_writer.rs +++ b/crates/iceberg/src/writer/file_writer/parquet_writer.rs @@ -534,28 +534,28 @@ impl ParquetWriter { let mut partition_literals: Vec> = Vec::new(); for field in table_spec.fields() { - if let (Some(lower), Some(upper)) = ( - lower_bounds.get(&field.source_id), - upper_bounds.get(&field.source_id), - ) { - if !field.transform.preserves_order() { + let source_id = field.source_id()?; + if let (Some(lower), Some(upper)) = + (lower_bounds.get(&source_id), upper_bounds.get(&source_id)) + { + if !field.transform().preserves_order() { return Err(invalid_data!( "cannot infer partition value for non linear partition field (needs to preserve order): {} with transform {}", - field.name, - field.transform + field.name(), + field.transform() )); } if lower != upper { return Err(invalid_data!( "multiple partition values for field {}: lower: {:?}, upper: {:?}", - field.name, + field.name(), lower, upper )); } - let transform_fn = create_transform_function(&field.transform)?; + let transform_fn = create_transform_function(&field.transform())?; let transform_literal = Literal::from(transform_fn.transform_literal_result(lower)?);