This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/main/pr-3282-fc682d7a050cdc05e2b90ccf04192034dc8f1a1c in repository https://gitbox.apache.org/repos/asf/iceberg-rust.git
commit e9c22bf7148a2fadd8dc94879371866dea89c883 Author: Minh Vu <[email protected]> AuthorDate: Fri Oct 2 09:57:45 2026 +0000 perf(arrow): skip rebuilding projected arrays without parent nulls (#3282) * perf(arrow): skip rebuilding top-level projected arrays * perf(arrow): drop in-crate benchmark harness * perf(arrow): skip rebuilding nested projected arrays * perf(arrow): tighten projector zero-copy path * perf(arrow): polish projector zero-copy path --- crates/iceberg/src/arrow/record_batch_projector.rs | 160 ++++++++++++++++++--- 1 file changed, 144 insertions(+), 16 deletions(-) diff --git a/crates/iceberg/src/arrow/record_batch_projector.rs b/crates/iceberg/src/arrow/record_batch_projector.rs index d0802ee53..6dda6b0f4 100644 --- a/crates/iceberg/src/arrow/record_batch_projector.rs +++ b/crates/iceberg/src/arrow/record_batch_projector.rs @@ -17,7 +17,7 @@ use std::sync::Arc; -use arrow_array::{ArrayRef, RecordBatch, StructArray, make_array}; +use arrow_array::{Array, ArrayRef, RecordBatch, StructArray, make_array}; use arrow_buffer::NullBuffer; use arrow_schema::{DataType, Field, FieldRef, Fields, Schema, SchemaRef}; use parquet::arrow::PARQUET_FIELD_ID_META_KEY; @@ -169,28 +169,33 @@ impl RecordBatchProjector { /// Do projection with columns pub fn project_column(&self, batch: &[ArrayRef]) -> Result<Vec<ArrayRef>> { - self.field_indices - .iter() - .map(|index_vec| Self::get_column_by_field_index(batch, index_vec)) - .collect::<Result<Vec<_>>>() + let mut columns = Vec::with_capacity(self.field_indices.len()); + for index_vec in &self.field_indices { + columns.push(Self::get_column_by_field_index(batch, index_vec)?); + } + Ok(columns) } fn get_column_by_field_index(batch: &[ArrayRef], field_index: &[usize]) -> Result<ArrayRef> { let mut rev_iterator = field_index.iter().rev(); - let mut array = batch[*rev_iterator.next().unwrap()].clone(); - let mut null_buffer = array.logical_nulls(); + let mut array = &batch[*rev_iterator.next().unwrap()]; + let mut ancestor_nulls = None; for idx in rev_iterator { - array = array + let struct_array = array .as_any() .downcast_ref::<StructArray>() - .ok_or(Error::new( - ErrorKind::Unexpected, - "Cannot convert Array to StructArray", - ))? - .column(*idx) - .clone(); - null_buffer = NullBuffer::union(null_buffer.as_ref(), array.logical_nulls().as_ref()); + .ok_or_else(|| { + Error::new(ErrorKind::Unexpected, "Cannot convert Array to StructArray") + })?; + ancestor_nulls = NullBuffer::union(ancestor_nulls.as_ref(), struct_array.nulls()); + array = struct_array.column(*idx); } + let Some(ancestor_nulls) = ancestor_nulls else { + // Iceberg leaf types handled here have logical_nulls() == nulls(), so no null + // materialization is needed when no ancestor contributes nulls. + return Ok(Arc::clone(array)); + }; + let null_buffer = NullBuffer::union(Some(&ancestor_nulls), array.logical_nulls().as_ref()); Ok(make_array( array.to_data().into_builder().nulls(null_buffer).build()?, )) @@ -201,13 +206,40 @@ impl RecordBatchProjector { mod test { use std::sync::Arc; - use arrow_array::{ArrayRef, Int32Array, RecordBatch, StringArray, StructArray}; + use arrow_array::{Array, ArrayRef, Int32Array, RecordBatch, StringArray, StructArray}; + use arrow_buffer::NullBuffer; use arrow_schema::{DataType, Field, Fields, Schema}; use crate::arrow::record_batch_projector::RecordBatchProjector; use crate::spec::{NestedField, PrimitiveType, Schema as IcebergSchema, Type}; use crate::{Error, ErrorKind}; + fn nested_projector() -> (RecordBatchProjector, Field, Field) { + let leaf_field = Field::new("leaf", DataType::Int32, true); + let inner_field = Field::new( + "inner", + DataType::Struct(Fields::from(vec![leaf_field.clone()])), + true, + ); + let outer_field = Field::new( + "outer", + DataType::Struct(Fields::from(vec![inner_field.clone()])), + true, + ); + let schema = Arc::new(Schema::new(vec![outer_field])); + let field_id_fetch_func = |field: &Field| -> crate::Result<Option<i64>> { + match field.name().as_str() { + "outer" => Ok(Some(1)), + "inner" => Ok(Some(2)), + "leaf" => Ok(Some(3)), + _ => Ok(None), + } + }; + let projector = + RecordBatchProjector::new(schema, &[3], field_id_fetch_func, |_| true).unwrap(); + (projector, inner_field, leaf_field) + } + #[test] fn test_record_batch_projector_nested_level() { let inner_fields = vec![ @@ -271,6 +303,102 @@ mod test { assert_eq!(projected_inner_int_array.values(), &[4, 5, 6]); } + #[test] + fn test_record_batch_projector_top_level_nullable_column() { + let iceberg_schema = IcebergSchema::builder() + .with_schema_id(0) + .with_fields(vec![ + NestedField::optional(1, "id", Type::Primitive(PrimitiveType::Int)).into(), + ]) + .build() + .unwrap(); + let projector = + RecordBatchProjector::from_iceberg_schema(Arc::new(iceberg_schema), &[1]).unwrap(); + let input = Arc::new(Int32Array::from(vec![Some(10), None, Some(30)])) as ArrayRef; + + let projected = projector + .project_column(std::slice::from_ref(&input)) + .unwrap(); + assert!(Arc::ptr_eq(&projected[0], &input)); + let projected_array = projected[0].as_any().downcast_ref::<Int32Array>().unwrap(); + + assert_eq!(projected_array.value(0), 10); + assert_eq!(projected_array.null_count(), 1); + assert!(projected_array.is_null(1)); + assert_eq!(projected_array.value(2), 30); + } + + #[test] + fn test_record_batch_projector_nested_nullable_leaf_without_parent_nulls() { + let (projector, inner_field, leaf_field) = nested_projector(); + let leaf = Arc::new(Int32Array::from(vec![Some(10), None, Some(30)])) as ArrayRef; + let inner = Arc::new(StructArray::new( + Fields::from(vec![leaf_field]), + vec![leaf.clone()], + None, + )) as ArrayRef; + let outer = Arc::new(StructArray::new( + Fields::from(vec![inner_field]), + vec![inner], + None, + )) as ArrayRef; + + let projected = projector.project_column(&[outer]).unwrap(); + assert!(Arc::ptr_eq(&projected[0], &leaf)); + } + + #[test] + fn test_record_batch_projector_nested_all_valid_slice_reuses_leaf() { + let (projector, inner_field, leaf_field) = nested_projector(); + let leaf = Arc::new(Int32Array::from(vec![Some(0), Some(10), None, Some(30)])) as ArrayRef; + let inner = Arc::new(StructArray::new( + Fields::from(vec![leaf_field]), + vec![leaf], + None, + )) as ArrayRef; + let outer = StructArray::new( + Fields::from(vec![inner_field]), + vec![inner], + Some(NullBuffer::from(vec![false, true, true, true])), + ) + .slice(1, 3); + assert_eq!(outer.null_count(), 0); + assert!(outer.nulls().is_some()); + + let sliced_inner = outer + .column(0) + .as_any() + .downcast_ref::<StructArray>() + .unwrap(); + let sliced_leaf = sliced_inner.column(0).clone(); + let projected = projector + .project_column(&[Arc::new(outer) as ArrayRef]) + .unwrap(); + assert!(Arc::ptr_eq(&projected[0], &sliced_leaf)); + } + + #[test] + fn test_record_batch_projector_propagates_nested_parent_nulls() { + let (projector, inner_field, leaf_field) = nested_projector(); + let leaf = Arc::new(Int32Array::from(vec![Some(10), Some(20), Some(30), None])) as ArrayRef; + let inner = Arc::new(StructArray::new( + Fields::from(vec![leaf_field]), + vec![leaf], + Some(NullBuffer::from(vec![true, true, false, true])), + )) as ArrayRef; + let outer = Arc::new(StructArray::new( + Fields::from(vec![inner_field]), + vec![inner], + Some(NullBuffer::from(vec![true, false, true, true])), + )) as ArrayRef; + + let projected = projector.project_column(&[outer]).unwrap(); + assert_eq!( + projected[0].as_ref(), + &Int32Array::from(vec![Some(10), None, None, None]) + ); + } + #[test] fn test_field_not_found() { let inner_fields = vec