This is an automated email from the ASF dual-hosted git repository.
CTTY pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/iceberg-rust.git
The following commit(s) were added to refs/heads/main by this push:
new 3d84c8135 feat: add support for _partition metadata column (#2668)
3d84c8135 is described below
commit 3d84c81353b1b23b6e4ae8eea8f8a021cc6927a7
Author: Parth Chandra <[email protected]>
AuthorDate: Wed Jul 29 15:28:39 2026 -0700
feat: add support for _partition metadata column (#2668)
## Which issue does this PR close?
- part of #https://github.com/apache/iceberg-rust/issues/2607.
## What changes are included in this PR?
Implements the _partition metadata column for table scans. This is a
struct column whose type is the union of all partition fields across all
partition specs (handling partition evolution). Each row gets the
partition values for its data file.
- Adds compute_unified_partition_type() to compute the union of
partition fields across all specs (equivalent to Java's
Partitioning.partitionType())
- Adds PartitionColumnConstant and build_partition_column_constant() for
pre-computing the struct values per file
- Adds ColumnSource::AddStructConstant variant to RecordBatchTransformer
for materializing struct columns
- Threads the unified partition type through scan planning and populates
the constant in FileScanTask
- Pipeline detects RESERVED_FIELD_ID_PARTITION in projected fields and
injects the struct constant
## Are these changes tested?
Because we do not have write support yet, I made the corresponding
change to comet and then tested by adding tests in Comet which uses
iceberg-java to write files and then iceberg-rust to read them back.
https://github.com/parthchandra/datafusion-comet/blob/iceberg-metadata-columns/spark/src/test/resources/sql-tests/iceberg/metadata_column_partition.sql
---
crates/iceberg/public-api.txt | 5 +-
crates/iceberg/src/arrow/mod.rs | 1 +
crates/iceberg/src/arrow/reader/pipeline.rs | 36 +-
crates/iceberg/src/arrow/reader/row_filter.rs | 2 +
.../iceberg/src/arrow/record_batch_transformer.rs | 671 ++++++++++++++++-----
crates/iceberg/src/arrow/value.rs | 74 +--
crates/iceberg/src/lib.rs | 1 +
crates/iceberg/src/partitioning.rs | 449 ++++++++++++++
crates/iceberg/src/scan/context.rs | 10 +-
crates/iceberg/src/scan/mod.rs | 20 +-
crates/iceberg/src/scan/task.rs | 18 +-
11 files changed, 1093 insertions(+), 194 deletions(-)
diff --git a/crates/iceberg/public-api.txt b/crates/iceberg/public-api.txt
index 610675c83..8a6227729 100644
--- a/crates/iceberg/public-api.txt
+++ b/crates/iceberg/public-api.txt
@@ -1158,6 +1158,8 @@ pub fn
iceberg::metadata_columns::partition_field(partition_fields: alloc::vec::
pub fn iceberg::metadata_columns::pos_field() -> &'static
iceberg::spec::NestedFieldRef
pub fn iceberg::metadata_columns::row_id_field() -> &'static
iceberg::spec::NestedFieldRef
pub fn iceberg::metadata_columns::spec_id_field() -> &'static
iceberg::spec::NestedFieldRef
+pub mod iceberg::partitioning
+pub fn
iceberg::partitioning::compute_unified_partition_type<'a>(partition_specs: impl
core::iter::traits::iterator::Iterator<Item = &'a
iceberg::spec::PartitionSpec>, schema: &iceberg::spec::Schema) ->
iceberg::Result<iceberg::spec::StructType>
pub mod iceberg::puffin
pub enum iceberg::puffin::CompressionCodec
pub iceberg::puffin::CompressionCodec::Gzip(u8)
@@ -1274,6 +1276,7 @@ pub iceberg::scan::FileScanTask::project_field_ids:
alloc::vec::Vec<i32>
pub iceberg::scan::FileScanTask::record_count: core::option::Option<u64>
pub iceberg::scan::FileScanTask::schema: iceberg::spec::SchemaRef
pub iceberg::scan::FileScanTask::start: u64
+pub iceberg::scan::FileScanTask::unified_partition_type:
core::option::Option<alloc::sync::Arc<iceberg::spec::StructType>>
impl iceberg::scan::FileScanTask
pub fn iceberg::scan::FileScanTask::data_file_path(&self) -> &str
pub fn iceberg::scan::FileScanTask::predicate(&self) ->
core::option::Option<&iceberg::expr::BoundPredicate>
@@ -1288,7 +1291,7 @@ impl core::fmt::Debug for iceberg::scan::FileScanTask
pub fn iceberg::scan::FileScanTask::fmt(&self, f: &mut
core::fmt::Formatter<'_>) -> core::fmt::Result
impl core::marker::StructuralPartialEq for iceberg::scan::FileScanTask
impl iceberg::scan::FileScanTask
-pub fn iceberg::scan::FileScanTask::builder() -> FileScanTaskBuilder<((), (),
(), (), (), (), (), (), (), (), (), (), (), (), ())>
+pub fn iceberg::scan::FileScanTask::builder() -> FileScanTaskBuilder<((), (),
(), (), (), (), (), (), (), (), (), (), (), (), (), ())>
impl serde_core::ser::Serialize for iceberg::scan::FileScanTask
pub fn iceberg::scan::FileScanTask::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::scan::FileScanTask
diff --git a/crates/iceberg/src/arrow/mod.rs b/crates/iceberg/src/arrow/mod.rs
index 4e0286cbd..28cf6e94a 100644
--- a/crates/iceberg/src/arrow/mod.rs
+++ b/crates/iceberg/src/arrow/mod.rs
@@ -32,6 +32,7 @@ mod reader;
/// RecordBatch projection utilities
pub mod record_batch_projector;
pub(crate) mod record_batch_transformer;
+pub(crate) use record_batch_transformer::build_partition_constant;
mod scan_metrics;
mod value;
diff --git a/crates/iceberg/src/arrow/reader/pipeline.rs
b/crates/iceberg/src/arrow/reader/pipeline.rs
index 579f04726..e7bcdf5fb 100644
--- a/crates/iceberg/src/arrow/reader/pipeline.rs
+++ b/crates/iceberg/src/arrow/reader/pipeline.rs
@@ -34,6 +34,7 @@ use super::{
ArrowFileReader, ArrowReader, ParquetReadOptions,
add_fallback_field_ids_to_arrow_schema,
apply_name_mapping_to_arrow_schema,
};
+use crate::arrow::build_partition_constant;
use crate::arrow::caching_delete_file_loader::CachingDeleteFileLoader;
use crate::arrow::int96::coerce_int96_timestamps;
use crate::arrow::record_batch_transformer::RecordBatchTransformerBuilder;
@@ -42,11 +43,11 @@ use crate::encryption::StandardKeyMetadata;
use crate::error::Result;
use crate::io::{FileIO, FileMetadata, FileRead};
use crate::metadata_columns::{
- RESERVED_COL_NAME_POS, RESERVED_FIELD_ID_FILE, RESERVED_FIELD_ID_POS,
- RESERVED_FIELD_ID_SPEC_ID, is_metadata_field,
+ RESERVED_COL_NAME_POS, RESERVED_FIELD_ID_FILE, RESERVED_FIELD_ID_PARTITION,
+ RESERVED_FIELD_ID_POS, RESERVED_FIELD_ID_SPEC_ID, is_metadata_field,
};
use crate::scan::{ArrowRecordBatchStream, FileScanTask, FileScanTaskStream};
-use crate::spec::Datum;
+use crate::spec::{Datum, PartitionSpec, Struct};
use crate::{Error, ErrorKind};
impl ArrowReader {
@@ -311,6 +312,35 @@ impl FileScanTaskReader {
record_batch_transformer_builder.with_virtual_field(RESERVED_FIELD_ID_POS);
}
+ // Add the _partition metadata struct column if it's in the projected
fields.
+ // Computed lazily here at read time from the unified partition type +
task's spec + data.
+ if task
+ .project_field_ids()
+ .contains(&RESERVED_FIELD_ID_PARTITION)
+ && let Some(unified_type) = &task.unified_partition_type
+ {
+ let (spec, partition_data) = match (&task.partition_spec,
&task.partition) {
+ (Some(spec), Some(data)) => (spec.clone(), data.clone()),
+ // A missing spec/data is only acceptable when there are no
partition
+ // fields to fill (unpartitioned table). If the unified type
has fields
+ // but we lack a spec or data, the task is inconsistent and we
cannot
+ // build the _partition column.
+ _ if unified_type.fields().is_empty() => {
+ (Arc::new(PartitionSpec::unpartition_spec()),
Struct::empty())
+ }
+ _ => {
+ return Err(Error::new(
+ ErrorKind::Unexpected,
+ "cannot build _partition column: unified partition
type has fields \
+ but the scan task is missing its partition spec or
data",
+ ));
+ }
+ };
+ let constant = build_partition_constant(unified_type, &spec,
&partition_data)?;
+ record_batch_transformer_builder =
+
record_batch_transformer_builder.with_partition_constant(constant);
+ }
+
let mut record_batch_transformer =
record_batch_transformer_builder.build();
if let Some(batch_size) = self.batch_size {
diff --git a/crates/iceberg/src/arrow/reader/row_filter.rs
b/crates/iceberg/src/arrow/reader/row_filter.rs
index a94c159d1..d82b9a0af 100644
--- a/crates/iceberg/src/arrow/reader/row_filter.rs
+++ b/crates/iceberg/src/arrow/reader/row_filter.rs
@@ -1147,6 +1147,7 @@ mod tests {
partition: None,
partition_spec: None,
name_mapping: None,
+ unified_partition_type: None,
case_sensitive: false,
key_metadata: None,
};
@@ -1246,6 +1247,7 @@ mod tests {
partition: None,
partition_spec: None,
name_mapping: None,
+ unified_partition_type: None,
case_sensitive: false,
key_metadata: None,
};
diff --git a/crates/iceberg/src/arrow/record_batch_transformer.rs
b/crates/iceberg/src/arrow/record_batch_transformer.rs
index ec91eb7b3..bf7e89910 100644
--- a/crates/iceberg/src/arrow/record_batch_transformer.rs
+++ b/crates/iceberg/src/arrow/record_batch_transformer.rs
@@ -20,21 +20,39 @@ use std::sync::Arc;
use arrow_array::{
Array as ArrowArray, ArrayRef, Int32Array, RecordBatch,
RecordBatchOptions, RunArray,
+ StructArray,
};
use arrow_cast::cast;
use arrow_schema::{
- DataType, Field, FieldRef, Schema as ArrowSchema, SchemaRef as
ArrowSchemaRef, SchemaRef,
+ DataType, Field, FieldRef, Fields, Schema as ArrowSchema, SchemaRef as
ArrowSchemaRef,
+ SchemaRef,
};
use parquet::arrow::PARQUET_FIELD_ID_META_KEY;
use crate::arrow::value::{create_primitive_array_repeated,
create_primitive_array_single_element};
use crate::arrow::{datum_to_arrow_type_with_ree, schema_to_arrow_schema,
type_to_arrow_type};
-use crate::metadata_columns::get_metadata_field;
+use crate::metadata_columns::{
+ RESERVED_COL_NAME_PARTITION, RESERVED_FIELD_ID_PARTITION,
get_metadata_field,
+};
use crate::spec::{
- Datum, Literal, PartitionSpec, PrimitiveLiteral, Schema as IcebergSchema,
Struct, Transform,
+ Datum, Literal, PartitionSpec, PrimitiveLiteral, Schema as IcebergSchema,
Struct, StructType,
+ Transform,
};
use crate::{Error, ErrorKind, Result};
+/// Create an Arrow Field with PARQUET_FIELD_ID_META_KEY metadata attached.
+fn field_with_id(
+ name: impl Into<String>,
+ data_type: DataType,
+ nullable: bool,
+ field_id: i32,
+) -> Field {
+ Field::new(name, data_type, nullable).with_metadata(HashMap::from([(
+ PARQUET_FIELD_ID_META_KEY.to_string(),
+ field_id.to_string(),
+ )]))
+}
+
/// Build a map of field ID to constant value (as Datum) for
identity-partitioned fields.
///
/// Implements Iceberg spec "Column Projection" rule #1: use partition
metadata constants
@@ -140,6 +158,13 @@ pub(crate) enum ColumnSource {
target_type: DataType,
value: Option<PrimitiveLiteral>,
},
+
+ // A struct column where each child is a constant primitive value.
+ // Used for the _partition metadata column.
+ AddStructConstant {
+ fields: Fields,
+ child_values: Vec<Option<PrimitiveLiteral>>,
+ },
// The iceberg spec refers to other permissible schema evolution actions
// (see https://iceberg.apache.org/spec/#schema-evolution):
// renaming fields, deleting fields and reordering fields.
@@ -186,16 +211,68 @@ enum SchemaComparison {
/// Builder for RecordBatchTransformer to improve ergonomics when constructing
with optional parameters.
///
-/// Constant fields are pre-computed for both virtual/metadata fields (like
_file) and
-/// identity-partitioned fields to avoid duplicate work during batch
processing.
+/// All per-file constants (scalar metadata like `_file`, identity partition
values,
+/// and the `_partition` struct) are stored in a single `constant_fields` map
keyed by
+/// field_id. This unified representation (via [`ColumnConstant`]) means the
+/// transformer handles all constant columns through one code path.
#[derive(Debug)]
pub(crate) struct RecordBatchTransformerBuilder {
snapshot_schema: Arc<IcebergSchema>,
projected_iceberg_field_ids: Vec<i32>,
- constant_fields: HashMap<i32, Datum>,
+ constant_fields: HashMap<i32, ColumnConstant>,
virtual_fields: HashSet<i32>,
}
+/// A per-file constant value for a column.
+///
+/// Covers both scalar constants (metadata columns like `_file` and `_spec_id`,
+/// as well as identity partition source fields) and the struct `_partition`
column.
+#[derive(Debug, Clone, PartialEq)]
+pub(crate) enum ColumnConstant {
+ /// A scalar constant (e.g., `_file` path, `_spec_id`, identity partition
values).
+ /// The Datum carries both the Iceberg type and the value.
+ Scalar(Datum),
+ /// A struct constant (the `_partition` column). Each child is a primitive
constant
+ /// or null (for partition evolution gaps).
+ Struct(StructConstant),
+}
+
+/// Pre-computed data for a struct constant column.
+#[derive(Debug, Clone, PartialEq)]
+pub(crate) struct StructConstant {
+ fields: Fields,
+ child_values: Vec<Option<PrimitiveLiteral>>,
+}
+
+impl StructConstant {
+ /// Create a new StructConstant, validating that fields and child_values
have
+ /// the same length.
+ pub(crate) fn new(fields: Fields, child_values:
Vec<Option<PrimitiveLiteral>>) -> Result<Self> {
+ if fields.len() != child_values.len() {
+ return Err(Error::new(
+ ErrorKind::DataInvalid,
+ format!(
+ "StructConstant: fields length ({}) != child_values length
({})",
+ fields.len(),
+ child_values.len()
+ ),
+ ));
+ }
+ Ok(Self {
+ fields,
+ child_values,
+ })
+ }
+
+ pub(crate) fn fields(&self) -> &Fields {
+ &self.fields
+ }
+
+ pub(crate) fn child_values(&self) -> &[Option<PrimitiveLiteral>] {
+ &self.child_values
+ }
+}
+
impl RecordBatchTransformerBuilder {
pub(crate) fn new(
snapshot_schema: Arc<IcebergSchema>,
@@ -209,14 +286,11 @@ impl RecordBatchTransformerBuilder {
}
}
- /// Add a constant value for a specific field ID.
+ /// Add a scalar constant value for a specific field ID.
/// This is used for virtual/metadata fields like _file that have constant
values per batch.
- ///
- /// # Arguments
- /// * `field_id` - The field ID to associate with the constant
- /// * `datum` - The constant value (with type) for this field
pub(crate) fn with_constant(mut self, field_id: i32, datum: Datum) -> Self
{
- self.constant_fields.insert(field_id, datum);
+ self.constant_fields
+ .insert(field_id, ColumnConstant::Scalar(datum));
self
}
@@ -230,13 +304,12 @@ impl RecordBatchTransformerBuilder {
partition_spec: Arc<PartitionSpec>,
partition_data: Struct,
) -> Result<Self> {
- // Compute partition constants for identity-transformed fields
(already returns Datum)
let partition_constants =
constants_map(&partition_spec, &partition_data,
&self.snapshot_schema)?;
- // Add partition constants to constant_fields
for (field_id, datum) in partition_constants {
- self.constant_fields.insert(field_id, datum);
+ self.constant_fields
+ .insert(field_id, ColumnConstant::Scalar(datum));
}
Ok(self)
@@ -249,6 +322,15 @@ impl RecordBatchTransformerBuilder {
self
}
+ /// Set a pre-computed _partition column constant directly.
+ pub(crate) fn with_partition_constant(mut self, partition_column:
StructConstant) -> Self {
+ self.constant_fields.insert(
+ RESERVED_FIELD_ID_PARTITION,
+ ColumnConstant::Struct(partition_column),
+ );
+ self
+ }
+
pub(crate) fn build(self) -> RecordBatchTransformer {
RecordBatchTransformer {
snapshot_schema: self.snapshot_schema,
@@ -294,10 +376,8 @@ impl RecordBatchTransformerBuilder {
pub(crate) struct RecordBatchTransformer {
snapshot_schema: Arc<IcebergSchema>,
projected_iceberg_field_ids: Vec<i32>,
- // Pre-computed constant field information: field_id -> Datum
- // Includes both virtual/metadata fields (like _file) and
identity-partitioned fields
- // Datum holds both the Iceberg type and the value
- constant_fields: HashMap<i32, Datum>,
+ // Unified map of all per-file constant columns (metadata, identity
partition, _partition struct).
+ constant_fields: HashMap<i32, ColumnConstant>,
// Field IDs whose data is delivered by the source batch as an arrow-rs
virtual column
// (e.g. _pos via RowNumber). These fields bypass the snapshot-schema
lookup and the
@@ -364,7 +444,7 @@ impl RecordBatchTransformer {
source_schema: &ArrowSchemaRef,
snapshot_schema: &IcebergSchema,
projected_iceberg_field_ids: &[i32],
- constant_fields: &HashMap<i32, Datum>,
+ constant_fields: &HashMap<i32, ColumnConstant>,
virtual_fields: &HashSet<i32>,
) -> Result<BatchTransform> {
let mapped_unprojected_arrow_schema =
Arc::new(schema_to_arrow_schema(snapshot_schema)?);
@@ -376,39 +456,66 @@ impl RecordBatchTransformer {
let fields: Result<Vec<_>> = projected_iceberg_field_ids
.iter()
.map(|field_id| {
- // Metadata/virtual fields (like _file, _spec_id) don't exist
in the table
- // schema, so build their Arrow field from the metadata column
definition and
- // the pre-computed constant's type.
- if let Some(datum) = constant_fields.get(field_id)
- && let Ok(iceberg_field) = get_metadata_field(*field_id)
- {
- let arrow_type = datum_to_arrow_type_with_ree(datum);
- let arrow_field =
- Field::new(&iceberg_field.name, arrow_type,
!iceberg_field.required)
- .with_metadata(HashMap::from([(
- PARQUET_FIELD_ID_META_KEY.to_string(),
- iceberg_field.id.to_string(),
- )]));
- Ok(Arc::new(arrow_field))
- } else if virtual_fields.contains(field_id) {
+ match constant_fields.get(field_id) {
+ Some(ColumnConstant::Struct(pc))
+ if *field_id == RESERVED_FIELD_ID_PARTITION =>
+ {
+ let struct_type =
DataType::Struct(pc.fields().clone());
+ let nullable = pc.fields().is_empty();
+ let arrow_field = field_with_id(
+ RESERVED_COL_NAME_PARTITION,
+ struct_type,
+ nullable,
+ *field_id,
+ );
+ return Ok(Arc::new(arrow_field));
+ }
+ Some(ColumnConstant::Struct(_)) => {
+ return Err(Error::new(
+ ErrorKind::Unexpected,
+ format!(
+ "Struct column constants are only supported
for the `_partition` \
+ metadata column (field id
{RESERVED_FIELD_ID_PARTITION}), but \
+ one was set for field id {field_id}"
+ ),
+ ));
+ }
+ Some(ColumnConstant::Scalar(datum)) => {
+ if let Ok(iceberg_field) =
get_metadata_field(*field_id) {
+ let arrow_type =
datum_to_arrow_type_with_ree(datum);
+ let arrow_field = field_with_id(
+ &iceberg_field.name,
+ arrow_type,
+ !iceberg_field.required,
+ iceberg_field.id,
+ );
+ return Ok(Arc::new(arrow_field));
+ }
+ // Identity partition constant -- fall through to use
the
+ // mapped schema field (read from file if present).
+ }
+ None => {}
+ }
+
+ if virtual_fields.contains(field_id) {
let virtual_field = get_metadata_field(*field_id)?;
let arrow_type =
type_to_arrow_type(&virtual_field.field_type)?;
- let arrow_field =
- Field::new(&virtual_field.name, arrow_type,
!virtual_field.required)
- .with_metadata(HashMap::from([(
- PARQUET_FIELD_ID_META_KEY.to_string(),
- virtual_field.id.to_string(),
- )]));
- Ok(Arc::new(arrow_field))
- } else {
- // Regular fields and identity-partitioned constant fields
both exist in the
- // table schema, so use the mapped Arrow field as-is.
- Ok(field_id_to_mapped_schema_map
- .get(field_id)
- .ok_or(Error::new(ErrorKind::Unexpected, "field not
found"))?
- .0
- .clone())
+ let arrow_field = field_with_id(
+ &virtual_field.name,
+ arrow_type,
+ !virtual_field.required,
+ virtual_field.id,
+ );
+ return Ok(Arc::new(arrow_field));
}
+
+ // Regular fields and identity-partitioned constant fields
both exist in the
+ // table schema, so use the mapped Arrow field as-is.
+ Ok(field_id_to_mapped_schema_map
+ .get(field_id)
+ .ok_or(Error::new(ErrorKind::Unexpected, "field not
found"))?
+ .0
+ .clone())
})
.collect();
@@ -492,7 +599,7 @@ impl RecordBatchTransformer {
snapshot_schema: &IcebergSchema,
projected_iceberg_field_ids: &[i32],
field_id_to_mapped_schema_map: HashMap<i32, (FieldRef, usize)>,
- constant_fields: &HashMap<i32, Datum>,
+ constant_fields: &HashMap<i32, ColumnConstant>,
virtual_fields: &HashSet<i32>,
) -> Result<Vec<ColumnSource>> {
let field_id_to_source_schema_map =
@@ -501,39 +608,62 @@ impl RecordBatchTransformer {
projected_iceberg_field_ids
.iter()
.map(|field_id| {
- // Check if this is a constant field (metadata/virtual or
identity-partitioned).
- //
- // Metadata/virtual fields (like _file, _spec_id) never exist
in the data file,
- // so they always use their pre-computed constant value.
- //
- // For identity-partitioned fields, the Iceberg spec's "Column
Projection" rules
- // only apply to "field ids which are not present in a data
file". When the column
- // IS present in the Parquet file, it must be read from the
file; the partition
- // metadata constant is only a fallback for when the column is
absent (e.g. add_files).
- if let Some(datum) = constant_fields.get(field_id) {
- let is_metadata_field =
get_metadata_field(*field_id).is_ok();
- let present_in_file =
field_id_to_source_schema_map.contains_key(field_id);
-
- if is_metadata_field || !present_in_file {
- let arrow_type = if is_metadata_field {
- datum_to_arrow_type_with_ree(datum)
- } else {
- field_id_to_mapped_schema_map
- .get(field_id)
- .ok_or(Error::new(
- ErrorKind::Unexpected,
- "could not find field in schema",
- ))?
- .0
- .data_type()
- .clone()
- };
-
- return Ok(ColumnSource::Add {
- value: Some(datum.literal().clone()),
- target_type: arrow_type,
+ match constant_fields.get(field_id) {
+ Some(ColumnConstant::Struct(pc))
+ if *field_id == RESERVED_FIELD_ID_PARTITION =>
+ {
+ return Ok(ColumnSource::AddStructConstant {
+ fields: pc.fields().clone(),
+ child_values: pc.child_values().to_vec(),
});
}
+ Some(ColumnConstant::Struct(_)) => {
+ return Err(Error::new(
+ ErrorKind::Unexpected,
+ format!(
+ "Struct column constants are only supported
for the `_partition` \
+ metadata column (field id
{RESERVED_FIELD_ID_PARTITION}), but \
+ one was set for field id {field_id}"
+ ),
+ ));
+ }
+ Some(ColumnConstant::Scalar(datum)) => {
+ // Metadata/virtual fields (like _file, _spec_id)
never exist in the
+ // data file, so they always use their pre-computed
constant value.
+ //
+ // For identity-partitioned fields, the Iceberg spec's
"Column
+ // Projection" rules only apply to "field ids which
are not present in
+ // a data file". When the column IS present in the
Parquet file, it
+ // must be read from the file; the partition metadata
constant is only
+ // a fallback for when the column is absent (e.g.
add_files).
+ let is_metadata =
get_metadata_field(*field_id).is_ok();
+ let present_in_file =
+
field_id_to_source_schema_map.contains_key(field_id);
+
+ if is_metadata || !present_in_file {
+ let arrow_type = if is_metadata {
+ datum_to_arrow_type_with_ree(datum)
+ } else {
+ field_id_to_mapped_schema_map
+ .get(field_id)
+ .ok_or(Error::new(
+ ErrorKind::Unexpected,
+ "could not find field in schema",
+ ))?
+ .0
+ .data_type()
+ .clone()
+ };
+
+ return Ok(ColumnSource::Add {
+ value: Some(datum.literal().clone()),
+ target_type: arrow_type,
+ });
+ }
+ // Identity partition field present in the file --
fall through
+ // to read from the file instead of using the constant.
+ }
+ None => {}
}
if virtual_fields.contains(field_id) {
@@ -682,6 +812,11 @@ impl RecordBatchTransformer {
ColumnSource::Add { target_type, value } => {
Self::create_column(target_type, value, num_rows)?
}
+
+ ColumnSource::AddStructConstant {
+ fields,
+ child_values,
+ } => Self::create_struct_column(fields, child_values,
num_rows)?,
})
})
.collect()
@@ -723,6 +858,89 @@ impl RecordBatchTransformer {
create_primitive_array_repeated(target_type, prim_lit, num_rows)
}
}
+
+ fn create_struct_column(
+ fields: &Fields,
+ child_values: &[Option<PrimitiveLiteral>],
+ num_rows: usize,
+ ) -> Result<ArrayRef> {
+ if fields.is_empty() {
+ let nulls = arrow_buffer::NullBuffer::new_null(num_rows);
+ return Ok(Arc::new(StructArray::new_empty_fields(
+ num_rows,
+ Some(nulls),
+ )));
+ }
+
+ let child_arrays: Vec<ArrayRef> = fields
+ .iter()
+ .zip(child_values.iter())
+ .map(|(field, value)| {
+ create_primitive_array_repeated(field.data_type(), value,
num_rows)
+ })
+ .collect::<Result<_>>()?;
+
+ Ok(Arc::new(StructArray::try_new(
+ fields.clone(),
+ child_arrays,
+ None,
+ )?))
+ }
+}
+
+/// Builds a [`StructConstant`] from the unified partition type and a file's
+/// partition spec/data.
+///
+/// If `unified_partition_type` has no fields (unpartitioned table), returns
an empty constant
+/// that renders as a null struct column.
+///
+/// For each field in the unified partition type:
+/// - If it corresponds to a field in this file's partition spec, use the
value from partition_data
+/// - Otherwise (partition evolution), use null
+pub(crate) fn build_partition_constant(
+ unified_partition_type: &StructType,
+ partition_spec: &PartitionSpec,
+ partition_data: &Struct,
+) -> Result<StructConstant> {
+ // Unpartitioned table: empty struct rendered as null
+ if unified_partition_type.fields().is_empty() {
+ return StructConstant::new(Fields::empty(), vec![]);
+ }
+
+ use crate::arrow::type_to_arrow_type;
+
+ let spec_fields = partition_spec.fields();
+
+ let mut arrow_fields =
Vec::with_capacity(unified_partition_type.fields().len());
+ let mut child_values =
Vec::with_capacity(unified_partition_type.fields().len());
+
+ for unified_field in unified_partition_type.fields() {
+ let arrow_type = type_to_arrow_type(&unified_field.field_type)?;
+ // Don't attach PARQUET_FIELD_ID_META_KEY metadata to child fields --
+ // Spark's output schema for _partition doesn't include it, and Arrow's
+ // Field::eq checks metadata, causing the schema adapter to insert an
+ // unnecessary cast operation per batch.
+ // Always nullable: partition evolution means any field may be absent
+ // for files written under a different spec.
+ let arrow_field = Field::new(&unified_field.name, arrow_type, true);
+ arrow_fields.push(Arc::new(arrow_field));
+
+ // 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)
+ .and_then(|pos| {
+ // Get the value from partition_data at this position
+ match &partition_data[pos] {
+ Some(Literal::Primitive(prim)) => Some(prim.clone()),
+ _ => None,
+ }
+ });
+
+ child_values.push(value);
+ }
+
+ StructConstant::new(Fields::from(arrow_fields), child_values)
}
#[cfg(test)]
@@ -735,8 +953,9 @@ mod test {
StringArray,
};
use arrow_schema::{DataType, Field, Schema as ArrowSchema};
- use parquet::arrow::PARQUET_FIELD_ID_META_KEY;
+ use super::field_with_id;
+ use crate::arrow::build_partition_constant;
use crate::arrow::record_batch_transformer::{
RecordBatchTransformer, RecordBatchTransformerBuilder,
};
@@ -866,8 +1085,8 @@ mod test {
.build();
let file_schema = Arc::new(ArrowSchema::new(vec![
- simple_field("id", DataType::Int32, false, "1"),
- simple_field("name", DataType::Utf8, true, "2"),
+ field_with_id("id", DataType::Int32, false, 1),
+ field_with_id("name", DataType::Utf8, true, 2),
]));
let file_batch = RecordBatch::try_new(file_schema, vec![
@@ -934,11 +1153,11 @@ mod test {
RecordBatchTransformerBuilder::new(snapshot_schema,
&projected_iceberg_field_ids)
.build();
- let file_schema = Arc::new(ArrowSchema::new(vec![simple_field(
+ let file_schema = Arc::new(ArrowSchema::new(vec![field_with_id(
"id",
DataType::Int32,
false,
- "1",
+ 1,
)]));
let file_batch =
RecordBatch::try_new(file_schema,
vec![Arc::new(Int32Array::from(vec![1, 2, 3]))])
@@ -988,8 +1207,8 @@ mod test {
.build();
let file_schema = Arc::new(ArrowSchema::new(vec![
- simple_field("id", DataType::Int32, false, "1"),
- simple_field("data", DataType::Utf8, false, "2"),
+ field_with_id("id", DataType::Int32, false, 1),
+ field_with_id("data", DataType::Utf8, false, 2),
]));
let file_batch = RecordBatch::try_new(file_schema, vec![
@@ -1109,38 +1328,30 @@ mod test {
fn arrow_schema_already_same_as_target() -> Arc<ArrowSchema> {
Arc::new(ArrowSchema::new(vec![
- simple_field("a", DataType::Utf8, true, "10"),
- simple_field("b", DataType::Int64, false, "11"),
- simple_field("c", DataType::Float64, false, "12"),
- simple_field("e", DataType::Utf8, true, "14"),
- simple_field("f", DataType::Utf8, false, "15"),
+ field_with_id("a", DataType::Utf8, true, 10),
+ field_with_id("b", DataType::Int64, false, 11),
+ field_with_id("c", DataType::Float64, false, 12),
+ field_with_id("e", DataType::Utf8, true, 14),
+ field_with_id("f", DataType::Utf8, false, 15),
]))
}
fn arrow_schema_promotion_addition_and_renaming_required() ->
Arc<ArrowSchema> {
Arc::new(ArrowSchema::new(vec![
- simple_field("b", DataType::Int32, false, "11"),
- simple_field("c", DataType::Float32, false, "12"),
- simple_field("d", DataType::Int32, false, "13"),
- simple_field("e_old", DataType::Utf8, true, "14"),
+ field_with_id("b", DataType::Int32, false, 11),
+ field_with_id("c", DataType::Float32, false, 12),
+ field_with_id("d", DataType::Int32, false, 13),
+ field_with_id("e_old", DataType::Utf8, true, 14),
]))
}
fn arrow_schema_no_promotion_addition_or_renaming_required() ->
Arc<ArrowSchema> {
Arc::new(ArrowSchema::new(vec![
- simple_field("d", DataType::Int32, false, "13"),
- simple_field("e", DataType::Utf8, true, "14"),
+ field_with_id("d", DataType::Int32, false, 13),
+ field_with_id("e", DataType::Utf8, true, 14),
]))
}
- /// Create a simple arrow field with metadata.
- fn simple_field(name: &str, ty: DataType, nullable: bool, value: &str) ->
Field {
- Field::new(name, ty, nullable).with_metadata(HashMap::from([(
- PARQUET_FIELD_ID_META_KEY.to_string(),
- value.to_string(),
- )]))
- }
-
/// Test for add_files with Parquet files that have NO field IDs (Hive
tables).
///
/// This reproduces the scenario from Iceberg spec where:
@@ -1320,8 +1531,8 @@ mod test {
// Parquet file contains both id and name columns
let parquet_schema = Arc::new(ArrowSchema::new(vec![
- simple_field("id", DataType::Int32, false, "1"),
- simple_field("name", DataType::Utf8, true, "2"),
+ field_with_id("id", DataType::Int32, false, 1),
+ field_with_id("name", DataType::Utf8, true, 2),
]));
let projected_field_ids = [1, 2]; // id, name
@@ -1441,8 +1652,8 @@ mod test {
// Parquet file contains only id and name (dept is in partition path)
let parquet_schema = Arc::new(ArrowSchema::new(vec![
- simple_field("id", DataType::Int32, false, "1"),
- simple_field("name", DataType::Utf8, true, "3"),
+ field_with_id("id", DataType::Int32, false, 1),
+ field_with_id("name", DataType::Utf8, true, 3),
]));
let projected_field_ids = [1, 2, 3]; // id, dept, name
@@ -1531,8 +1742,8 @@ mod test {
let partition_data =
Struct::from_iter(vec![Some(Literal::string("engineering"))]);
let parquet_schema = Arc::new(ArrowSchema::new(vec![
- simple_field("id", DataType::Int32, false, "1"),
- simple_field("name", DataType::Utf8, true, "3"),
+ field_with_id("id", DataType::Int32, false, 1),
+ field_with_id("name", DataType::Utf8, true, 3),
]));
let projected_field_ids = [1, 3]; // id, name -- dept is gone
@@ -1627,8 +1838,8 @@ mod test {
// Parquet file has OLD column name "id" but SAME field_id=1
// Field-ID-based mapping should find this despite name mismatch
let parquet_schema = Arc::new(ArrowSchema::new(vec![
- simple_field("id", DataType::Int32, false, "1"),
- simple_field("name", DataType::Utf8, true, "2"),
+ field_with_id("id", DataType::Int32, false, 1),
+ field_with_id("name", DataType::Utf8, true, 2),
]));
let projected_field_ids = [1, 2]; // row_id (field_id=1), name
(field_id=2)
@@ -1732,8 +1943,8 @@ mod test {
// Has id (field_id=1) and data (field_id=3, assigned by ArrowReader
via name mapping)
// Missing: dept (in partition), category (has default), notes (no
default)
let parquet_schema = Arc::new(ArrowSchema::new(vec![
- simple_field("id", DataType::Int32, false, "1"),
- simple_field("data", DataType::Utf8, false, "3"),
+ field_with_id("id", DataType::Int32, false, 1),
+ field_with_id("data", DataType::Utf8, false, 3),
]));
let projected_field_ids = [1, 2, 3, 4, 5]; // id, dept, data,
category, notes
@@ -1824,11 +2035,11 @@ mod test {
// Partition has null value for the data column
let partition_data = Struct::from_iter(vec![None]);
- let file_schema = Arc::new(ArrowSchema::new(vec![simple_field(
+ let file_schema = Arc::new(ArrowSchema::new(vec![field_with_id(
"id",
DataType::Int32,
true,
- "1",
+ 1,
)]));
let projected_field_ids = [1, 2];
@@ -1879,14 +2090,9 @@ mod test {
// Simulate what arrow-rs's virtual-column reader produces: file
columns
// followed by the _pos Int64 column carrying absolute file row
indices.
let parquet_schema = Arc::new(ArrowSchema::new(vec![
- simple_field("id", DataType::Int32, false, "1"),
- simple_field("name", DataType::Utf8, true, "2"),
- simple_field(
- "_pos",
- DataType::Int64,
- false,
- &RESERVED_FIELD_ID_POS.to_string(),
- ),
+ field_with_id("id", DataType::Int32, false, 1),
+ field_with_id("name", DataType::Utf8, true, 2),
+ field_with_id("_pos", DataType::Int64, false,
RESERVED_FIELD_ID_POS),
]));
let projected_field_ids = [1, 2, RESERVED_FIELD_ID_POS];
@@ -1954,14 +2160,9 @@ mod test {
);
let parquet_schema = Arc::new(ArrowSchema::new(vec![
- simple_field("id_long", DataType::Int64, false, "1"),
- simple_field("name", DataType::Utf8, true, "2"),
- simple_field(
- "_pos",
- DataType::Int64,
- false,
- &RESERVED_FIELD_ID_POS.to_string(),
- ),
+ field_with_id("id_long", DataType::Int64, false, 1),
+ field_with_id("name", DataType::Utf8, true, 2),
+ field_with_id("_pos", DataType::Int64, false,
RESERVED_FIELD_ID_POS),
]));
let projected_field_ids = [RESERVED_FIELD_ID_POS, 2, 1];
@@ -2008,4 +2209,202 @@ mod test {
assert_eq!(id_col.len(), 3);
assert_eq!(id_col.values(), &[100, 200, 300]);
}
+
+ #[test]
+ fn partition_column_struct_constant() {
+ use arrow_array::StructArray;
+
+ use crate::metadata_columns::RESERVED_FIELD_ID_PARTITION;
+ use crate::spec::Transform;
+
+ let snapshot_schema = Arc::new(
+ Schema::builder()
+ .with_schema_id(0)
+ .with_fields(vec![
+ NestedField::required(1, "id",
Type::Primitive(PrimitiveType::Int)).into(),
+ NestedField::optional(2, "name",
Type::Primitive(PrimitiveType::String)).into(),
+ ])
+ .build()
+ .unwrap(),
+ );
+
+ // Partition spec: identity(id)
+ let partition_spec = Arc::new(
+ crate::spec::PartitionSpec::builder(snapshot_schema.clone())
+ .with_spec_id(0)
+ .add_partition_field("id", "id", Transform::Identity)
+ .unwrap()
+ .build()
+ .unwrap(),
+ );
+
+ // Unified partition type: just one field (same as the spec's type)
+ let unified_partition_type =
partition_spec.partition_type(&snapshot_schema).unwrap();
+
+ // Partition data: id=42
+ let partition_data = Struct::from_iter(vec![Some(Literal::int(42))]);
+
+ // Parquet file has both columns
+ let parquet_schema = Arc::new(ArrowSchema::new(vec![
+ field_with_id("id", DataType::Int32, false, 1),
+ field_with_id("name", DataType::Utf8, true, 2),
+ ]));
+
+ // Project id, name, and _partition
+ let projected_field_ids = [1, 2, RESERVED_FIELD_ID_PARTITION];
+
+ let partition_column =
+ build_partition_constant(&unified_partition_type, &partition_spec,
&partition_data)
+ .unwrap();
+ let mut transformer =
+ RecordBatchTransformerBuilder::new(snapshot_schema,
&projected_field_ids)
+ .with_partition_constant(partition_column)
+ .build();
+
+ let parquet_batch = RecordBatch::try_new(parquet_schema, vec![
+ Arc::new(Int32Array::from(vec![100, 200, 300])),
+ Arc::new(StringArray::from(vec!["a", "b", "c"])),
+ ])
+ .unwrap();
+
+ let result = transformer.process_record_batch(parquet_batch).unwrap();
+
+ assert_eq!(result.num_columns(), 3);
+ assert_eq!(result.num_rows(), 3);
+
+ // id column from file
+ let id_col = result
+ .column(0)
+ .as_any()
+ .downcast_ref::<Int32Array>()
+ .unwrap();
+ assert_eq!(id_col.values(), &[100, 200, 300]);
+
+ // name column from file
+ let name_col = result
+ .column(1)
+ .as_any()
+ .downcast_ref::<StringArray>()
+ .unwrap();
+ assert_eq!(name_col.value(0), "a");
+
+ // _partition struct column
+ let partition_col = result
+ .column(2)
+ .as_any()
+ .downcast_ref::<StructArray>()
+ .unwrap();
+ assert_eq!(partition_col.num_columns(), 1);
+ assert_eq!(partition_col.len(), 3);
+
+ let inner = partition_col
+ .column(0)
+ .as_any()
+ .downcast_ref::<Int32Array>()
+ .unwrap();
+ assert_eq!(inner.value(0), 42);
+ assert_eq!(inner.value(1), 42);
+ assert_eq!(inner.value(2), 42);
+ }
+
+ #[test]
+ fn partition_column_with_evolution() {
+ use arrow_array::StructArray;
+
+ use crate::metadata_columns::RESERVED_FIELD_ID_PARTITION;
+ use crate::partitioning::compute_unified_partition_type;
+ use crate::spec::Transform;
+
+ // Schema with two fields that could be partition sources
+ let snapshot_schema = Arc::new(
+ Schema::builder()
+ .with_schema_id(0)
+ .with_fields(vec![
+ NestedField::required(1, "year",
Type::Primitive(PrimitiveType::Int)).into(),
+ NestedField::required(2, "month",
Type::Primitive(PrimitiveType::Int)).into(),
+ NestedField::optional(3, "data",
Type::Primitive(PrimitiveType::String)).into(),
+ ])
+ .build()
+ .unwrap(),
+ );
+
+ // Old spec: partition by year only
+ let spec_v0 =
crate::spec::PartitionSpec::builder(snapshot_schema.clone())
+ .with_spec_id(0)
+ .add_partition_field("year", "year", Transform::Identity)
+ .unwrap()
+ .build()
+ .unwrap();
+
+ // New spec: partition by year and month
+ let spec_v1 =
crate::spec::PartitionSpec::builder(snapshot_schema.clone())
+ .with_spec_id(1)
+ .add_partition_field("year", "year", Transform::Identity)
+ .unwrap()
+ .add_partition_field("month", "month", Transform::Identity)
+ .unwrap()
+ .build()
+ .unwrap();
+
+ // Unified type includes both year and month
+ let unified_partition_type =
+ compute_unified_partition_type([&spec_v0, &spec_v1].into_iter(),
&snapshot_schema)
+ .unwrap();
+
+ assert_eq!(unified_partition_type.fields().len(), 2);
+
+ // File written with spec_v0 (only has year=2023)
+ let partition_data = Struct::from_iter(vec![Some(Literal::int(2023))]);
+
+ let parquet_schema = Arc::new(ArrowSchema::new(vec![field_with_id(
+ "data",
+ DataType::Utf8,
+ true,
+ 3,
+ )]));
+
+ let projected_field_ids = [3, RESERVED_FIELD_ID_PARTITION];
+
+ let partition_column =
+ build_partition_constant(&unified_partition_type, &spec_v0,
&partition_data).unwrap();
+ let mut transformer =
+ RecordBatchTransformerBuilder::new(snapshot_schema,
&projected_field_ids)
+ .with_partition_constant(partition_column)
+ .build();
+
+ let parquet_batch =
+ RecordBatch::try_new(parquet_schema,
vec![Arc::new(StringArray::from(vec![
+ "hello", "world",
+ ]))])
+ .unwrap();
+
+ let result = transformer.process_record_batch(parquet_batch).unwrap();
+
+ assert_eq!(result.num_columns(), 2);
+ assert_eq!(result.num_rows(), 2);
+
+ // _partition struct has 2 fields: year (present) and month (null for
this old spec file)
+ let partition_col = result
+ .column(1)
+ .as_any()
+ .downcast_ref::<StructArray>()
+ .unwrap();
+ assert_eq!(partition_col.num_columns(), 2);
+
+ let year_col = partition_col
+ .column(0)
+ .as_any()
+ .downcast_ref::<Int32Array>()
+ .unwrap();
+ assert_eq!(year_col.value(0), 2023);
+ assert_eq!(year_col.value(1), 2023);
+
+ let month_col = partition_col
+ .column(1)
+ .as_any()
+ .downcast_ref::<Int32Array>()
+ .unwrap();
+ assert!(month_col.is_null(0));
+ assert!(month_col.is_null(1));
+ }
}
diff --git a/crates/iceberg/src/arrow/value.rs
b/crates/iceberg/src/arrow/value.rs
index 34f6ce6e9..d22e565a4 100644
--- a/crates/iceberg/src/arrow/value.rs
+++ b/crates/iceberg/src/arrow/value.rs
@@ -21,7 +21,7 @@ use arrow_array::{
Array, ArrayRef, BinaryArray, BooleanArray, Date32Array, Decimal128Array,
FixedSizeBinaryArray,
FixedSizeListArray, Float32Array, Float64Array, Int32Array, Int64Array,
LargeBinaryArray,
LargeListArray, LargeStringArray, ListArray, MapArray, StringArray,
StructArray,
- Time64MicrosecondArray, TimestampMicrosecondArray,
TimestampNanosecondArray,
+ Time64MicrosecondArray, TimestampMicrosecondArray,
TimestampNanosecondArray, new_null_array,
};
use arrow_buffer::NullBuffer;
use arrow_schema::{DataType, FieldRef, TimeUnit};
@@ -820,37 +820,22 @@ pub(crate) fn create_primitive_array_repeated(
num_rows: usize,
) -> Result<ArrayRef> {
Ok(match (data_type, prim_lit) {
+ // --- Primitive Some arms ---
(DataType::Boolean, Some(PrimitiveLiteral::Boolean(value))) => {
Arc::new(BooleanArray::from(vec![*value; num_rows]))
}
- (DataType::Boolean, None) => {
- let vals: Vec<Option<bool>> = vec![None; num_rows];
- Arc::new(BooleanArray::from(vals))
- }
(DataType::Int32, Some(PrimitiveLiteral::Int(value))) => {
Arc::new(Int32Array::from(vec![*value; num_rows]))
}
- (DataType::Int32, None) => {
- let vals: Vec<Option<i32>> = vec![None; num_rows];
- Arc::new(Int32Array::from(vals))
- }
(DataType::Date32, Some(PrimitiveLiteral::Int(value))) => {
Arc::new(Date32Array::from(vec![*value; num_rows]))
}
- (DataType::Date32, None) => {
- let vals: Vec<Option<i32>> = vec![None; num_rows];
- Arc::new(Date32Array::from(vals))
- }
(DataType::Int64, Some(PrimitiveLiteral::Int(value))) => {
Arc::new(Int64Array::from(vec![i64::from(*value); num_rows]))
}
(DataType::Int64, Some(PrimitiveLiteral::Long(value))) => {
Arc::new(Int64Array::from(vec![*value; num_rows]))
}
- (DataType::Int64, None) => {
- let vals: Vec<Option<i64>> = vec![None; num_rows];
- Arc::new(Int64Array::from(vals))
- }
(
DataType::Timestamp(TimeUnit::Microsecond, timezone),
Some(PrimitiveLiteral::Long(value)),
@@ -862,15 +847,6 @@ pub(crate) fn create_primitive_array_repeated(
Arc::new(array)
}
}
- (DataType::Timestamp(TimeUnit::Microsecond, timezone), None) => {
- let vals: Vec<Option<i64>> = vec![None; num_rows];
- let array = TimestampMicrosecondArray::from(vals);
- if let Some(timezone) = timezone {
- Arc::new(array.with_timezone(timezone.clone()))
- } else {
- Arc::new(array)
- }
- }
(
DataType::Timestamp(TimeUnit::Nanosecond, timezone),
Some(PrimitiveLiteral::Long(value)),
@@ -882,42 +858,32 @@ pub(crate) fn create_primitive_array_repeated(
Arc::new(array)
}
}
- (DataType::Timestamp(TimeUnit::Nanosecond, timezone), None) => {
- let vals: Vec<Option<i64>> = vec![None; num_rows];
- let array = TimestampNanosecondArray::from(vals);
- if let Some(timezone) = timezone {
- Arc::new(array.with_timezone(timezone.clone()))
- } else {
- Arc::new(array)
- }
- }
(DataType::Float32, Some(PrimitiveLiteral::Float(value))) => {
Arc::new(Float32Array::from(vec![value.0; num_rows]))
}
- (DataType::Float32, None) => {
- let vals: Vec<Option<f32>> = vec![None; num_rows];
- Arc::new(Float32Array::from(vals))
- }
(DataType::Float64, Some(PrimitiveLiteral::Double(value))) => {
Arc::new(Float64Array::from(vec![value.0; num_rows]))
}
- (DataType::Float64, None) => {
- let vals: Vec<Option<f64>> = vec![None; num_rows];
- Arc::new(Float64Array::from(vals))
- }
(DataType::Utf8, Some(PrimitiveLiteral::String(value))) => {
Arc::new(StringArray::from(vec![value.clone(); num_rows]))
}
- (DataType::Utf8, None) => {
- let vals: Vec<Option<String>> = vec![None; num_rows];
- Arc::new(StringArray::from(vals))
- }
(DataType::Binary, Some(PrimitiveLiteral::Binary(value))) => {
Arc::new(BinaryArray::from_vec(vec![value; num_rows]))
}
- (DataType::Binary, None) => {
- let vals: Vec<Option<&[u8]>> = vec![None; num_rows];
- Arc::new(BinaryArray::from_opt_vec(vals))
+ (DataType::LargeBinary, Some(PrimitiveLiteral::Binary(value))) => {
+ Arc::new(LargeBinaryArray::from_vec(vec![value; num_rows]))
+ }
+ (DataType::FixedSizeBinary(len),
Some(PrimitiveLiteral::Binary(value))) => {
+ let repeated: Vec<&[u8]> = vec![value.as_slice(); num_rows];
+
Arc::new(FixedSizeBinaryArray::try_from_iter(repeated.into_iter()).map_err(|e| {
+ Error::new(
+ ErrorKind::DataInvalid,
+ format!("Failed to create FixedSizeBinary({len}) array:
{e}"),
+ )
+ })?)
+ }
+ (DataType::Time64(TimeUnit::Microsecond),
Some(PrimitiveLiteral::Long(value))) => {
+ Arc::new(Time64MicrosecondArray::from(vec![*value; num_rows]))
}
(DataType::Decimal128(precision, scale),
Some(PrimitiveLiteral::Int128(value))) => {
Arc::new(
@@ -947,6 +913,8 @@ pub(crate) fn create_primitive_array_repeated(
})?,
)
}
+
+ // --- Special-case None arms ---
(DataType::Decimal128(precision, scale), None) => {
let vals: Vec<Option<i128>> = vec![None; num_rows];
Arc::new(
@@ -963,7 +931,7 @@ pub(crate) fn create_primitive_array_repeated(
)
}
(DataType::Struct(fields), None) => {
- // Create a StructArray filled with nulls
+ // Create a StructArray filled with nulls, recursively creating
null children
let null_arrays: Vec<ArrayRef> = fields
.iter()
.map(|field|
create_primitive_array_repeated(field.data_type(), &None, num_rows))
@@ -976,6 +944,10 @@ pub(crate) fn create_primitive_array_repeated(
))
}
(DataType::Null, _) => Arc::new(arrow_array::NullArray::new(num_rows)),
+
+ // --- Catch-all null arm: use arrow-rs new_null_array for any
remaining DataType ---
+ (dt, None) => new_null_array(dt, num_rows),
+
(dt, _) => {
return Err(Error::new(
ErrorKind::Unexpected,
diff --git a/crates/iceberg/src/lib.rs b/crates/iceberg/src/lib.rs
index 4e346460f..301992d15 100644
--- a/crates/iceberg/src/lib.rs
+++ b/crates/iceberg/src/lib.rs
@@ -100,6 +100,7 @@ pub mod writer;
mod delete_vector;
pub mod metadata_columns;
+pub mod partitioning;
pub mod puffin;
/// Utility functions and modules.
pub mod util;
diff --git a/crates/iceberg/src/partitioning.rs
b/crates/iceberg/src/partitioning.rs
new file mode 100644
index 000000000..ae8a58c6a
--- /dev/null
+++ b/crates/iceberg/src/partitioning.rs
@@ -0,0 +1,449 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Partition type utilities for Iceberg tables.
+
+use std::cmp::Reverse;
+use std::collections::{HashMap, HashSet};
+
+use crate::spec::{
+ NestedField, NestedFieldRef, PartitionField, PartitionSpec, Schema,
StructType, Transform, Type,
+};
+use crate::{Error, ErrorKind, Result};
+
+/// Computes the unified partition type across all partition specs in the
table.
+///
+/// This is equivalent to Java's `Partitioning.partitionType(table)`. The
result is a
+/// StructType containing all partition fields ever used across all specs,
enabling correct
+/// representation of the `_partition` metadata column when partition
evolution has occurred.
+///
+/// Matches Java's `buildPartitionProjectionType` behavior:
+/// - Specs are sorted by spec_id in descending order (newer specs first), so
newer field
+/// names take precedence when deduplicating by field_id.
+/// - Unknown transforms cause an error.
+/// - Fields whose source column was dropped from the schema are skipped.
+/// - Two specs defining the same field_id must be compatible (same source,
compatible
+/// transforms); V1 tables do not guarantee field ids are unique across
specs.
+/// - When a newer spec marks a field as Void (dropped) but an older spec has
it with a
+/// real transform, the older spec's type is preserved while the newer
spec's name is kept.
+/// - Fields are deduplicated by field_id; each unique field_id appears
exactly once.
+/// - Output fields are sorted by field_id ascending.
+///
+/// # Arguments
+/// * `partition_specs` - Iterator over all partition specs in the table
+/// * `schema` - The current table schema (needed to determine result types of
transforms)
+pub fn compute_unified_partition_type<'a>(
+ partition_specs: impl Iterator<Item = &'a PartitionSpec>,
+ schema: &Schema,
+) -> Result<StructType> {
+ let mut specs: Vec<&PartitionSpec> = partition_specs.collect();
+ specs.sort_by_key(|s| Reverse(s.spec_id()));
+
+ let active_field_ids = all_active_field_ids(specs.iter().copied(), schema);
+
+ let mut field_map: HashMap<i32, &PartitionField> = HashMap::new();
+ let mut type_map: HashMap<i32, Type> = HashMap::new();
+ let mut name_map: HashMap<i32, String> = HashMap::new();
+
+ for spec in &specs {
+ for field in spec.fields() {
+ 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) {
+ return Err(Error::new(
+ ErrorKind::DataInvalid,
+ format!(
+ "Partition field '{}' uses an unknown transform whose
result type \
+ cannot be determined",
+ field.name
+ ),
+ ));
+ }
+
+ if !active_field_ids.contains(&field_id) {
+ continue;
+ }
+
+ 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)?;
+ field_map.insert(field_id, field);
+ type_map.insert(field_id, res_type);
+ name_map.insert(field_id, field.name.clone());
+ }
+ Some(existing) => {
+ // V1 tables do not guarantee field ids are unique across
specs, so two
+ // specs may define the same field id. They must be
compatible.
+ if !equivalent_ignoring_names(field, existing) {
+ return Err(Error::new(
+ ErrorKind::DataInvalid,
+ format!(
+ "Conflicting partition fields for field id
{field_id}: \
+ '{}' and '{}'",
+ 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)?;
+ field_map.insert(field_id, field);
+ type_map.insert(field_id, res_type);
+ }
+ }
+ }
+ }
+ }
+
+ let mut field_ids: Vec<i32> = field_map.keys().copied().collect();
+ field_ids.sort();
+
+ let struct_fields = field_ids
+ .into_iter()
+ .map(|fid| -> Result<NestedFieldRef> {
+ let name = name_map.get(&fid).ok_or_else(|| {
+ Error::new(
+ ErrorKind::Unexpected,
+ format!("Missing name for partition field {fid}"),
+ )
+ })?;
+ let ty = type_map.remove(&fid).ok_or_else(|| {
+ Error::new(
+ ErrorKind::Unexpected,
+ format!("Missing type for partition field {fid}"),
+ )
+ })?;
+ Ok(NestedField::optional(fid, name, ty).into())
+ })
+ .collect::<Result<Vec<_>>>()?;
+
+ Ok(StructType::new(struct_fields))
+}
+
+fn is_void_transform(field: &PartitionField) -> bool {
+ 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)
+}
+
+/// Transforms are compatible if they are equal, or if either is Void (a
dropped field).
+/// Matches Java's `Partitioning.compatibleTransforms`.
+fn compatible_transforms(t1: &Transform, t2: &Transform) -> bool {
+ t1 == t2 || matches!(t1, Transform::Void) || matches!(t2, Transform::Void)
+}
+
+fn all_active_field_ids<'a>(
+ partition_specs: impl Iterator<Item = &'a PartitionSpec>,
+ schema: &Schema,
+) -> 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)
+ .collect()
+}
+
+#[cfg(test)]
+mod tests {
+ use std::sync::Arc;
+
+ use super::*;
+ use crate::spec::{NestedField, PrimitiveType, Transform, Type,
UnboundPartitionSpec};
+
+ fn test_schema() -> Schema {
+ Schema::builder()
+ .with_fields(vec![
+ NestedField::required(1, "id",
Type::Primitive(PrimitiveType::Int)).into(),
+ NestedField::required(2, "data",
Type::Primitive(PrimitiveType::String)).into(),
+ NestedField::required(3, "ts",
Type::Primitive(PrimitiveType::Timestamp)).into(),
+ NestedField::required(4, "category",
Type::Primitive(PrimitiveType::String)).into(),
+ ])
+ .build()
+ .unwrap()
+ }
+
+ fn build_spec(
+ schema: &Schema,
+ spec_id: i32,
+ fields: Vec<(i32, &str, Transform)>,
+ ) -> PartitionSpec {
+ 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)
+ .unwrap();
+ }
+ builder.build().bind(schema.clone()).unwrap()
+ }
+
+ #[test]
+ fn test_single_spec_identity() {
+ let schema = test_schema();
+ let spec = build_spec(&schema, 0, vec![(4, "category",
Transform::Identity)]);
+
+ let result = compute_unified_partition_type([&spec].into_iter(),
&schema).unwrap();
+ assert_eq!(result.fields().len(), 1);
+ assert_eq!(result.fields()[0].name, "category");
+ assert_eq!(
+ *result.fields()[0].field_type,
+ Type::Primitive(PrimitiveType::String)
+ );
+ }
+
+ #[test]
+ fn test_single_spec_with_year_transform() {
+ let schema = test_schema();
+ let spec = build_spec(&schema, 0, vec![(3, "ts_year",
Transform::Year)]);
+
+ let result = compute_unified_partition_type([&spec].into_iter(),
&schema).unwrap();
+ assert_eq!(result.fields().len(), 1);
+ assert_eq!(result.fields()[0].name, "ts_year");
+ assert_eq!(
+ *result.fields()[0].field_type,
+ Type::Primitive(PrimitiveType::Int)
+ );
+ }
+
+ #[test]
+ fn test_unpartitioned() {
+ let schema = test_schema();
+ let spec = PartitionSpec::unpartition_spec();
+ let result = compute_unified_partition_type([&spec].into_iter(),
&schema).unwrap();
+ assert!(result.fields().is_empty());
+ }
+
+ #[test]
+ fn test_multiple_fields_sorted_by_id() {
+ let schema = test_schema();
+ let spec = build_spec(&schema, 0, vec![
+ (3, "ts_year", Transform::Year),
+ (4, "category", Transform::Identity),
+ ]);
+
+ let result = compute_unified_partition_type([&spec].into_iter(),
&schema).unwrap();
+ assert_eq!(result.fields().len(), 2);
+ assert!(result.fields()[0].id < result.fields()[1].id);
+ }
+
+ #[test]
+ fn test_newer_name_takes_precedence() {
+ let schema = test_schema();
+
+ // 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,
+ })
+ .unwrap()
+ .build()
+ .unwrap();
+
+ // 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,
+ })
+ .unwrap()
+ .build()
+ .unwrap();
+
+ let result =
+ compute_unified_partition_type([&spec_v0, &spec_v1].into_iter(),
&schema).unwrap();
+ assert_eq!(result.fields().len(), 1);
+ assert_eq!(result.fields()[0].name, "cat_new");
+ }
+
+ #[test]
+ fn test_void_replaced_by_older_non_void() {
+ let schema = test_schema();
+
+ // 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,
+ })
+ .unwrap()
+ .build()
+ .unwrap();
+
+ // 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,
+ })
+ .unwrap()
+ .build()
+ .unwrap();
+
+ let result =
+ compute_unified_partition_type([&spec_v0, &spec_v1].into_iter(),
&schema).unwrap();
+
+ assert_eq!(result.fields().len(), 1);
+ // Name from newer spec
+ assert_eq!(result.fields()[0].name, "category_v2");
+ // Type from older non-void spec
+ assert_eq!(
+ *result.fields()[0].field_type,
+ Type::Primitive(PrimitiveType::String)
+ );
+ }
+
+ #[test]
+ fn test_dropped_source_column_skipped() {
+ // Schema without field 4 (category was dropped)
+ let schema = Schema::builder()
+ .with_fields(vec![
+ NestedField::required(1, "id",
Type::Primitive(PrimitiveType::Int)).into(),
+ NestedField::required(2, "data",
Type::Primitive(PrimitiveType::String)).into(),
+ ])
+ .build()
+ .unwrap();
+
+ // Spec references source_id=4 which no longer exists in the schema.
+ // Deserialize directly since the builder rejects unknown source
columns.
+ let spec = serde_json::from_value::<PartitionSpec>(serde_json::json!({
+ "spec-id": 0,
+ "fields": [{
+ "source-id": 4,
+ "field-id": 1000,
+ "name": "category",
+ "transform": "identity"
+ }]
+ }))
+ .unwrap();
+
+ let result = compute_unified_partition_type([&spec].into_iter(),
&schema).unwrap();
+ assert!(result.fields().is_empty());
+ }
+
+ #[test]
+ fn test_evolution_adds_new_field() {
+ let schema = test_schema();
+
+ // Spec 0: partition by category
+ let spec_v0 = build_spec(&schema, 0, vec![(4, "category",
Transform::Identity)]);
+
+ // 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,
+ })
+ .unwrap()
+ .add_partition_field("ts", "ts_year", Transform::Year)
+ .unwrap()
+ .build()
+ .unwrap();
+
+ let result =
+ compute_unified_partition_type([&spec_v0, &spec_v1].into_iter(),
&schema).unwrap();
+ assert_eq!(result.fields().len(), 2);
+ }
+
+ #[test]
+ fn test_unknown_transform_errors() {
+ let schema = test_schema();
+
+ // A spec using an unknown transform. Deserialize directly since the
builder
+ // validates transforms.
+ let spec = serde_json::from_value::<PartitionSpec>(serde_json::json!({
+ "spec-id": 0,
+ "fields": [{
+ "source-id": 4,
+ "field-id": 1000,
+ "name": "category",
+ "transform": "unknown"
+ }]
+ }))
+ .unwrap();
+
+ let err = compute_unified_partition_type([&spec].into_iter(),
&schema).unwrap_err();
+ assert_eq!(err.kind(), ErrorKind::DataInvalid);
+ }
+
+ #[test]
+ fn test_conflicting_partition_fields_error() {
+ let schema = test_schema();
+
+ // Spec 0: field id 1000 -> source 4 (category), identity
+ let spec_v0 =
serde_json::from_value::<PartitionSpec>(serde_json::json!({
+ "spec-id": 0,
+ "fields": [{
+ "source-id": 4,
+ "field-id": 1000,
+ "name": "category",
+ "transform": "identity"
+ }]
+ }))
+ .unwrap();
+
+ // Spec 1: field id 1000 reused for a different source (ts) and
transform (year).
+ // This conflicts with spec 0 and must be rejected.
+ let spec_v1 =
serde_json::from_value::<PartitionSpec>(serde_json::json!({
+ "spec-id": 1,
+ "fields": [{
+ "source-id": 3,
+ "field-id": 1000,
+ "name": "ts_year",
+ "transform": "year"
+ }]
+ }))
+ .unwrap();
+
+ let err =
+ compute_unified_partition_type([&spec_v0, &spec_v1].into_iter(),
&schema).unwrap_err();
+ assert_eq!(err.kind(), ErrorKind::DataInvalid);
+ }
+}
diff --git a/crates/iceberg/src/scan/context.rs
b/crates/iceberg/src/scan/context.rs
index 9181d0272..176bcba45 100644
--- a/crates/iceberg/src/scan/context.rs
+++ b/crates/iceberg/src/scan/context.rs
@@ -29,7 +29,7 @@ use crate::scan::{
};
use crate::spec::{
ManifestContentType, ManifestEntryRef, ManifestFile, ManifestList,
NameMapping,
- PartitionSpecRef, SchemaRef, SnapshotRef, TableMetadataRef,
+ PartitionSpecRef, SchemaRef, SnapshotRef, StructType, TableMetadataRef,
};
use crate::{Error, ErrorKind, Result};
@@ -49,6 +49,7 @@ pub(crate) struct ManifestFileContext {
name_mapping: Option<Arc<NameMapping>>,
case_sensitive: bool,
partition_spec: Option<PartitionSpecRef>,
+ unified_partition_type: Option<Arc<StructType>>,
}
/// Wraps a [`ManifestEntryRef`] alongside the objects that are needed
@@ -65,6 +66,7 @@ pub(crate) struct ManifestEntryContext {
pub name_mapping: Option<Arc<NameMapping>>,
pub case_sensitive: bool,
pub partition_spec: Option<PartitionSpecRef>,
+ pub unified_partition_type: Option<Arc<StructType>>,
}
impl ManifestFileContext {
@@ -83,6 +85,7 @@ impl ManifestFileContext {
name_mapping,
case_sensitive,
partition_spec,
+ unified_partition_type,
} = self;
let manifest = object_cache.get_manifest(&manifest_file).await?;
@@ -100,6 +103,7 @@ impl ManifestFileContext {
name_mapping: name_mapping.clone(),
case_sensitive,
partition_spec: partition_spec.clone(),
+ unified_partition_type: unified_partition_type.clone(),
};
sender
@@ -141,6 +145,7 @@ impl ManifestEntryContext {
.with_partition(Some(self.manifest_entry.data_file.partition.clone()))
.with_partition_spec(self.partition_spec.clone())
.with_name_mapping(self.name_mapping)
+ .with_unified_partition_type(self.unified_partition_type.clone())
.with_case_sensitive(self.case_sensitive)
.with_key_metadata(self.manifest_entry.data_file.key_metadata().map(Box::from))
.build())
@@ -165,6 +170,8 @@ pub(crate) struct PlanContext {
pub partition_filter_cache: Arc<PartitionFilterCache>,
pub manifest_evaluator_cache: Arc<ManifestEvaluatorCache>,
pub expression_evaluator_cache: Arc<ExpressionEvaluatorCache>,
+
+ pub unified_partition_type: Option<Arc<StructType>>,
}
impl PlanContext {
@@ -294,6 +301,7 @@ impl PlanContext {
.table_metadata
.partition_spec_by_id(manifest_file.partition_spec_id)
.cloned(),
+ unified_partition_type: self.unified_partition_type.clone(),
}
}
}
diff --git a/crates/iceberg/src/scan/mod.rs b/crates/iceberg/src/scan/mod.rs
index 19f46c2cb..a1c85c711 100644
--- a/crates/iceberg/src/scan/mod.rs
+++ b/crates/iceberg/src/scan/mod.rs
@@ -37,7 +37,10 @@ use crate::delete_file_index::DeleteFileIndex;
use
crate::expr::visitors::inclusive_metrics_evaluator::InclusiveMetricsEvaluator;
use crate::expr::{Bind, BoundPredicate, Predicate};
use crate::io::FileIO;
-use crate::metadata_columns::{get_metadata_field_id, is_metadata_column_name};
+use crate::metadata_columns::{
+ RESERVED_FIELD_ID_PARTITION, get_metadata_field_id,
is_metadata_column_name,
+};
+use crate::partitioning::compute_unified_partition_type;
use crate::runtime::Runtime;
use crate::spec::{DEFAULT_SCHEMA_NAME_MAPPING, DataContentType, NameMapping,
SnapshotRef};
use crate::table::Table;
@@ -300,6 +303,20 @@ impl<'a> TableScanBuilder<'a> {
.transpose()?
.map(Arc::new);
+ // Compute unified partition type if _partition is projected
+ let unified_partition_type = if
field_ids.contains(&RESERVED_FIELD_ID_PARTITION) {
+ let partition_type = compute_unified_partition_type(
+ self.table
+ .metadata()
+ .partition_specs_iter()
+ .map(|s| s.as_ref()),
+ &schema,
+ )?;
+ Some(Arc::new(partition_type))
+ } else {
+ None
+ };
+
let plan_context = PlanContext {
snapshot,
table_metadata: self.table.metadata_ref(),
@@ -313,6 +330,7 @@ impl<'a> TableScanBuilder<'a> {
partition_filter_cache: Arc::new(PartitionFilterCache::new()),
manifest_evaluator_cache: Arc::new(ManifestEvaluatorCache::new()),
expression_evaluator_cache:
Arc::new(ExpressionEvaluatorCache::new()),
+ unified_partition_type,
};
Ok(TableScan {
diff --git a/crates/iceberg/src/scan/task.rs b/crates/iceberg/src/scan/task.rs
index faeac51be..24a603687 100644
--- a/crates/iceberg/src/scan/task.rs
+++ b/crates/iceberg/src/scan/task.rs
@@ -25,7 +25,7 @@ use crate::Result;
use crate::expr::BoundPredicate;
use crate::spec::{
DataContentType, DataFileFormat, ManifestEntryRef, NameMapping,
PartitionSpec, Schema,
- SchemaRef, Struct,
+ SchemaRef, Struct, StructType,
};
/// A stream of [`FileScanTask`].
@@ -116,6 +116,22 @@ pub struct FileScanTask {
#[builder(default)]
pub name_mapping: Option<Arc<NameMapping>>,
+ /// The unified partition type across all specs in the table.
+ /// When `RESERVED_FIELD_ID_PARTITION` is in the projected field IDs, the
reader
+ /// uses this type along with the task's partition_spec and partition data
to
+ /// materialize the `_partition` struct column at read time.
+ ///
+ /// This is a table-level value (same for all tasks in a scan), stored
per-task
+ /// so that readers are self-contained without needing back-pointers to
table
+ /// metadata. The cost is one Arc clone per task.
+ /// Serde: not yet implemented (same pattern as partition, partition_spec,
name_mapping).
+ #[serde(default)]
+ #[serde(skip_serializing_if = "Option::is_none")]
+ #[serde(serialize_with = "serialize_not_implemented")]
+ #[serde(deserialize_with = "deserialize_not_implemented")]
+ #[builder(default)]
+ pub unified_partition_type: Option<Arc<StructType>>,
+
/// Whether this scan task should treat column names as case-sensitive
when binding predicates.
pub case_sensitive: bool,