comphead commented on code in PR #3282:
URL: https://github.com/apache/iceberg-rust/pull/3282#discussion_r4127161068
##########
crates/iceberg/src/arrow/record_batch_projector.rs:
##########
@@ -176,21 +176,32 @@ impl RecordBatchProjector {
}
fn get_column_by_field_index(batch: &[ArrayRef], field_index: &[usize]) ->
Result<ArrayRef> {
+ if let [index] = field_index {
+ return Ok(batch[*index].clone());
+ }
+
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 parent_null_buffer = 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());
+ ))?;
+ parent_null_buffer = NullBuffer::union(
+ parent_null_buffer.as_ref(),
+ struct_array.logical_nulls().as_ref(),
+ );
+ array = struct_array.column(*idx).clone();
Review Comment:
Two costs on the zero-copy path here:
- `ok_or(Error::new(..))` builds the error on every successful downcast.
`Error::new` allocates a `String` and calls `Backtrace::capture()`, which walks
the stack when `RUST_BACKTRACE` is set. `ok_or_else` avoids that, and Clippy's
`or_fun_call` flags this line.
- `StructArray` doesn't override `logical_nulls`, so `logical_nulls()`
clones the buffer only to borrow it. `nulls()` borrows. Walking with
`&ArrayRef` and cloning once at the end also saves an `Arc` clone per level.
Together these measured 28.1 ns → 6.6 ns at 1 level and 108 ns → 18.7 ns at
4 levels, with no parent nulls. Nit: this accumulates every enclosing struct,
so `ancestor_nulls` may read more accurately.
##########
crates/iceberg/src/arrow/record_batch_projector.rs:
##########
@@ -176,21 +176,32 @@ impl RecordBatchProjector {
}
fn get_column_by_field_index(batch: &[ArrayRef], field_index: &[usize]) ->
Result<ArrayRef> {
+ if let [index] = field_index {
+ return Ok(batch[*index].clone());
+ }
+
Review Comment:
This early return looks redundant after the third commit. For a one-element
path the loop below doesn't run, `parent_null_buffer` stays `None`, and the
`let-else` returns the same `Arc`.
```suggestion
```
##########
crates/iceberg/src/arrow/record_batch_projector.rs:
##########
@@ -176,21 +176,32 @@ impl RecordBatchProjector {
}
fn get_column_by_field_index(batch: &[ArrayRef], field_index: &[usize]) ->
Result<ArrayRef> {
+ if let [index] = field_index {
+ return Ok(batch[*index].clone());
+ }
+
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 parent_null_buffer = 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());
+ ))?;
+ parent_null_buffer = NullBuffer::union(
+ parent_null_buffer.as_ref(),
+ struct_array.logical_nulls().as_ref(),
+ );
+ array = struct_array.column(*idx).clone();
}
+ let Some(parent_null_buffer) = parent_null_buffer else {
+ return Ok(array);
+ };
Review Comment:
Consider also skipping the rebuild when the leaf already carries these
nulls. arrow-rs's Parquet reader derives child validity from definition levels,
so nullable children read from Parquet already contain every ancestor null. A
round trip confirmed this for `Int32` and `Utf8` leaves. Required leaves come
back without a null buffer and still take the rebuild.
With the changes above plus a `contains` check, the parent-nulls path
measured 259 ns → 117 ns for `Int32` and 18.3 µs → 106 ns for a 240 KiB `Utf8`
leaf. A failed check costs about 20 ns.
A short comment may help here too, since `NullBuffer::union` returning
`None` for all-valid masks is what keeps sliced batches on the zero-copy path.
<details><summary>Function with these suggestions applied</summary>
```rust
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()];
let mut parent_null_buffer = None;
for idx in rev_iterator {
parent_null_buffer = NullBuffer::union(parent_null_buffer.as_ref(),
array.nulls());
array = array
.as_any()
.downcast_ref::<StructArray>()
.ok_or_else(|| {
Error::new(ErrorKind::Unexpected, "Cannot convert Array to
StructArray")
})?
.column(*idx);
}
// Only nulls in an enclosing struct force a rebuild.
`NullBuffer::union` returns `None`
// when no input has a null, so all-valid masks (e.g. from slicing) keep
the array as-is.
let Some(parent_null_buffer) = parent_null_buffer else {
return Ok(array.clone());
};
// Readers such as arrow-rs's Parquet reader already carry ancestor
nulls in nullable
// children, so the rebuild would not change anything.
if array
.nulls()
.is_some_and(|nulls| nulls.contains(&parent_null_buffer))
{
return Ok(array.clone());
}
let null_buffer =
NullBuffer::union(Some(&parent_null_buffer),
array.logical_nulls().as_ref());
Ok(make_array(
array.to_data().into_builder().nulls(null_buffer).build()?,
))
}
```
</details>
##########
crates/iceberg/src/arrow/record_batch_projector.rs:
##########
@@ -176,21 +176,32 @@ impl RecordBatchProjector {
}
fn get_column_by_field_index(batch: &[ArrayRef], field_index: &[usize]) ->
Result<ArrayRef> {
+ if let [index] = field_index {
+ return Ok(batch[*index].clone());
+ }
+
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 parent_null_buffer = 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());
+ ))?;
+ parent_null_buffer = NullBuffer::union(
+ parent_null_buffer.as_ref(),
+ struct_array.logical_nulls().as_ref(),
+ );
+ array = struct_array.column(*idx).clone();
}
+ let Some(parent_null_buffer) = parent_null_buffer else {
+ return Ok(array);
+ };
+ let null_buffer =
+ NullBuffer::union(Some(&parent_null_buffer),
array.logical_nulls().as_ref());
Ok(make_array(
array.to_data().into_builder().nulls(null_buffer).build()?,
))
Review Comment:
Optional, and it predates this PR: `build()` runs full validation, which
recounts the nulls and checks UTF-8 across the whole values buffer.
`arrow_select::nullif` computes the mask and the null count in one pass and
uses arrow's own unchecked build, so no `unsafe` is needed here.
```rust
use arrow_array::BooleanArray;
use arrow_select::nullif::nullif;
let is_parent_null = BooleanArray::new(!parent_null_buffer.inner(), None);
Ok(nullif(array.as_ref(), &is_parent_null)?)
```
On top of the changes above, a 240 KiB `Utf8` leaf built in memory measured
19.8 µs → 268 ns. For `Int32` it was 233 ns → 260 ns. Note that `nullif` reads
the leaf's physical nulls, passes `Null` leaves through where `build()`
currently errors, and would need a guard for `Union` leaves to keep today's
error.
##########
crates/iceberg/src/arrow/record_batch_projector.rs:
##########
@@ -271,6 +309,77 @@ mod test {
assert_eq!(projected_inner_int_array.values(), &[4, 5, 6]);
}
+ #[test]
+ fn test_record_batch_projector_top_level_nullable_column() {
Review Comment:
This test also passes on the pre-PR code because it only checks values, and
`col0` of `test_equality_delete_with_nullable_field` already covers a nullable
top-level column. Asserting `Arc::ptr_eq(&projected[0], &input)` would pin the
zero-copy behavior.
One option is a single identity test in place of this test, the two below,
and `nested_projector`. Its last block assumes the `contains` check suggested
above, and the `Array` test import becomes unused.
<details><summary>Test</summary>
```rust
#[test]
fn test_record_batch_projector_reuses_arrays() {
let leaf_fields = Fields::from(vec![Field::new("leaf", DataType::Int32,
true)]);
let schema = Arc::new(Schema::new(vec![
Field::new("id", DataType::Int32, true),
Field::new("parent", DataType::Struct(leaf_fields.clone()), true),
]));
let field_id_fetch_func = |field: &Field| match field.name().as_str() {
"id" => Ok(Some(1)),
"parent" => Ok(Some(2)),
"leaf" => Ok(Some(3)),
_ => Err(Error::new(ErrorKind::Unexpected, "Field id not found")),
};
let projector =
RecordBatchProjector::new(schema, &[1, 3], field_id_fetch_func, |_|
true).unwrap();
// Leaf nulls alone do not force a rebuild.
let id = Arc::new(Int32Array::from(vec![Some(1), None])) as ArrayRef;
let leaf = Arc::new(Int32Array::from(vec![None, Some(2)])) as ArrayRef;
let parent = StructArray::new(leaf_fields.clone(), vec![leaf.clone()],
None);
let projected = projector
.project_column(&[id.clone(), Arc::new(parent)])
.unwrap();
assert!(Arc::ptr_eq(&projected[0], &id));
assert!(Arc::ptr_eq(&projected[1], &leaf));
// Slicing past the parent's only null leaves an all-valid mask, which
is not a rebuild either.
let parent = StructArray::new(
leaf_fields.clone(),
vec![Arc::new(Int32Array::from(vec![Some(0), None, Some(2)])) as
ArrayRef],
Some(NullBuffer::from(vec![false, true, true])),
)
.slice(1, 2);
let leaf = parent.column(0).clone();
let projected = projector
.project_column(&[id.clone(), Arc::new(parent)])
.unwrap();
assert!(Arc::ptr_eq(&projected[1], &leaf));
// A leaf that already carries its parent's nulls, as read from Parquet,
is kept too.
let leaf = Arc::new(Int32Array::from(vec![None, Some(2)])) as ArrayRef;
let parent = StructArray::new(
leaf_fields,
vec![leaf.clone()],
Some(NullBuffer::from(vec![false, true])),
);
let projected = projector.project_column(&[id,
Arc::new(parent)]).unwrap();
assert!(Arc::ptr_eq(&projected[1], &leaf));
}
```
</details>
##########
crates/iceberg/src/arrow/record_batch_projector.rs:
##########
@@ -271,6 +309,77 @@ 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(&[input]).unwrap();
+ 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, schema, 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],
+ Some(NullBuffer::new_valid(3)),
Review Comment:
`StructArray::new` drops a null buffer that has no nulls (`nulls.filter(|n|
n.null_count() > 0)`), so this is stored as `None` and the all-valid case isn't
exercised. A sliced parent keeps `Some`, for example `StructArray::new(fields,
vec![leaf], Some(NullBuffer::from(vec![false, true, true]))).slice(1, 2)`.
##########
crates/iceberg/src/arrow/record_batch_projector.rs:
##########
@@ -271,6 +309,77 @@ 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(&[input]).unwrap();
+ 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, schema, 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],
+ Some(NullBuffer::new_valid(3)),
+ )) as ArrayRef;
+ let batch = RecordBatch::try_new(schema, vec![outer]).unwrap();
+
+ let projected = projector.project_column(batch.columns()).unwrap();
+ assert!(Arc::ptr_eq(&projected[0], &leaf));
+ let projected_leaf =
projected[0].as_any().downcast_ref::<Int32Array>().unwrap();
+ assert_eq!(projected_leaf.value(0), 10);
+ assert!(projected_leaf.is_null(1));
+ assert_eq!(projected_leaf.value(2), 30);
+ }
+
+ #[test]
+ fn test_record_batch_projector_propagates_nested_parent_nulls() {
Review Comment:
This also passes on the pre-PR code and overlaps `col1` and `col2` of
`test_equality_delete_with_nullable_field`, which cover one- and two-level
parent nulls plus leaf nulls. If it stays, the five assertions could be one
`assert_eq!(projected[0].as_ref(), &Int32Array::from(vec![Some(10), None, None,
None]))`, and the `RecordBatch` isn't needed since `project_column` takes
`&[ArrayRef]`.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]