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);
+}