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";