This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] 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 e9c22bf71 perf(arrow): skip rebuilding projected arrays without parent
nulls (#3282)
e9c22bf71 is described below
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![