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 dd3bc614 feat: support more scalar types in incremental diff (#940)
dd3bc614 is described below

commit dd3bc614f7642d1731ebd355e3bc0e9bc20a8c3f
Author: Jingsong Lee <[email protected]>
AuthorDate: Thu Sep 24 19:55:27 2026 +0800

    feat: support more scalar types in incremental diff (#940)
---
 crates/paimon/src/table/table_read.rs              |  65 +++++++++--
 crates/paimon/tests/incremental_batch_scan_test.rs | 123 +++++++++++++++++++++
 2 files changed, 176 insertions(+), 12 deletions(-)

diff --git a/crates/paimon/src/table/table_read.rs 
b/crates/paimon/src/table/table_read.rs
index 3313df0e..d783ee00 100644
--- a/crates/paimon/src/table/table_read.rs
+++ b/crates/paimon/src/table/table_read.rs
@@ -1482,6 +1482,12 @@ fn is_diff_supported_type(data_type: &DataType) -> bool {
             | DataType::Char(_)
             | DataType::VarChar(_)
             | DataType::Date(_)
+            | DataType::Binary(_)
+            | DataType::VarBinary(_)
+            | DataType::Decimal(_)
+            | DataType::Time(_)
+            | DataType::Timestamp(_)
+            | DataType::LocalZonedTimestamp(_)
     )
 }
 
@@ -1557,8 +1563,11 @@ fn scalar_compare(
     right_row: usize,
 ) -> crate::Result<Ordering> {
     use arrow_array::{
-        BooleanArray, Date32Array, Float32Array, Float64Array, Int16Array, 
Int32Array, Int64Array,
-        Int8Array, StringArray, UInt16Array, UInt32Array, UInt64Array, 
UInt8Array,
+        BinaryArray, BooleanArray, Date32Array, Decimal128Array, Float32Array, 
Float64Array,
+        Int16Array, Int32Array, Int64Array, Int8Array, StringArray, 
Time32MillisecondArray,
+        Time32SecondArray, Time64MicrosecondArray, Time64NanosecondArray,
+        TimestampMicrosecondArray, TimestampMillisecondArray, 
TimestampNanosecondArray,
+        TimestampSecondArray, UInt16Array, UInt32Array, UInt64Array, 
UInt8Array,
     };
 
     match (left.is_null(left_row), right.is_null(right_row)) {
@@ -1589,6 +1598,35 @@ fn scalar_compare(
     compare!(UInt64Array, |a: &UInt64Array, r| a.value(r));
     compare!(BooleanArray, |a: &BooleanArray, r| a.value(r));
     compare!(Date32Array, |a: &Date32Array, r| a.value(r));
+    compare!(Decimal128Array, |a: &Decimal128Array, r| a.value(r));
+    compare!(Time32SecondArray, |a: &Time32SecondArray, r| a.value(r));
+    compare!(Time32MillisecondArray, |a: &Time32MillisecondArray, r| a
+        .value(r));
+    compare!(Time64MicrosecondArray, |a: &Time64MicrosecondArray, r| a
+        .value(r));
+    compare!(Time64NanosecondArray, |a: &Time64NanosecondArray, r| a
+        .value(r));
+    compare!(TimestampSecondArray, |a: &TimestampSecondArray, r| a
+        .value(r));
+    compare!(
+        TimestampMillisecondArray,
+        |a: &TimestampMillisecondArray, r| a.value(r)
+    );
+    compare!(
+        TimestampMicrosecondArray,
+        |a: &TimestampMicrosecondArray, r| a.value(r)
+    );
+    compare!(
+        TimestampNanosecondArray,
+        |a: &TimestampNanosecondArray, r| a.value(r)
+    );
+
+    if let (Some(a), Some(b)) = (
+        left.as_any().downcast_ref::<BinaryArray>(),
+        right.as_any().downcast_ref::<BinaryArray>(),
+    ) {
+        return Ok(a.value(left_row).cmp(b.value(right_row)));
+    }
 
     if let (Some(a), Some(b)) = (
         left.as_any().downcast_ref::<StringArray>(),
@@ -2009,8 +2047,8 @@ mod tests {
     }
 
     #[test]
-    fn test_diff_rejects_types_without_comparator_support() {
-        use crate::spec::{ArrayType, DecimalType, IntType, TimestampType};
+    fn test_diff_accepts_scalar_types_and_rejects_nested() {
+        use crate::spec::{ArrayType, BinaryType, DecimalType, IntType, 
TimeType, TimestampType};
 
         let decimal = DataField::new(
             1,
@@ -2027,18 +2065,21 @@ mod tests {
             "created_at".to_string(),
             DataType::Timestamp(TimestampType::new(6).unwrap()),
         );
-        assert!(matches!(
-            ensure_diff_supported_read_type(&[decimal]),
-            Err(crate::Error::Unsupported { message }) if 
message.contains("amount")
-        ));
+        let binary = DataField::new(
+            4,
+            "payload".to_string(),
+            DataType::Binary(BinaryType::new(8).unwrap()),
+        );
+        let time = DataField::new(
+            5,
+            "time".to_string(),
+            DataType::Time(TimeType::new(3).unwrap()),
+        );
+        ensure_diff_supported_read_type(&[decimal, timestamp, binary, 
time]).unwrap();
         assert!(matches!(
             ensure_diff_supported_read_type(&[nested]),
             Err(crate::Error::Unsupported { message }) if 
message.contains("tags")
         ));
-        assert!(matches!(
-            ensure_diff_supported_read_type(&[timestamp]),
-            Err(crate::Error::Unsupported { message }) if 
message.contains("created_at")
-        ));
     }
 
     #[test]
diff --git a/crates/paimon/tests/incremental_batch_scan_test.rs 
b/crates/paimon/tests/incremental_batch_scan_test.rs
index 4e11ada3..193afb12 100644
--- a/crates/paimon/tests/incremental_batch_scan_test.rs
+++ b/crates/paimon/tests/incremental_batch_scan_test.rs
@@ -566,6 +566,129 @@ async fn 
diff_identical_rows_are_skipped_from_after_image() {
     assert_eq!(rows, vec![(2, 20)]);
 }
 
+#[tokio::test]
+async fn diff_detects_binary_decimal_and_temporal_changes() {
+    use arrow_array::{
+        BinaryArray, Decimal128Array, Time32MillisecondArray, 
TimestampMicrosecondArray,
+    };
+    use arrow_schema::{DataType as ArrowDataType, Field, Schema as 
ArrowSchema, TimeUnit};
+    use paimon::spec::{
+        DataType, DecimalType, IntType, Schema, TableSchema, TimeType, 
TimestampType, VarBinaryType,
+    };
+    use std::sync::Arc;
+
+    let table_path = "memory:/incremental_batch/diff_scalar_types";
+    let schema = Schema::builder()
+        .column("id", DataType::Int(IntType::new()))
+        .column(
+            "payload",
+            DataType::VarBinary(VarBinaryType::new(8).unwrap()),
+        )
+        .column(
+            "amount",
+            DataType::Decimal(DecimalType::new(10, 2).unwrap()),
+        )
+        .column("event_time", DataType::Time(TimeType::new(3).unwrap()))
+        .column(
+            "created_at",
+            DataType::Timestamp(TimestampType::new(6).unwrap()),
+        )
+        .primary_key(["id"])
+        .option("bucket", "1")
+        .build()
+        .unwrap();
+    let (file_io, table) = memory_table(table_path, TableSchema::new(0, 
&schema));
+    setup_dirs(&file_io, table_path).await;
+    persist_table_schema(&file_io, table_path, table.schema()).await;
+
+    let make_batch = |ids: Vec<i32>,
+                      payloads: Vec<&[u8]>,
+                      amounts: Vec<i128>,
+                      times: Vec<i32>,
+                      timestamps: Vec<i64>| {
+        RecordBatch::try_new(
+            Arc::new(ArrowSchema::new(vec![
+                Field::new("id", ArrowDataType::Int32, false),
+                Field::new("payload", ArrowDataType::Binary, false),
+                Field::new("amount", ArrowDataType::Decimal128(10, 2), false),
+                Field::new(
+                    "event_time",
+                    ArrowDataType::Time32(TimeUnit::Millisecond),
+                    false,
+                ),
+                Field::new(
+                    "created_at",
+                    ArrowDataType::Timestamp(TimeUnit::Microsecond, None),
+                    false,
+                ),
+            ])),
+            vec![
+                Arc::new(Int32Array::from(ids)),
+                Arc::new(BinaryArray::from_iter_values(payloads)),
+                Arc::new(
+                    Decimal128Array::from(amounts)
+                        .with_precision_and_scale(10, 2)
+                        .unwrap(),
+                ),
+                Arc::new(Time32MillisecondArray::from(times)),
+                Arc::new(TimestampMicrosecondArray::from(timestamps)),
+            ],
+        )
+        .unwrap()
+    };
+
+    write_batch(
+        &table,
+        &make_batch(
+            vec![1, 2, 3, 4, 5],
+            vec![b"a", b"b", b"c", b"d", b"e"],
+            vec![100, 200, 300, 400, 500],
+            vec![1, 2, 3, 4, 5],
+            vec![1000, 2000, 3000, 4000, 5000],
+        ),
+    )
+    .await;
+    write_batch(
+        &table,
+        &make_batch(
+            vec![1, 2, 3, 4, 5, 6],
+            vec![b"z", b"b", b"c", b"d", b"e", b"f"],
+            vec![100, 250, 300, 400, 500, 600],
+            vec![1, 2, 3, 9, 5, 6],
+            vec![1000, 2000, 9000, 4000, 5000, 6000],
+        ),
+    )
+    .await;
+
+    let builder = table.new_read_builder();
+    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();
+    let mut ids: Vec<i32> = batches
+        .iter()
+        .flat_map(|batch| {
+            let ids = batch
+                .column(0)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap();
+            ids.values().iter().copied().collect::<Vec<_>>()
+        })
+        .collect();
+    ids.sort_unstable();
+    assert_eq!(ids, vec![1, 2, 3, 4, 6]);
+}
+
 #[tokio::test]
 async fn diff_projection_without_primary_key_still_compares_full_rows() {
     let table_path = "memory:/incremental_batch/diff_projection_without_pk";

Reply via email to