This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git


The following commit(s) were added to refs/heads/main by this push:
     new dd6c95ae Support nested and vector values and deletion vectors in 
batch Diff (#944)
dd6c95ae is described below

commit dd6c95ae6791f00642a388bd0b940c97313a80cf
Author: Jingsong Lee <[email protected]>
AuthorDate: Thu Sep 24 20:46:40 2026 +0800

    Support nested and vector values and deletion vectors in batch Diff (#944)
---
 crates/paimon/src/table/table_read.rs              | 263 ++++++-
 crates/paimon/src/table/table_scan.rs              |  59 +-
 .../paimon/tests/incremental_diff_extended_test.rs | 804 +++++++++++++++++++++
 3 files changed, 1071 insertions(+), 55 deletions(-)

diff --git a/crates/paimon/src/table/table_read.rs 
b/crates/paimon/src/table/table_read.rs
index d783ee00..81a5b200 100644
--- a/crates/paimon/src/table/table_read.rs
+++ b/crates/paimon/src/table/table_read.rs
@@ -34,8 +34,8 @@ use crate::spec::{
 };
 use crate::DataSplit;
 use arrow_array::{
-    builder::StringBuilder, Array, ArrayRef, RecordBatch, RecordBatchOptions, 
StringArray,
-    UInt32Array,
+    builder::StringBuilder, Array, ArrayRef, FixedSizeListArray, ListArray, 
MapArray, RecordBatch,
+    RecordBatchOptions, StringArray, StructArray, UInt32Array,
 };
 use arrow_schema::Schema as ArrowSchema;
 use arrow_select::concat::concat as arrow_concat;
@@ -648,7 +648,8 @@ impl<'a> PaimonTableRead<'a> {
         let audit_schema = audit_schema_for_read_type(&self.read_type, 
include_sequence)?;
 
         let mut diff_read_type = self.table.schema().fields().to_vec();
-        ensure_diff_supported_read_type(&diff_read_type)?;
+        let key_indices = primary_key_indices(self.table, &diff_read_type)?;
+        ensure_diff_supported_read_type(&diff_read_type, &key_indices)?;
         if include_sequence {
             diff_read_type.insert(
                 0,
@@ -734,8 +735,8 @@ impl<'a> PaimonTableRead<'a> {
         after: &[DataSplit],
     ) -> crate::Result<ArrowRecordBatchStream> {
         let diff_read_type = self.table.schema().fields().to_vec();
-        ensure_diff_supported_read_type(&diff_read_type)?;
         let key_indices = primary_key_indices(self.table, &diff_read_type)?;
+        ensure_diff_supported_read_type(&diff_read_type, &key_indices)?;
         let value_indices = value_indices_for_diff(self.table, 
&diff_read_type);
         let output_schema = build_target_arrow_schema(&self.read_type)?;
         let output_col_indices = self
@@ -817,16 +818,6 @@ impl<'a> PaimonTableRead<'a> {
         if splits.is_empty() {
             return Ok(Box::pin(futures::stream::empty()));
         }
-        for split in splits {
-            if split
-                .data_deletion_files()
-                .is_some_and(|files| files.iter().any(|file| file.is_some()))
-            {
-                return Err(crate::Error::Unsupported {
-                    message: "Batch incremental Diff does not support deletion 
vectors".to_string(),
-                });
-            }
-        }
         let budget = self.parquet_read_budget()?;
         let parquet_read_budget = budget
             .has_resources()
@@ -1454,9 +1445,17 @@ fn primary_key_indices(table: &Table, read_type: 
&[DataField]) -> crate::Result<
     Ok(indices)
 }
 
-fn ensure_diff_supported_read_type(read_type: &[DataField]) -> 
crate::Result<()> {
-    for field in read_type {
-        if !is_diff_supported_type(field.data_type()) {
+fn ensure_diff_supported_read_type(
+    read_type: &[DataField],
+    key_indices: &[usize],
+) -> crate::Result<()> {
+    for (index, field) in read_type.iter().enumerate() {
+        let supported = if key_indices.contains(&index) {
+            is_diff_supported_scalar_type(field.data_type())
+        } else {
+            is_diff_supported_value_type(field.data_type())
+        };
+        if !supported {
             return Err(crate::Error::Unsupported {
                 message: format!(
                     "Batch incremental Diff does not support column '{}' of 
type {:?}",
@@ -1469,7 +1468,24 @@ fn ensure_diff_supported_read_type(read_type: 
&[DataField]) -> crate::Result<()>
     Ok(())
 }
 
-fn is_diff_supported_type(data_type: &DataType) -> bool {
+fn is_diff_supported_value_type(data_type: &DataType) -> bool {
+    match data_type {
+        DataType::Row(row) => row
+            .fields()
+            .iter()
+            .all(|field| is_diff_supported_value_type(field.data_type())),
+        DataType::Array(array) => 
is_diff_supported_value_type(array.element_type()),
+        DataType::Map(map) => {
+            is_diff_supported_value_type(map.key_type())
+                && is_diff_supported_value_type(map.value_type())
+        }
+        DataType::Multiset(multiset) => 
is_diff_supported_value_type(multiset.element_type()),
+        DataType::Vector(vector) => 
is_diff_supported_scalar_type(vector.element_type()),
+        _ => is_diff_supported_scalar_type(data_type),
+    }
+}
+
+fn is_diff_supported_scalar_type(data_type: &DataType) -> bool {
     matches!(
         data_type,
         DataType::Boolean(_)
@@ -1543,19 +1559,123 @@ fn rows_equal_at(
     indices: &[usize],
 ) -> crate::Result<bool> {
     for &idx in indices {
-        let ord = scalar_compare(
+        if !value_equal_at(
             left_batch.column(idx),
             left_row,
             right_batch.column(idx),
             right_row,
-        )?;
-        if ord != Ordering::Equal {
+        )? {
             return Ok(false);
         }
     }
     Ok(true)
 }
 
+/// Compare logical values, ignoring unused children beneath NULL parents.
+/// Scalar comparisons also preserve Java's NaN and signed-zero behavior.
+fn value_equal_at(
+    left: &dyn Array,
+    left_row: usize,
+    right: &dyn Array,
+    right_row: usize,
+) -> crate::Result<bool> {
+    match (left.is_null(left_row), right.is_null(right_row)) {
+        (true, true) => return Ok(true),
+        (true, false) | (false, true) => return Ok(false),
+        (false, false) => {}
+    }
+
+    if let (Some(left), Some(right)) = (
+        left.as_any().downcast_ref::<StructArray>(),
+        right.as_any().downcast_ref::<StructArray>(),
+    ) {
+        if left.num_columns() != right.num_columns() {
+            return Ok(false);
+        }
+        for (left_child, right_child) in 
left.columns().iter().zip(right.columns()) {
+            if !value_equal_at(
+                left_child.as_ref(),
+                left_row,
+                right_child.as_ref(),
+                right_row,
+            )? {
+                return Ok(false);
+            }
+        }
+        return Ok(true);
+    }
+
+    if let (Some(left), Some(right)) = (
+        left.as_any().downcast_ref::<ListArray>(),
+        right.as_any().downcast_ref::<ListArray>(),
+    ) {
+        let left_offsets = left.value_offsets();
+        let right_offsets = right.value_offsets();
+        let (left_start, left_end) = (left_offsets[left_row], 
left_offsets[left_row + 1]);
+        let (right_start, right_end) = (right_offsets[right_row], 
right_offsets[right_row + 1]);
+        if left_end - left_start != right_end - right_start {
+            return Ok(false);
+        }
+        for offset in 0..(left_end - left_start) {
+            if !value_equal_at(
+                left.values().as_ref(),
+                (left_start + offset) as usize,
+                right.values().as_ref(),
+                (right_start + offset) as usize,
+            )? {
+                return Ok(false);
+            }
+        }
+        return Ok(true);
+    }
+
+    if let (Some(left), Some(right)) = (
+        left.as_any().downcast_ref::<MapArray>(),
+        right.as_any().downcast_ref::<MapArray>(),
+    ) {
+        let left_offsets = left.value_offsets();
+        let right_offsets = right.value_offsets();
+        let (left_start, left_end) = (left_offsets[left_row], 
left_offsets[left_row + 1]);
+        let (right_start, right_end) = (right_offsets[right_row], 
right_offsets[right_row + 1]);
+        if left_end - left_start != right_end - right_start {
+            return Ok(false);
+        }
+        for offset in 0..(left_end - left_start) {
+            if !value_equal_at(
+                left.entries(),
+                (left_start + offset) as usize,
+                right.entries(),
+                (right_start + offset) as usize,
+            )? {
+                return Ok(false);
+            }
+        }
+        return Ok(true);
+    }
+
+    if let (Some(left), Some(right)) = (
+        left.as_any().downcast_ref::<FixedSizeListArray>(),
+        right.as_any().downcast_ref::<FixedSizeListArray>(),
+    ) {
+        if left.value_length() != right.value_length() {
+            return Ok(false);
+        }
+        for offset in 0..left.value_length() {
+            if !value_equal_at(
+                left.values().as_ref(),
+                (left.value_offset(left_row) + offset) as usize,
+                right.values().as_ref(),
+                (right.value_offset(right_row) + offset) as usize,
+            )? {
+                return Ok(false);
+            }
+        }
+        return Ok(true);
+    }
+
+    Ok(scalar_compare(left, left_row, right, right_row)? == Ordering::Equal)
+}
+
 fn scalar_compare(
     left: &dyn Array,
     left_row: usize,
@@ -2047,8 +2167,11 @@ mod tests {
     }
 
     #[test]
-    fn test_diff_accepts_scalar_types_and_rejects_nested() {
-        use crate::spec::{ArrayType, BinaryType, DecimalType, IntType, 
TimeType, TimestampType};
+    fn test_diff_accepts_nested_values_but_not_nested_keys() {
+        use crate::spec::{
+            ArrayType, BinaryType, BlobType, DecimalType, IntType, MapType, 
MultisetType, RowType,
+            TimeType, TimestampType, VarCharType,
+        };
 
         let decimal = DataField::new(
             1,
@@ -2075,11 +2198,101 @@ mod tests {
             "time".to_string(),
             DataType::Time(TimeType::new(3).unwrap()),
         );
-        ensure_diff_supported_read_type(&[decimal, timestamp, binary, 
time]).unwrap();
+        ensure_diff_supported_read_type(&[decimal, timestamp, binary, time], 
&[0]).unwrap();
+        ensure_diff_supported_read_type(std::slice::from_ref(&nested), 
&[]).unwrap();
         assert!(matches!(
-            ensure_diff_supported_read_type(&[nested]),
+            ensure_diff_supported_read_type(&[nested], &[0]),
             Err(crate::Error::Unsupported { message }) if 
message.contains("tags")
         ));
+
+        let complex = DataField::new(
+            6,
+            "complex".to_string(),
+            DataType::Array(ArrayType::new(DataType::Row(RowType::new(vec![
+                DataField::new(
+                    7,
+                    "counts".to_string(),
+                    DataType::Map(MapType::new(
+                        DataType::VarChar(VarCharType::string_type()),
+                        
DataType::Multiset(MultisetType::new(DataType::Int(IntType::new()))),
+                    )),
+                ),
+            ])))),
+        );
+        ensure_diff_supported_read_type(&[complex], &[]).unwrap();
+
+        let unsupported = DataField::new(
+            8,
+            "blobs".to_string(),
+            DataType::Array(ArrayType::new(DataType::Blob(BlobType::new()))),
+        );
+        assert!(matches!(
+            ensure_diff_supported_read_type(&[unsupported], &[]),
+            Err(crate::Error::Unsupported { message }) if 
message.contains("blobs")
+        ));
+    }
+
+    #[test]
+    fn test_diff_nested_struct_ignores_children_of_null_parent() {
+        use arrow_array::{Int32Array, StructArray};
+        use arrow_buffer::NullBuffer;
+        use arrow_schema::{DataType as ArrowDataType, Field};
+
+        let field = Arc::new(Field::new("child", ArrowDataType::Int32, true));
+        let left = StructArray::try_new(
+            vec![Arc::clone(&field)].into(),
+            vec![Arc::new(Int32Array::from(vec![Some(10), Some(20), None]))],
+            Some(NullBuffer::from(vec![false, true, true])),
+        )
+        .unwrap();
+        let right = StructArray::try_new(
+            vec![field].into(),
+            vec![Arc::new(Int32Array::from(vec![Some(99), Some(21), None]))],
+            Some(NullBuffer::from(vec![false, true, true])),
+        )
+        .unwrap();
+        assert!(value_equal_at(&left, 0, &right, 0).unwrap());
+        assert!(!value_equal_at(&left, 1, &right, 1).unwrap());
+        assert!(value_equal_at(&left, 2, &right, 2).unwrap());
+        assert!(!value_equal_at(&left, 0, &right, 1).unwrap());
+    }
+
+    #[test]
+    fn test_diff_nested_list_compares_sliced_values_and_null_elements() {
+        use arrow_array::{Float32Array, ListArray};
+        use arrow_buffer::{OffsetBuffer, ScalarBuffer};
+        use arrow_schema::{DataType as ArrowDataType, Field};
+
+        let element = Arc::new(Field::new("element", ArrowDataType::Float32, 
true));
+        let left = ListArray::new(
+            Arc::clone(&element),
+            OffsetBuffer::new(ScalarBuffer::from(vec![0, 1, 3, 4, 5])),
+            Arc::new(Float32Array::from(vec![
+                Some(9.0),
+                Some(f32::NAN),
+                None,
+                Some(-0.0),
+                Some(8.0),
+            ])),
+            None,
+        );
+        let right = ListArray::new(
+            element,
+            OffsetBuffer::new(ScalarBuffer::from(vec![0, 1, 3, 4, 5])),
+            Arc::new(Float32Array::from(vec![
+                Some(7.0),
+                Some(f32::from_bits(0xffc0_0001)),
+                None,
+                Some(0.0),
+                Some(6.0),
+            ])),
+            None,
+        );
+        let left_slice = left.slice(1, 2);
+        let right_slice = right.slice(1, 2);
+        assert!(value_equal_at(&left_slice, 0, &right_slice, 0).unwrap());
+        assert!(!value_equal_at(&left_slice, 1, &right_slice, 1).unwrap());
+        assert!(!value_equal_at(&left, 0, &right, 0).unwrap());
     }
 
     #[test]
diff --git a/crates/paimon/src/table/table_scan.rs 
b/crates/paimon/src/table/table_scan.rs
index d1c614b6..3358203c 100644
--- a/crates/paimon/src/table/table_scan.rs
+++ b/crates/paimon/src/table/table_scan.rs
@@ -2018,13 +2018,6 @@ impl<'a> PaimonTableScan<'a> {
     ) -> crate::Result<(Plan, Plan)> {
         self.ensure_query_auth_allowed()?;
         let core_options = CoreOptions::new(self.table.schema().options());
-        if core_options.deletion_vectors_enabled() {
-            return Err(crate::Error::Unsupported {
-                message:
-                    "Batch incremental Diff does not support tables with 
deletion-vectors.enabled=true"
-                        .to_string(),
-            });
-        }
         // Both forms: without `first_row_id` a filter stays an ordinary data
         // predicate instead of becoming a row range.
         if self.row_ranges.is_some()
@@ -2049,11 +2042,34 @@ impl<'a> PaimonTableScan<'a> {
             .await?;
         let before_entries = 
full_state_scan.plan_manifest_entries(before).await?;
         let after_entries = 
full_state_scan.plan_manifest_entries(after).await?;
+        // Each side is a complete snapshot state. A DV index belongs to that
+        // snapshot, so resolve the two index manifests independently before
+        // building splits. Reusing the endpoint's DVs for the older state can
+        // hide rows that were live at the start of the Diff.
+        let deletion_vectors_needed = core_options.deletion_vectors_enabled();
+        let (before_index_entries, after_index_entries) = futures::try_join!(
+            full_state_scan.read_index_manifest_entries(before, false, 
deletion_vectors_needed),
+            full_state_scan.read_index_manifest_entries(after, false, 
deletion_vectors_needed)
+        )?;
         let before_plan = full_state_scan
-            .plan_snapshot_from_entries(before.clone(), before_entries, None, 
None, None, None)
+            .plan_snapshot_from_entries(
+                before.clone(),
+                before_entries,
+                None,
+                before_index_entries,
+                None,
+                None,
+            )
             .await?;
         let after_plan = full_state_scan
-            .plan_snapshot_from_entries(after.clone(), after_entries, None, 
None, None, None)
+            .plan_snapshot_from_entries(
+                after.clone(),
+                after_entries,
+                None,
+                after_index_entries,
+                None,
+                None,
+            )
             .await?;
         Ok((before_plan, after_plan))
     }
@@ -4799,17 +4815,14 @@ mod tests {
         )
     }
 
-    fn diff_test_table(table_path: &str, deletion_vectors_enabled: bool) -> 
Table {
+    fn diff_test_table(table_path: &str) -> Table {
         let file_io = FileIOBuilder::new("memory").build().unwrap();
-        let mut schema = PaimonSchema::builder()
+        let schema = PaimonSchema::builder()
             .column("id", DataType::Int(IntType::new()))
             .column("value", DataType::Int(IntType::new()))
             .primary_key(["id"])
             .option("bucket", "1")
             .option("merge-engine", "deduplicate");
-        if deletion_vectors_enabled {
-            schema = schema.option("deletion-vectors.enabled", "true");
-        }
         Table::new(
             file_io,
             Identifier::new("test_db", "diff_gate"),
@@ -6214,23 +6227,9 @@ mod tests {
         );
     }
 
-    #[tokio::test]
-    async fn test_diff_rejects_deletion_vectors_enabled() {
-        let table = diff_test_table("memory:/diff_dv_gate", true);
-        let scan = PaimonTableScan::new(&table, None, Vec::new(), None, None, 
None);
-        let before = diff_snapshot(1);
-        let after = diff_snapshot(2);
-
-        let err = scan.plan_snapshot_diff(&before, &after).await.unwrap_err();
-        assert!(
-            matches!(err, crate::Error::Unsupported { ref message } if 
message.contains("deletion-vectors.enabled=true")),
-            "Diff must fail closed on deletion-vector tables"
-        );
-    }
-
     #[tokio::test]
     async fn test_diff_rejects_bucket_rescale_from_snapshot_schemas() {
-        let table = diff_test_table("memory:/diff_bucket_rescale_gate", false);
+        let table = diff_test_table("memory:/diff_bucket_rescale_gate");
         let before_schema = table.schema().clone();
         let after_schema = before_schema
             .apply_changes(vec![crate::spec::SchemaChange::set_option(
@@ -6254,7 +6253,7 @@ mod tests {
 
     #[tokio::test]
     async fn test_diff_rejects_row_id_filters() {
-        let table = diff_test_table("memory:/diff_row_id_gate", false);
+        let table = diff_test_table("memory:/diff_row_id_gate");
         let mut builder = table.new_read_builder();
         let filter = Predicate::Leaf {
             column: crate::spec::ROW_ID_FIELD_NAME.to_string(),
diff --git a/crates/paimon/tests/incremental_diff_extended_test.rs 
b/crates/paimon/tests/incremental_diff_extended_test.rs
new file mode 100644
index 00000000..cf67ce42
--- /dev/null
+++ b/crates/paimon/tests/incremental_diff_extended_test.rs
@@ -0,0 +1,804 @@
+// 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.
+
+mod common;
+
+use std::sync::Arc;
+
+use arrow_array::{
+    Array, ArrayRef, Int32Array, ListArray, MapArray, RecordBatch, 
StringArray, StructArray,
+};
+use arrow_buffer::{NullBuffer, OffsetBuffer, ScalarBuffer};
+use arrow_schema::DataType as ArrowDataType;
+use bytes::Bytes;
+use futures::TryStreamExt;
+use indexmap::IndexMap;
+use paimon::arrow::build_target_arrow_schema;
+use paimon::spec::{
+    ArrayType, DataField, DataType, DeletionVectorMeta, FileKind, 
IndexFileMeta, IndexManifest,
+    IndexManifestEntry, IntType, MapType, MultisetType, RowType, Schema, 
TableSchema, VarCharType,
+};
+use paimon::table::{IncrementalPlan, IncrementalScanMode, IncrementalSplit, 
Table};
+use roaring::RoaringBitmap;
+
+use common::incremental_helpers::{
+    make_batch, memory_table, persist_table_schema, pk_schema, setup_dirs, 
write_batch,
+};
+
+type StringMap = Option<Vec<(&'static str, Option<i32>)>>;
+type StringBag = Option<Vec<(&'static str, i32)>>;
+
+#[derive(Clone)]
+struct NestedRow {
+    id: i32,
+    payload: Option<Option<i32>>,
+    items: Option<Vec<Option<i32>>>,
+    attributes: StringMap,
+    bag: StringBag,
+}
+
+impl NestedRow {
+    fn new(id: i32) -> Self {
+        Self {
+            id,
+            payload: None,
+            items: None,
+            attributes: None,
+            bag: None,
+        }
+    }
+
+    fn payload(mut self, value: Option<i32>) -> Self {
+        self.payload = Some(value);
+        self
+    }
+
+    fn items(mut self, values: Vec<Option<i32>>) -> Self {
+        self.items = Some(values);
+        self
+    }
+
+    fn attributes(mut self, values: Vec<(&'static str, Option<i32>)>) -> Self {
+        self.attributes = Some(values);
+        self
+    }
+
+    fn bag(mut self, values: Vec<(&'static str, i32)>) -> Self {
+        self.bag = Some(values);
+        self
+    }
+}
+
+fn nested_schema() -> TableSchema {
+    let schema = Schema::builder()
+        .column("id", DataType::Int(IntType::new()))
+        .column(
+            "payload",
+            DataType::Row(RowType::new(vec![DataField::new(
+                2,
+                "child".to_string(),
+                DataType::Int(IntType::new()),
+            )])),
+        )
+        .column(
+            "items",
+            DataType::Array(ArrayType::new(DataType::Int(IntType::new()))),
+        )
+        .column(
+            "attributes",
+            DataType::Map(MapType::new(
+                DataType::VarChar(VarCharType::string_type()),
+                DataType::Int(IntType::new()),
+            )),
+        )
+        .column(
+            "bag",
+            DataType::Multiset(MultisetType::new(DataType::VarChar(
+                VarCharType::string_type(),
+            ))),
+        )
+        .primary_key(["id"])
+        .option("bucket", "1")
+        .option("merge-engine", "deduplicate")
+        .build()
+        .unwrap();
+    TableSchema::new(0, &schema)
+}
+
+fn map_array(arrow_type: &ArrowDataType, values: &[StringMap]) -> MapArray {
+    let ArrowDataType::Map(entries_field, sorted) = arrow_type else {
+        panic!("expected Arrow Map")
+    };
+    let ArrowDataType::Struct(entry_fields) = entries_field.data_type() else {
+        panic!("expected Map entries struct")
+    };
+    let mut offsets = vec![0_i32];
+    let mut keys = Vec::new();
+    let mut items = Vec::new();
+    let mut valid = Vec::new();
+    for map in values {
+        valid.push(map.is_some());
+        if let Some(map) = map {
+            for (key, value) in map {
+                keys.push(*key);
+                items.push(*value);
+            }
+        }
+        offsets.push(keys.len() as i32);
+    }
+    let entries = StructArray::try_new(
+        entry_fields.clone(),
+        vec![
+            Arc::new(StringArray::from(keys)) as ArrayRef,
+            Arc::new(Int32Array::from(items)) as ArrayRef,
+        ],
+        None,
+    )
+    .unwrap();
+    MapArray::try_new(
+        Arc::clone(entries_field),
+        OffsetBuffer::new(ScalarBuffer::from(offsets)),
+        entries,
+        Some(NullBuffer::from(valid)),
+        *sorted,
+    )
+    .unwrap()
+}
+
+fn nested_batch(table: &Table, rows: &[NestedRow]) -> RecordBatch {
+    let schema = build_target_arrow_schema(table.schema().fields()).unwrap();
+    let ArrowDataType::Struct(payload_fields) = schema.field(1).data_type() 
else {
+        panic!("expected payload struct")
+    };
+    let child = Int32Array::from(
+        rows.iter()
+            .map(|row| row.payload.flatten())
+            .collect::<Vec<_>>(),
+    );
+    let payload = StructArray::try_new(
+        payload_fields.clone(),
+        vec![Arc::new(child)],
+        Some(NullBuffer::from(
+            rows.iter()
+                .map(|row| row.payload.is_some())
+                .collect::<Vec<_>>(),
+        )),
+    )
+    .unwrap();
+
+    let ArrowDataType::List(element_field) = schema.field(2).data_type() else {
+        panic!("expected items list")
+    };
+    let mut offsets = vec![0_i32];
+    let mut elements: Vec<Option<i32>> = Vec::new();
+    let mut list_valid = Vec::new();
+    for row in rows {
+        list_valid.push(row.items.is_some());
+        if let Some(items) = &row.items {
+            elements.extend(items);
+        }
+        offsets.push(elements.len() as i32);
+    }
+    let items = ListArray::new(
+        Arc::clone(element_field),
+        OffsetBuffer::new(ScalarBuffer::from(offsets)),
+        Arc::new(Int32Array::from(elements)),
+        Some(NullBuffer::from(list_valid)),
+    );
+
+    let attributes = rows
+        .iter()
+        .map(|row| row.attributes.clone())
+        .collect::<Vec<_>>();
+    let bag = rows
+        .iter()
+        .map(|row| {
+            row.bag.as_ref().map(|values| {
+                values
+                    .iter()
+                    .map(|(key, count)| (*key, Some(*count)))
+                    .collect::<Vec<_>>()
+            })
+        })
+        .collect::<Vec<_>>();
+    RecordBatch::try_new(
+        schema.clone(),
+        vec![
+            Arc::new(Int32Array::from(
+                rows.iter().map(|row| row.id).collect::<Vec<_>>(),
+            )),
+            Arc::new(payload),
+            Arc::new(items),
+            Arc::new(map_array(schema.field(3).data_type(), &attributes)),
+            Arc::new(map_array(schema.field(4).data_type(), &bag)),
+        ],
+    )
+    .unwrap()
+}
+
+async fn nested_table(path: &str) -> Table {
+    let (file_io, table) = memory_table(path, nested_schema());
+    setup_dirs(&file_io, path).await;
+    persist_table_schema(&file_io, path, table.schema()).await;
+    table
+}
+
+async fn diff_plan(table: &Table, start: i64, end: i64) -> IncrementalPlan {
+    table
+        .new_read_builder()
+        .new_incremental_scan(IncrementalScanMode::Diff, start, end)
+        .plan()
+        .await
+        .unwrap()
+}
+
+async fn diff_batches(table: &Table, plan: &IncrementalPlan) -> 
Vec<RecordBatch> {
+    table
+        .new_read_builder()
+        .new_read()
+        .unwrap()
+        .to_incremental_arrow(plan)
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap()
+}
+
+fn ids_at(batches: &[RecordBatch], index: usize) -> Vec<i32> {
+    let mut ids = batches
+        .iter()
+        .flat_map(|batch| {
+            batch
+                .column(index)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap()
+                .values()
+                .iter()
+                .copied()
+                .collect::<Vec<_>>()
+        })
+        .collect::<Vec<_>>();
+    ids.sort_unstable();
+    ids
+}
+
+#[tokio::test]
+async fn diff_detects_each_nested_value_type_and_skips_unchanged_rows() {
+    let table = nested_table("memory:/incremental_diff/nested_values").await;
+    let before = vec![
+        NestedRow::new(1).payload(Some(10)),
+        NestedRow::new(2).items(vec![Some(1), None]),
+        NestedRow::new(3).attributes(vec![("a", Some(1))]),
+        NestedRow::new(4).bag(vec![("x", 2)]),
+        NestedRow::new(5)
+            .payload(None)
+            .items(vec![])
+            .attributes(vec![])
+            .bag(vec![]),
+        NestedRow::new(6),
+    ];
+    let after = vec![
+        NestedRow::new(1).payload(Some(11)),
+        NestedRow::new(2).items(vec![Some(1), Some(2)]),
+        NestedRow::new(3).attributes(vec![("a", Some(2))]),
+        NestedRow::new(4).bag(vec![("x", 3)]),
+        before[4].clone(),
+        before[5].clone(),
+        NestedRow::new(7).items(vec![]),
+    ];
+    write_batch(&table, &nested_batch(&table, &before)).await;
+    write_batch(&table, &nested_batch(&table, &after)).await;
+
+    let plan = diff_plan(&table, 1, 2).await;
+    let batches = diff_batches(&table, &plan).await;
+    assert_eq!(ids_at(&batches, 0), vec![1, 2, 3, 4, 7]);
+
+    let audit_batches: Vec<RecordBatch> = table
+        .new_read_builder()
+        .new_read()
+        .unwrap()
+        .to_audit_log_arrow(&plan)
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap();
+    let mut changes = Vec::new();
+    for batch in audit_batches {
+        let kinds = batch
+            .column(0)
+            .as_any()
+            .downcast_ref::<StringArray>()
+            .unwrap();
+        let ids = batch
+            .column(1)
+            .as_any()
+            .downcast_ref::<Int32Array>()
+            .unwrap();
+        for row in 0..batch.num_rows() {
+            changes.push((ids.value(row), kinds.value(row).to_string()));
+        }
+    }
+    changes.sort();
+    assert_eq!(
+        changes,
+        vec![
+            (1, "+U".to_string()),
+            (1, "-U".to_string()),
+            (2, "+U".to_string()),
+            (2, "-U".to_string()),
+            (3, "+U".to_string()),
+            (3, "-U".to_string()),
+            (4, "+U".to_string()),
+            (4, "-U".to_string()),
+            (7, "+I".to_string()),
+        ]
+    );
+}
+
+#[tokio::test]
+async fn diff_nested_null_empty_and_projection_semantics() {
+    let table = nested_table("memory:/incremental_diff/nested_nulls").await;
+    let before = vec![
+        NestedRow::new(1),
+        NestedRow::new(2).payload(None),
+        NestedRow::new(3).items(vec![]),
+        NestedRow::new(4).attributes(vec![]),
+        NestedRow::new(5).bag(vec![]),
+        NestedRow::new(6).items(vec![None]),
+    ];
+    let after = vec![
+        NestedRow::new(1).payload(None),
+        NestedRow::new(2).payload(Some(0)),
+        NestedRow::new(3),
+        NestedRow::new(4),
+        NestedRow::new(5),
+        NestedRow::new(6).items(vec![Some(0)]),
+    ];
+    write_batch(&table, &nested_batch(&table, &before)).await;
+    write_batch(&table, &nested_batch(&table, &after)).await;
+
+    assert_eq!(
+        ids_at(
+            &diff_batches(&table, &diff_plan(&table, 1, 2).await).await,
+            0
+        ),
+        vec![1, 2, 3, 4, 5, 6]
+    );
+
+    let mut builder = table.new_read_builder();
+    builder.with_projection(&["id"]).unwrap();
+    let plan = builder
+        .new_incremental_scan(IncrementalScanMode::Diff, 1, 2)
+        .plan()
+        .await
+        .unwrap();
+    let batches: Vec<RecordBatch> = builder
+        .new_read()
+        .unwrap()
+        .to_incremental_arrow(&plan)
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap();
+    assert_eq!(ids_at(&batches, 0), vec![1, 2, 3, 4, 5, 6]);
+}
+
+/// Install a real bitmap and index manifest for one historical snapshot. The
+/// data writer currently emits L0 files for DV tables; patching only snapshot
+/// index metadata models the materialized, Java-written per-file DV state.
+async fn install_snapshot_dv(
+    table: &Table,
+    snapshot_id: i64,
+    partition: &[u8],
+    bucket: i32,
+    data_file_name: &str,
+    deleted_positions: &[u32],
+) {
+    const MAGIC_NUMBER: i32 = 1581511376;
+    let mut bitmap = RoaringBitmap::new();
+    for position in deleted_positions {
+        bitmap.insert(*position);
+    }
+    let mut bitmap_bytes = Vec::new();
+    bitmap.serialize_into(&mut bitmap_bytes).unwrap();
+    let bitmap_length = 4 + bitmap_bytes.len() as i32;
+    let mut dv_bytes = Vec::new();
+    dv_bytes.extend_from_slice(&bitmap_length.to_be_bytes());
+    dv_bytes.extend_from_slice(&MAGIC_NUMBER.to_be_bytes());
+    dv_bytes.extend_from_slice(&bitmap_bytes);
+    dv_bytes.extend_from_slice(&0_i32.to_be_bytes());
+
+    let dv_name = format!("diff-dv-{snapshot_id}");
+    let dv_path = format!("{}/index/{dv_name}", table.location());
+    table
+        .file_io()
+        .mkdirs(&format!("{}/index", table.location()))
+        .await
+        .unwrap();
+    table
+        .file_io()
+        .new_output(&dv_path)
+        .unwrap()
+        .write(Bytes::from(dv_bytes.clone()))
+        .await
+        .unwrap();
+    let entry = IndexManifestEntry {
+        version: 1,
+        kind: FileKind::Add,
+        partition: partition.to_vec(),
+        bucket,
+        index_file: IndexFileMeta {
+            index_type: "DELETION_VECTORS".to_string(),
+            file_name: dv_name,
+            file_size: dv_bytes.len() as i64,
+            row_count: 1,
+            deletion_vectors_ranges: Some(IndexMap::from([(
+                data_file_name.to_string(),
+                DeletionVectorMeta {
+                    offset: 0,
+                    length: bitmap_length,
+                    cardinality: Some(deleted_positions.len() as i64),
+                },
+            )])),
+            external_path: None,
+            global_index_meta: None,
+        },
+    };
+    let manifest_name = format!("diff-index-{snapshot_id}");
+    IndexManifest::write(
+        table.file_io(),
+        &format!("{}/manifest/{manifest_name}", table.location()),
+        &[entry],
+    )
+    .await
+    .unwrap();
+
+    let manager = table.snapshot_manager();
+    let snapshot = manager.get_snapshot(snapshot_id).await.unwrap();
+    let mut metadata = serde_json::to_value(snapshot).unwrap();
+    metadata["indexManifest"] = serde_json::json!(manifest_name);
+    table
+        .file_io()
+        .new_output(&manager.snapshot_path(snapshot_id))
+        .unwrap()
+        .write(Bytes::from(serde_json::to_vec(&metadata).unwrap()))
+        .await
+        .unwrap();
+}
+
+async fn dv_table(path: &str) -> (Table, Vec<u8>, i32, String) {
+    let (file_io, table) = memory_table(
+        path,
+        pk_schema(&[
+            ("deletion-vectors.enabled", "true"),
+            ("deletion-vectors.merge-on-read", "true"),
+            ("merge-engine", "deduplicate"),
+        ]),
+    );
+    setup_dirs(&file_io, path).await;
+    persist_table_schema(&file_io, path, table.schema()).await;
+    write_batch(&table, &make_batch(vec![1, 2, 3], vec![10, 20, 30])).await;
+    let plan = table.new_read_builder().new_scan().plan().await.unwrap();
+    let split = &plan.splits()[0];
+    assert_eq!(split.data_files().len(), 1);
+    let location = (
+        split.partition().to_serialized_bytes(),
+        split.bucket(),
+        split.data_files()[0].file_name.clone(),
+    );
+    (table, location.0, location.1, location.2)
+}
+
+fn dv_attached(plan: &IncrementalPlan, before: bool) -> bool {
+    plan.splits().iter().any(|split| {
+        let IncrementalSplit::DiffPair {
+            before: before_splits,
+            after: after_splits,
+        } = split
+        else {
+            panic!("Diff plan must contain only split pairs")
+        };
+        let side = if before { before_splits } else { after_splits };
+        side.iter().any(|split| {
+            split
+                .data_deletion_files()
+                .is_some_and(|files| files.iter().any(Option::is_some))
+        })
+    })
+}
+
+#[tokio::test]
+async fn diff_uses_the_deletion_vectors_of_each_snapshot() {
+    let (table, partition, bucket, first_file) =
+        dv_table("memory:/incremental_diff/dv_history").await;
+    write_batch(&table, &make_batch(vec![4], vec![40])).await;
+    write_batch(&table, &make_batch(vec![2, 3], vec![25, 35])).await;
+    install_snapshot_dv(&table, 2, &partition, bucket, &first_file, 
&[1]).await;
+    install_snapshot_dv(&table, 3, &partition, bucket, &first_file, &[1, 
2]).await;
+
+    let first_diff = diff_plan(&table, 1, 2).await;
+    assert!(!dv_attached(&first_diff, true));
+    assert!(dv_attached(&first_diff, false));
+    let first_batches = diff_batches(&table, &first_diff).await;
+    assert_eq!(ids_at(&first_batches, 0), vec![4]);
+
+    let first_audit: Vec<RecordBatch> = table
+        .new_read_builder()
+        .new_read()
+        .unwrap()
+        .to_audit_log_arrow(&first_diff)
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap();
+    let mut first_events = first_audit
+        .iter()
+        .flat_map(|batch| {
+            let kinds = batch
+                .column(0)
+                .as_any()
+                .downcast_ref::<StringArray>()
+                .unwrap();
+            let ids = batch
+                .column(1)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap();
+            (0..batch.num_rows())
+                .map(|row| (ids.value(row), kinds.value(row).to_string()))
+                .collect::<Vec<_>>()
+        })
+        .collect::<Vec<_>>();
+    first_events.sort();
+    assert_eq!(
+        first_events,
+        vec![(2, "-D".to_string()), (4, "+I".to_string())]
+    );
+
+    let second_diff = diff_plan(&table, 2, 3).await;
+    assert!(dv_attached(&second_diff, true));
+    assert!(dv_attached(&second_diff, false));
+    let second_batches = diff_batches(&table, &second_diff).await;
+    assert_eq!(ids_at(&second_batches, 0), vec![2, 3]);
+    let mut values = second_batches
+        .iter()
+        .flat_map(|batch| {
+            let ids = batch
+                .column(0)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap();
+            let values = batch
+                .column(1)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap();
+            (0..batch.num_rows())
+                .map(|row| (ids.value(row), values.value(row)))
+                .collect::<Vec<_>>()
+        })
+        .collect::<Vec<_>>();
+    values.sort();
+    assert_eq!(values, vec![(2, 25), (3, 35)]);
+}
+
+async fn write_level_one_batch(table: &Table, batch: &RecordBatch) {
+    let builder = table.new_write_builder();
+    let mut writer = builder.new_write().unwrap();
+    writer.write_arrow_batch(batch).await.unwrap();
+    let mut messages = writer.prepare_commit().await.unwrap();
+    for message in &mut messages {
+        for file in &mut message.new_files {
+            file.level = 1;
+            file.delete_row_count = Some(0);
+        }
+    }
+    builder.new_commit().commit(messages).await.unwrap();
+}
+
+#[tokio::test]
+async fn diff_reads_materialized_dv_files_without_merge_on_read() {
+    let path = "memory:/incremental_diff/dv_materialized";
+    let (file_io, table) = memory_table(
+        path,
+        pk_schema(&[
+            ("deletion-vectors.enabled", "true"),
+            ("merge-engine", "deduplicate"),
+        ]),
+    );
+    setup_dirs(&file_io, path).await;
+    persist_table_schema(&file_io, path, table.schema()).await;
+    write_level_one_batch(&table, &make_batch(vec![1, 2, 3], vec![10, 20, 
30])).await;
+    let first_plan = table.new_read_builder().new_scan().plan().await.unwrap();
+    let first_split = &first_plan.splits()[0];
+    let partition = first_split.partition().to_serialized_bytes();
+    let bucket = first_split.bucket();
+    let first_file = first_split.data_files()[0].file_name.clone();
+    write_level_one_batch(&table, &make_batch(vec![4], vec![40])).await;
+    install_snapshot_dv(&table, 2, &partition, bucket, &first_file, 
&[1]).await;
+
+    let plan = diff_plan(&table, 1, 2).await;
+    assert!(!dv_attached(&plan, true));
+    assert!(dv_attached(&plan, false));
+    assert_eq!(ids_at(&diff_batches(&table, &plan).await, 0), vec![4]);
+
+    let audit_batches: Vec<RecordBatch> = table
+        .new_read_builder()
+        .new_read()
+        .unwrap()
+        .to_audit_log_arrow(&plan)
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap();
+    let mut events = Vec::new();
+    for batch in audit_batches {
+        let kinds = batch
+            .column(0)
+            .as_any()
+            .downcast_ref::<StringArray>()
+            .unwrap();
+        let ids = batch
+            .column(1)
+            .as_any()
+            .downcast_ref::<Int32Array>()
+            .unwrap();
+        for row in 0..batch.num_rows() {
+            events.push((ids.value(row), kinds.value(row).to_string()));
+        }
+    }
+    events.sort();
+    assert_eq!(events, vec![(2, "-D".to_string()), (4, "+I".to_string())]);
+}
+
+#[tokio::test]
+async fn diff_fails_when_a_historical_deletion_vector_is_missing() {
+    let (table, partition, bucket, first_file) =
+        dv_table("memory:/incremental_diff/dv_missing").await;
+    write_batch(&table, &make_batch(vec![4], vec![40])).await;
+    install_snapshot_dv(&table, 1, &partition, bucket, &first_file, 
&[1]).await;
+    let plan = diff_plan(&table, 1, 2).await;
+    assert!(dv_attached(&plan, true));
+
+    table
+        .file_io()
+        .delete_file(&format!("{}/index/diff-dv-1", table.location()))
+        .await
+        .unwrap();
+    let result = table
+        .new_read_builder()
+        .new_read()
+        .unwrap()
+        .to_incremental_arrow(&plan)
+        .unwrap()
+        .try_collect::<Vec<_>>()
+        .await;
+    assert!(
+        result.is_err(),
+        "a missing historical DV must not be ignored"
+    );
+}
+
+fn vector_batch(table: &Table, rows: &[(i32, Option<[f32; 3]>)]) -> 
RecordBatch {
+    use arrow_array::builder::{FixedSizeListBuilder, Float32Builder};
+
+    let schema = build_target_arrow_schema(table.schema().fields()).unwrap();
+    let ArrowDataType::FixedSizeList(element, length) = 
schema.field(1).data_type() else {
+        panic!("expected VECTOR to map to a fixed-size list")
+    };
+    let mut builder =
+        FixedSizeListBuilder::new(Float32Builder::new(), 
*length).with_field(Arc::clone(element));
+    for (_, vector) in rows {
+        if let Some(vector) = vector {
+            for value in vector {
+                builder.values().append_value(*value);
+            }
+            builder.append(true);
+        } else {
+            for _ in 0..*length {
+                builder.values().append_null();
+            }
+            builder.append(false);
+        }
+    }
+    RecordBatch::try_new(
+        schema,
+        vec![
+            Arc::new(Int32Array::from(
+                rows.iter().map(|(id, _)| *id).collect::<Vec<_>>(),
+            )),
+            Arc::new(builder.finish()),
+        ],
+    )
+    .unwrap()
+}
+
+#[tokio::test]
+async fn diff_compares_vectors_by_elements_and_preserves_float_semantics() {
+    use paimon::spec::{FloatType, VectorType};
+
+    let path = "memory:/incremental_diff/vector_values";
+    let schema = Schema::builder()
+        .column("id", DataType::Int(IntType::new()))
+        .column(
+            "embedding",
+            DataType::Vector(VectorType::new(3, 
DataType::Float(FloatType::new())).unwrap()),
+        )
+        .primary_key(["id"])
+        .option("bucket", "1")
+        .option("merge-engine", "deduplicate")
+        .build()
+        .unwrap();
+    let (file_io, table) = memory_table(path, TableSchema::new(0, &schema));
+    setup_dirs(&file_io, path).await;
+    persist_table_schema(&file_io, path, table.schema()).await;
+
+    let before = [
+        (1, Some([1.0, 2.0, 3.0])),
+        (2, Some([f32::NAN, 2.0, 3.0])),
+        (3, Some([-0.0, 2.0, 3.0])),
+        (4, None),
+        (5, Some([5.0, 5.0, 5.0])),
+    ];
+    let after = [
+        (1, Some([1.0, 2.0, 4.0])),
+        (2, Some([f32::from_bits(0xffc0_0001), 2.0, 3.0])),
+        (3, Some([0.0, 2.0, 3.0])),
+        (4, Some([0.0, 0.0, 0.0])),
+        (5, Some([5.0, 5.0, 5.0])),
+        (6, Some([6.0, 6.0, 6.0])),
+    ];
+    write_batch(&table, &vector_batch(&table, &before)).await;
+    write_batch(&table, &vector_batch(&table, &after)).await;
+
+    let plan = diff_plan(&table, 1, 2).await;
+    assert_eq!(
+        ids_at(&diff_batches(&table, &plan).await, 0),
+        vec![1, 3, 4, 6]
+    );
+    let audit: Vec<RecordBatch> = table
+        .new_read_builder()
+        .new_read()
+        .unwrap()
+        .to_audit_log_arrow(&plan)
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap();
+    let mut events = audit
+        .iter()
+        .flat_map(|batch| {
+            let kinds = batch
+                .column(0)
+                .as_any()
+                .downcast_ref::<StringArray>()
+                .unwrap();
+            let ids = batch
+                .column(1)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap();
+            (0..batch.num_rows())
+                .map(|row| (ids.value(row), kinds.value(row).to_string()))
+                .collect::<Vec<_>>()
+        })
+        .collect::<Vec<_>>();
+    events.sort();
+    assert_eq!(events.len(), 7);
+    assert_eq!(events.iter().filter(|(id, _)| *id == 2).count(), 0);
+    assert_eq!(events.iter().filter(|(id, _)| *id == 6).count(), 1);
+}

Reply via email to