Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 7 additions & 6 deletions crates/iceberg/public-api.txt
Original file line number Diff line number Diff line change
Expand Up @@ -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<Self, <__D as serde_core::de::Deserializer>::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<i32>
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<iceberg::spec::PartitionField> for iceberg::Result<iceberg::spec::PartitionField>
pub fn iceberg::Result<iceberg::spec::PartitionField>::from(field: iceberg::spec::PartitionField) -> Self
impl core::convert::From<iceberg::spec::PartitionField> 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
Expand Down
6 changes: 3 additions & 3 deletions crates/iceberg/src/arrow/partition_value_calculator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -73,15 +73,15 @@ impl PartitionValueCalculator {
let transform_functions: Vec<BoxedTransformFunction> = partition_spec
.fields()
.iter()
.map(|pf| create_transform_function(&pf.transform))
.map(|pf| create_transform_function(&pf.transform()))
.collect::<Result<Vec<_>>>()?;

// Extract source field IDs for projection
let source_field_ids: Vec<i32> = partition_spec
.fields()
.iter()
.map(|pf| pf.source_id)
.collect();
.map(|pf| pf.source_id())
.collect::<Result<_>>()?;

// Create projector for extracting source columns
let projector = RecordBatchProjector::from_iceberg_schema(
Expand Down
14 changes: 8 additions & 6 deletions crates/iceberg/src/arrow/record_batch_transformer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
};

Expand All @@ -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
),
));
}
Expand All @@ -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
),
));
}
Expand Down Expand Up @@ -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] {
Expand Down
5 changes: 3 additions & 2 deletions crates/iceberg/src/expr/visitors/inclusive_projection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<PartitionField> = 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())
}
}
Expand All @@ -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 {
Expand Down
5 changes: 3 additions & 2 deletions crates/iceberg/src/expr/visitors/strict_projection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<PartitionField> = 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())
}
}
Expand Down Expand Up @@ -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
Expand Down
38 changes: 22 additions & 16 deletions crates/iceberg/src/partitioning.rs
Original file line number Diff line number Diff line change
Expand Up @@ -62,35 +62,36 @@ 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()
));
}

if !active_field_ids.contains(&field_id) {
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
Expand All @@ -99,16 +100,16 @@ 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()
));
}

// Use the correct type for dropped partitions in v1 tables: if the
// 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);
}
Expand Down Expand Up @@ -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).
Expand All @@ -167,8 +168,13 @@ fn all_active_field_ids<'a>(
) -> HashSet<i32> {
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()
}

Expand Down Expand Up @@ -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()
Expand Down
3 changes: 2 additions & 1 deletion crates/iceberg/src/scan/cache.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading