From a71b75560b16c986c442d092d4ee8617a2f30170 Mon Sep 17 00:00:00 2001 From: pandas886 <123344357+pandas886@users.noreply.github.com> Date: Sat, 11 Jul 2026 22:37:41 +0800 Subject: [PATCH 1/6] feat: support batch incremental diff reads --- crates/paimon/src/spec/core_options.rs | 28 + crates/paimon/src/table/audit_log_table.rs | 2 +- crates/paimon/src/table/incremental_scan.rs | 57 +- crates/paimon/src/table/table_read.rs | 755 +++++++++++++++++- crates/paimon/src/table/table_scan.rs | 126 +++ crates/paimon/tests/audit_log_table_test.rs | 278 ++++++- .../tests/incremental_batch_scan_test.rs | 111 ++- 7 files changed, 1323 insertions(+), 34 deletions(-) diff --git a/crates/paimon/src/spec/core_options.rs b/crates/paimon/src/spec/core_options.rs index fce04169..ce01b634 100644 --- a/crates/paimon/src/spec/core_options.rs +++ b/crates/paimon/src/spec/core_options.rs @@ -89,6 +89,8 @@ const IGNORE_DELETE_FALLBACK_KEYS: &[&str] = &[ "deduplicate.ignore-delete", "partial-update.ignore-delete", ]; +const DIFF_PARALLELISM_OPTION: &str = "diff.parallelism"; +const DEFAULT_DIFF_PARALLELISM: usize = 4; const DEFAULT_COMMIT_MAX_RETRIES: u32 = 10; const DEFAULT_COMMIT_TIMEOUT_MS: u64 = 120_000; const DEFAULT_COMMIT_MIN_RETRY_WAIT_MS: u64 = 1_000; @@ -512,6 +514,17 @@ impl<'a> CoreOptions<'a> { .is_some_and(|v| v.eq_ignore_ascii_case("true")) } + /// Parallelism for batch incremental Diff pair reads (`diff.parallelism`). + /// + /// Default is 4; values below 1 are clamped to 1. + pub fn diff_parallelism(&self) -> usize { + self.options + .get(DIFF_PARALLELISM_OPTION) + .and_then(|s| s.parse().ok()) + .unwrap_or(DEFAULT_DIFF_PARALLELISM) + .max(1) + } + pub fn data_evolution_enabled(&self) -> bool { self.options .get(DATA_EVOLUTION_ENABLED_OPTION) @@ -1622,6 +1635,21 @@ mod tests { ); } + #[test] + fn test_diff_parallelism_defaults() { + let options = HashMap::new(); + let core = CoreOptions::new(&options); + assert_eq!(core.diff_parallelism(), 4); + + let options = HashMap::from([(DIFF_PARALLELISM_OPTION.to_string(), "0".into())]); + let core = CoreOptions::new(&options); + assert_eq!(core.diff_parallelism(), 1); + + let options = HashMap::from([(DIFF_PARALLELISM_OPTION.to_string(), "8".into())]); + let core = CoreOptions::new(&options); + assert_eq!(core.diff_parallelism(), 8); + } + #[test] fn test_changelog_producer_accepts_known_values() { for (value, expected) in [ diff --git a/crates/paimon/src/table/audit_log_table.rs b/crates/paimon/src/table/audit_log_table.rs index e5ce7572..f048e9b1 100644 --- a/crates/paimon/src/table/audit_log_table.rs +++ b/crates/paimon/src/table/audit_log_table.rs @@ -27,7 +27,7 @@ use crate::spec::{ /// Incremental reads produce: /// - Delta: primary-key rows use physical `_VALUE_KIND`; append rows are `+I` /// - Changelog: kinds come from physical `_VALUE_KIND` (`+I`/`-U`/`+U`/`-D`) -/// - Diff: not implemented in this release +/// - Diff: before/after image comparison (`+I`/`-U`/`+U`/`-D`, equal keys skipped) #[derive(Debug, Clone)] pub struct AuditLogTable { wrapped: Table, diff --git a/crates/paimon/src/table/incremental_scan.rs b/crates/paimon/src/table/incremental_scan.rs index 3aa169e9..ea4b1fbe 100644 --- a/crates/paimon/src/table/incremental_scan.rs +++ b/crates/paimon/src/table/incremental_scan.rs @@ -35,10 +35,11 @@ pub enum IncrementalScanMode { /// Resolve to [`Delta`](Self::Delta) when `changelog-producer=none`, /// otherwise to [`Changelog`](Self::Changelog). Auto, - /// Diff before/after snapshots. + /// Diff before/after snapshot states for PK tables. /// - /// Not fully implemented in this release; planning returns - /// [`Error::Unsupported`](crate::Error::Unsupported). + /// Phase 1 supports only `merge-engine=deduplicate`. Planning compares the + /// full table state at `start_exclusive` vs `end_inclusive` and yields + /// per-(partition, bucket) [`IncrementalSplit::DiffPair`] units. Diff, } @@ -223,9 +224,51 @@ impl<'a> IncrementalScan<'a> { } async fn plan_diff(&self, mode: IncrementalScanMode) -> crate::Result { - let _ = mode; - Err(crate::Error::Unsupported { - message: "Batch incremental Diff scan is not implemented yet".to_string(), - }) + let core_options = CoreOptions::new(self.table.schema().options()); + if core_options.merge_engine()? != crate::spec::MergeEngine::Deduplicate { + return Err(crate::Error::Unsupported { + message: "Batch incremental Diff only supports merge-engine=deduplicate in Phase 1" + .to_string(), + }); + } + let before = self + .snapshot_manager + .get_snapshot(self.start_exclusive) + .await?; + let after = self + .snapshot_manager + .get_snapshot(self.end_inclusive) + .await?; + let (before_plan, after_plan) = self.scan.plan_snapshot_diff(&before, &after).await?; + + use std::collections::BTreeMap; + type PBKey = (Vec, i32); + + let mut before_map: BTreeMap> = BTreeMap::new(); + for split in before_plan.splits() { + let key = (split.partition().to_serialized_bytes(), split.bucket()); + before_map.entry(key).or_default().push(split.clone()); + } + + let mut after_map: BTreeMap> = BTreeMap::new(); + for split in after_plan.splits() { + let key = (split.partition().to_serialized_bytes(), split.bucket()); + after_map.entry(key).or_default().push(split.clone()); + } + + let mut keys: std::collections::BTreeSet = before_map.keys().cloned().collect(); + keys.extend(after_map.keys().cloned()); + + let mut splits = Vec::new(); + for key in keys { + let before = before_map.remove(&key).unwrap_or_default(); + let after = after_map.remove(&key).unwrap_or_default(); + if before.is_empty() && after.is_empty() { + continue; + } + splits.push(IncrementalSplit::DiffPair { before, after }); + } + + Ok(IncrementalPlan::new(mode, splits)) } } diff --git a/crates/paimon/src/table/table_read.rs b/crates/paimon/src/table/table_read.rs index 6c57f6b6..6f78a8af 100644 --- a/crates/paimon/src/table/table_read.rs +++ b/crates/paimon/src/table/table_read.rs @@ -29,9 +29,12 @@ use crate::spec::{ VALUE_KIND_FIELD_ID, VALUE_KIND_FIELD_NAME, }; use crate::DataSplit; -use arrow_array::{Array, ArrayRef, RecordBatch, StringArray}; +use arrow_array::{builder::StringBuilder, Array, ArrayRef, RecordBatch, StringArray, UInt32Array}; use arrow_schema::Schema as ArrowSchema; -use futures::StreamExt; +use arrow_select::concat::concat as arrow_concat; +use arrow_select::take::take; +use futures::{stream, StreamExt}; +use std::cmp::Ordering; use std::sync::Arc; /// Table read: reads data from splits (e.g. produced by [TableScan::plan]). @@ -136,8 +139,8 @@ impl<'a> TableRead<'a> { /// Returns an [`ArrowRecordBatchStream`] for an incremental scan plan. /// - /// Only [`IncrementalSplit::Data`] is supported in this release. Diff - /// planning/read remains unimplemented. + /// Delta/Changelog use [`IncrementalSplit::Data`]. Diff uses + /// [`IncrementalSplit::DiffPair`] and emits after-image rows only. pub fn to_incremental_arrow( &self, plan: &IncrementalPlan, @@ -154,8 +157,8 @@ impl<'a> TableRead<'a> { /// /// Output schema is `rowkind` (+ optional `_SEQUENCE_NUMBER`) followed by /// the projected user columns. Primary-key Delta and Changelog rows take - /// kinds from `_VALUE_KIND`; append-only Delta rows are `+I`. Diff remains - /// unsupported. + /// kinds from `_VALUE_KIND`; append-only Delta rows are `+I`. Diff emits + /// `+I`/`-U`/`+U`/`-D` from before/after image comparison. pub fn to_audit_log_arrow( &self, plan: &IncrementalPlan, @@ -233,9 +236,7 @@ impl<'a> PaimonTableRead<'a> { plan: &IncrementalPlan, ) -> crate::Result { if plan.mode() == IncrementalScanMode::Diff { - return Err(crate::Error::Unsupported { - message: "Batch incremental Diff read not yet implemented".to_string(), - }); + return self.to_incremental_diff_arrow(plan); } let mut data_splits = Vec::new(); @@ -255,15 +256,53 @@ impl<'a> PaimonTableRead<'a> { self.new_data_file_reader()?.read(&data_splits) } + fn to_incremental_diff_arrow( + &self, + plan: &IncrementalPlan, + ) -> crate::Result { + let pairs: Vec<(Vec, Vec)> = plan + .splits() + .iter() + .filter_map(|s| match s { + IncrementalSplit::DiffPair { before, after } => { + Some((before.clone(), after.clone())) + } + _ => None, + }) + .collect(); + let parallel = CoreOptions::new(self.table.schema().options()).diff_parallelism(); + let table = self.table.clone(); + let read_type = self.read_type.clone(); + let data_predicates = self.data_predicates.clone(); + + Ok(Box::pin(async_stream::try_stream! { + let mut workers = stream::iter(pairs.into_iter().map(|(before, after)| { + let table = table.clone(); + let read_type = read_type.clone(); + let data_predicates = data_predicates.clone(); + async move { + let pair_read = + PaimonTableRead::new(&table, read_type, data_predicates); + pair_read.to_diff_after_image_stream(&before, &after) + } + })) + .buffer_unordered(parallel); + while let Some(stream_result) = workers.next().await { + let mut pair_stream = stream_result?; + while let Some(batch) = pair_stream.next().await { + yield batch?; + } + } + })) + } + /// Returns an audit-log stream for a planned incremental scan. pub fn to_audit_log_arrow( &self, plan: &IncrementalPlan, ) -> crate::Result { match plan.mode() { - IncrementalScanMode::Diff => Err(crate::Error::Unsupported { - message: "Batch incremental Diff audit read not yet implemented".to_string(), - }), + IncrementalScanMode::Diff => self.audit_diff_stream(plan), IncrementalScanMode::Delta => { self.audit_raw_stream(plan, !self.table.schema().primary_keys().is_empty()) } @@ -361,6 +400,232 @@ impl<'a> PaimonTableRead<'a> { })) } + fn audit_diff_stream(&self, plan: &IncrementalPlan) -> crate::Result { + let pairs: Vec<(Vec, Vec)> = plan + .splits() + .iter() + .filter_map(|s| match s { + IncrementalSplit::DiffPair { before, after } => { + Some((before.clone(), after.clone())) + } + _ => None, + }) + .collect(); + let parallel = CoreOptions::new(self.table.schema().options()).diff_parallelism(); + let table = self.table.clone(); + let read_type = self.read_type.clone(); + let data_predicates = self.data_predicates.clone(); + + Ok(Box::pin(async_stream::try_stream! { + let mut workers = stream::iter(pairs.into_iter().map(|(before, after)| { + let table = table.clone(); + let read_type = read_type.clone(); + let data_predicates = data_predicates.clone(); + async move { + let pair_read = PaimonTableRead::new(&table, read_type, data_predicates); + pair_read.to_audit_log_arrow_for_diff(&before, &after) + } + })) + .buffer_unordered(parallel); + while let Some(stream_result) = workers.next().await { + let mut pair_stream = stream_result?; + while let Some(batch) = pair_stream.next().await { + yield batch?; + } + } + })) + } + + fn to_audit_log_arrow_for_diff( + &self, + before: &[DataSplit], + after: &[DataSplit], + ) -> crate::Result { + ensure_diff_supported_read_type(&self.read_type)?; + let include_sequence = audit_sequence_number_enabled(self.table); + let audit_schema = audit_schema_for_read_type(&self.read_type, include_sequence)?; + + let mut diff_read_type = self.read_type.clone(); + if include_sequence { + diff_read_type.insert( + 0, + DataField::new( + SEQUENCE_NUMBER_FIELD_ID, + SEQUENCE_NUMBER_FIELD_NAME.to_string(), + DataType::BigInt(BigIntType::new()), + ), + ); + } + + let key_indices = primary_key_indices(self.table, &diff_read_type)?; + let value_indices = value_indices_for_diff(self.table, &diff_read_type); + + let before = before.to_vec(); + let after = after.to_vec(); + let table = self.table.clone(); + let read_type_for_output = self.read_type.clone(); + let data_predicates = self.data_predicates.clone(); + + Ok(Box::pin(async_stream::try_stream! { + let core_options = CoreOptions::new(table.schema().options()); + let pair_read = PaimonTableRead::new(&table, diff_read_type.clone(), data_predicates); + let before_stream = + pair_read.read_pk_sorted_for_diff_with_type(&before, &core_options, &diff_read_type)?; + let after_stream = + pair_read.read_pk_sorted_for_diff_with_type(&after, &core_options, &diff_read_type)?; + let mut bc = ArrowCursor::new(before_stream).await?; + let mut ac = ArrowCursor::new(after_stream).await?; + let mut data_col_indices: Option> = None; + let mut builder = AuditBatchBuilder::new(audit_schema.clone()); + + while bc.alive() || ac.alive() { + let indices = data_col_indices.get_or_insert_with(|| { + let sample = if bc.alive() { + bc.batch() + } else { + ac.batch() + }; + diff_output_col_indices(sample, &read_type_for_output, include_sequence) + .expect("diff output column indices") + }); + if !builder.has_data_columns() { + builder.set_data_col_indices(indices.clone()); + } + match cursor_cmp(&bc, &ac, &key_indices, &value_indices)? { + CursorOrd::BeforeOnly => { + builder.push("-D", bc.batch(), bc.row()); + bc.advance().await?; + } + CursorOrd::AfterOnly => { + builder.push("+I", ac.batch(), ac.row()); + ac.advance().await?; + } + CursorOrd::EqualSame => { + bc.advance().await?; + ac.advance().await?; + } + CursorOrd::EqualDiff => { + builder.push("-U", bc.batch(), bc.row()); + builder.push("+U", ac.batch(), ac.row()); + bc.advance().await?; + ac.advance().await?; + } + } + if builder.len() >= DIFF_BATCH_SIZE { + yield builder.flush()?; + } + } + if builder.len() > 0 { + yield builder.flush()?; + } + })) + } + + fn to_diff_after_image_stream( + &self, + before: &[DataSplit], + after: &[DataSplit], + ) -> crate::Result { + ensure_diff_supported_read_type(&self.read_type)?; + let key_indices = primary_key_indices(self.table, &self.read_type)?; + let value_indices = value_indices_for_diff(self.table, &self.read_type); + let output_schema = build_target_arrow_schema(&self.read_type)?; + let output_col_indices: Vec = (0..self.read_type.len()).collect(); + + let table = self.table.clone(); + let read_type = self.read_type.clone(); + let data_predicates = self.data_predicates.clone(); + let before = before.to_vec(); + let after = after.to_vec(); + + Ok(Box::pin(async_stream::try_stream! { + let core_options = CoreOptions::new(table.schema().options()); + let pair_read = PaimonTableRead::new(&table, read_type, data_predicates); + let before_stream = pair_read.read_pk_sorted_for_diff(&before, &core_options)?; + let after_stream = pair_read.read_pk_sorted_for_diff(&after, &core_options)?; + let mut bc = ArrowCursor::new(before_stream).await?; + let mut ac = ArrowCursor::new(after_stream).await?; + let mut builder = + DiffAfterImageBatchBuilder::new(output_schema.clone(), output_col_indices.clone()); + + while bc.alive() || ac.alive() { + match cursor_cmp(&bc, &ac, &key_indices, &value_indices)? { + CursorOrd::BeforeOnly => { + bc.advance().await?; + } + CursorOrd::AfterOnly => { + builder.push(ac.batch(), ac.row()); + ac.advance().await?; + } + CursorOrd::EqualSame => { + bc.advance().await?; + ac.advance().await?; + } + CursorOrd::EqualDiff => { + builder.push(ac.batch(), ac.row()); + bc.advance().await?; + ac.advance().await?; + } + } + if builder.len() >= DIFF_BATCH_SIZE { + yield builder.flush()?; + } + } + if builder.len() > 0 { + yield builder.flush()?; + } + })) + } + + fn read_pk_sorted_for_diff( + &self, + splits: &[DataSplit], + core_options: &CoreOptions, + ) -> crate::Result { + self.read_pk_sorted_for_diff_with_type(splits, core_options, &self.read_type) + } + + fn read_pk_sorted_for_diff_with_type( + &self, + splits: &[DataSplit], + core_options: &CoreOptions, + read_type: &[DataField], + ) -> crate::Result { + 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 reader = KeyValueFileReader::new( + self.table.file_io.clone(), + KeyValueReadConfig { + table_name: self.table.identifier().full_name(), + table_options: self.table.schema().options().clone(), + schema_manager: self.table.schema_manager().clone(), + table_schema_id: self.table.schema().id(), + table_fields: self.table.schema.fields().to_vec(), + read_type: read_type.to_vec(), + predicates: self.data_predicates.clone(), + primary_keys: self.table.schema.trimmed_primary_keys(), + merge_engine: core_options.merge_engine()?, + sequence_fields: core_options + .sequence_fields() + .iter() + .map(|s| s.to_string()) + .collect(), + }, + ); + reader.read(splits) + } + /// Returns an [`ArrowRecordBatchStream`]. pub fn to_arrow(&self, data_splits: &[DataSplit]) -> crate::Result { let has_primary_keys = !self.table.schema.primary_keys().is_empty(); @@ -609,6 +874,448 @@ fn rowkind_array_from_column(column: &dyn arrow_array::Array) -> crate::Result, + row: usize, +} + +impl ArrowCursor { + async fn new(stream: ArrowRecordBatchStream) -> crate::Result { + let mut cursor = Self { + stream, + batch: None, + row: 0, + }; + cursor.advance().await?; + Ok(cursor) + } + + fn alive(&self) -> bool { + self.batch.is_some() + } + + fn batch(&self) -> &RecordBatch { + self.batch.as_ref().expect("cursor must be alive") + } + + fn row(&self) -> usize { + self.row + } + + async fn advance(&mut self) -> crate::Result<()> { + loop { + if let Some(ref batch) = self.batch { + if self.row + 1 < batch.num_rows() { + self.row += 1; + return Ok(()); + } + } + match self.stream.next().await { + Some(Ok(batch)) if batch.num_rows() > 0 => { + self.batch = Some(batch); + self.row = 0; + return Ok(()); + } + Some(Ok(_)) => continue, + Some(Err(e)) => return Err(e), + None => { + self.batch = None; + return Ok(()); + } + } + } + } +} + +struct AuditBatchBuilder { + schema: Arc, + rowkind: StringBuilder, + row_indices: Vec<(usize, usize)>, + pinned_batches: Vec, + data_col_indices: Vec, + len: usize, +} + +impl AuditBatchBuilder { + fn new(schema: Arc) -> Self { + Self { + schema, + rowkind: StringBuilder::new(), + row_indices: Vec::new(), + pinned_batches: Vec::new(), + data_col_indices: Vec::new(), + len: 0, + } + } + + fn has_data_columns(&self) -> bool { + !self.data_col_indices.is_empty() + } + + fn set_data_col_indices(&mut self, indices: Vec) { + self.data_col_indices = indices; + } + + fn len(&self) -> usize { + self.len + } + + fn push(&mut self, kind: &str, batch: &RecordBatch, row: usize) { + self.rowkind.append_value(kind); + let batch_id = self.pin_batch(batch); + self.row_indices.push((batch_id, row)); + self.len += 1; + } + + fn pin_batch(&mut self, batch: &RecordBatch) -> usize { + if let Some(last) = self.pinned_batches.last() { + if std::ptr::eq(batch, last) { + return self.pinned_batches.len() - 1; + } + } + let batch_id = self.pinned_batches.len(); + self.pinned_batches.push(batch.clone()); + batch_id + } + + fn flush(&mut self) -> crate::Result { + let mut columns: Vec = vec![Arc::new(self.rowkind.finish())]; + self.rowkind = StringBuilder::new(); + for &col_idx in &self.data_col_indices { + let taken: Vec = self + .row_indices + .iter() + .map(|(batch_id, row)| { + take( + self.pinned_batches[*batch_id].column(col_idx).as_ref(), + &UInt32Array::from(vec![*row as u32]), + None, + ) + .map_err(|e| crate::Error::UnexpectedError { + message: format!("Failed to take audit diff column: {e}"), + source: Some(Box::new(e)), + }) + }) + .collect::>>()?; + let refs: Vec<&dyn Array> = taken.iter().map(|array| array.as_ref()).collect(); + columns.push( + arrow_concat(&refs).map_err(|e| crate::Error::UnexpectedError { + message: format!("Failed to concat audit diff column: {e}"), + source: Some(Box::new(e)), + })?, + ); + } + self.row_indices.clear(); + self.pinned_batches.clear(); + self.len = 0; + RecordBatch::try_new(self.schema.clone(), columns).map_err(|e| { + crate::Error::UnexpectedError { + message: format!("Failed to build audit diff batch: {e}"), + source: Some(Box::new(e)), + } + }) + } +} + +struct DiffAfterImageBatchBuilder { + schema: Arc, + row_indices: Vec<(usize, usize)>, + pinned_batches: Vec, + col_indices: Vec, + len: usize, +} + +impl DiffAfterImageBatchBuilder { + fn new(schema: Arc, col_indices: Vec) -> Self { + Self { + schema, + row_indices: Vec::new(), + pinned_batches: Vec::new(), + col_indices, + len: 0, + } + } + + fn len(&self) -> usize { + self.len + } + + fn push(&mut self, batch: &RecordBatch, row: usize) { + let batch_id = self.pin_batch(batch); + self.row_indices.push((batch_id, row)); + self.len += 1; + } + + fn pin_batch(&mut self, batch: &RecordBatch) -> usize { + if let Some(last) = self.pinned_batches.last() { + if std::ptr::eq(batch, last) { + return self.pinned_batches.len() - 1; + } + } + let batch_id = self.pinned_batches.len(); + self.pinned_batches.push(batch.clone()); + batch_id + } + + fn flush(&mut self) -> crate::Result { + let mut columns = Vec::with_capacity(self.col_indices.len()); + for &col_idx in &self.col_indices { + let taken: Vec = self + .row_indices + .iter() + .map(|(batch_id, row)| { + take( + self.pinned_batches[*batch_id].column(col_idx).as_ref(), + &UInt32Array::from(vec![*row as u32]), + None, + ) + .map_err(|e| crate::Error::UnexpectedError { + message: format!("Failed to take diff after-image column: {e}"), + source: Some(Box::new(e)), + }) + }) + .collect::>>()?; + let refs: Vec<&dyn Array> = taken.iter().map(|array| array.as_ref()).collect(); + columns.push( + arrow_concat(&refs).map_err(|e| crate::Error::UnexpectedError { + message: format!("Failed to concat diff after-image column: {e}"), + source: Some(Box::new(e)), + })?, + ); + } + self.row_indices.clear(); + self.pinned_batches.clear(); + self.len = 0; + RecordBatch::try_new(self.schema.clone(), columns).map_err(|e| { + crate::Error::UnexpectedError { + message: format!("Failed to build diff after-image batch: {e}"), + source: Some(Box::new(e)), + } + }) + } +} + +fn diff_output_col_indices( + batch: &RecordBatch, + read_type: &[DataField], + include_sequence: bool, +) -> crate::Result> { + let mut indices = Vec::with_capacity(read_type.len() + usize::from(include_sequence)); + if include_sequence { + indices.push( + batch + .schema() + .index_of(SEQUENCE_NUMBER_FIELD_NAME) + .map_err(|e| crate::Error::DataInvalid { + message: format!("Diff read missing _SEQUENCE_NUMBER: {e}"), + source: None, + })?, + ); + } + for field in read_type { + indices.push(batch.schema().index_of(field.name()).map_err(|e| { + crate::Error::DataInvalid { + message: format!("Diff read missing column '{}': {e}", field.name()), + source: None, + } + })?); + } + Ok(indices) +} + +fn value_indices_for_diff(table: &Table, fields: &[DataField]) -> Vec { + let primary_key_names = table.schema().trimmed_primary_keys(); + let primary_keys: std::collections::HashSet<&str> = + primary_key_names.iter().map(|key| key.as_str()).collect(); + fields + .iter() + .enumerate() + .filter(|(_, field)| { + field.name() != SEQUENCE_NUMBER_FIELD_NAME && !primary_keys.contains(field.name()) + }) + .map(|(index, _)| index) + .collect() +} + +fn primary_key_indices(table: &Table, read_type: &[DataField]) -> crate::Result> { + let mut indices = Vec::new(); + for pk in table.schema().trimmed_primary_keys() { + let idx = read_type + .iter() + .position(|field| field.name() == pk) + .ok_or_else(|| crate::Error::DataInvalid { + message: format!("Primary key column '{pk}' missing from read projection"), + source: None, + })?; + indices.push(idx); + } + 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()) { + return Err(crate::Error::Unsupported { + message: format!( + "Batch incremental Diff does not support nested or decimal column '{}'", + field.name() + ), + }); + } + } + Ok(()) +} + +fn is_diff_supported_type(data_type: &DataType) -> bool { + !matches!( + data_type, + DataType::Decimal(_) | DataType::Array(_) | DataType::Map(_) | DataType::Row(_) + ) +} + +fn cursor_cmp( + bc: &ArrowCursor, + ac: &ArrowCursor, + key_indices: &[usize], + value_indices: &[usize], +) -> crate::Result { + match (bc.alive(), ac.alive()) { + (false, false) => unreachable!("cursor_cmp called with both streams exhausted"), + (false, true) => return Ok(CursorOrd::AfterOnly), + (true, false) => return Ok(CursorOrd::BeforeOnly), + (true, true) => {} + } + match compare_pk(bc, ac, key_indices)? { + Ordering::Less => Ok(CursorOrd::BeforeOnly), + Ordering::Greater => Ok(CursorOrd::AfterOnly), + Ordering::Equal => { + if rows_equal_at(bc.batch(), bc.row(), ac.batch(), ac.row(), value_indices)? { + Ok(CursorOrd::EqualSame) + } else { + Ok(CursorOrd::EqualDiff) + } + } + } +} + +fn compare_pk( + bc: &ArrowCursor, + ac: &ArrowCursor, + key_indices: &[usize], +) -> crate::Result { + for &idx in key_indices { + let ord = scalar_compare( + bc.batch().column(idx), + bc.row(), + ac.batch().column(idx), + ac.row(), + )?; + if ord != Ordering::Equal { + return Ok(ord); + } + } + Ok(Ordering::Equal) +} + +fn rows_equal_at( + left_batch: &RecordBatch, + left_row: usize, + right_batch: &RecordBatch, + right_row: usize, + indices: &[usize], +) -> crate::Result { + for &idx in indices { + let ord = scalar_compare( + left_batch.column(idx), + left_row, + right_batch.column(idx), + right_row, + )?; + if ord != Ordering::Equal { + return Ok(false); + } + } + Ok(true) +} + +fn scalar_compare( + left: &dyn Array, + left_row: usize, + right: &dyn Array, + right_row: usize, +) -> crate::Result { + use arrow_array::{ + BooleanArray, Date32Array, Float32Array, Float64Array, Int16Array, Int32Array, Int64Array, + Int8Array, StringArray, UInt16Array, UInt32Array, UInt64Array, UInt8Array, + }; + + macro_rules! compare { + ($ty:ty, $getter:expr) => { + if let (Some(a), Some(b)) = ( + left.as_any().downcast_ref::<$ty>(), + right.as_any().downcast_ref::<$ty>(), + ) { + return Ok($getter(a, left_row).cmp(&$getter(b, right_row))); + } + }; + } + + compare!(Int8Array, |a: &Int8Array, r| a.value(r)); + compare!(Int16Array, |a: &Int16Array, r| a.value(r)); + compare!(Int32Array, |a: &Int32Array, r| a.value(r)); + compare!(Int64Array, |a: &Int64Array, r| a.value(r)); + compare!(UInt8Array, |a: &UInt8Array, r| a.value(r)); + compare!(UInt16Array, |a: &UInt16Array, r| a.value(r)); + compare!(UInt32Array, |a: &UInt32Array, r| a.value(r)); + compare!(UInt64Array, |a: &UInt64Array, r| a.value(r)); + compare!(BooleanArray, |a: &BooleanArray, r| a.value(r)); + compare!(Date32Array, |a: &Date32Array, r| a.value(r)); + + if let (Some(a), Some(b)) = ( + left.as_any().downcast_ref::(), + right.as_any().downcast_ref::(), + ) { + return Ok(a.value(left_row).cmp(b.value(right_row))); + } + + if let (Some(a), Some(b)) = ( + left.as_any().downcast_ref::(), + right.as_any().downcast_ref::(), + ) { + return Ok(a + .value(left_row) + .partial_cmp(&b.value(right_row)) + .unwrap_or(Ordering::Equal)); + } + if let (Some(a), Some(b)) = ( + left.as_any().downcast_ref::(), + right.as_any().downcast_ref::(), + ) { + return Ok(a + .value(left_row) + .partial_cmp(&b.value(right_row)) + .unwrap_or(Ordering::Equal)); + } + + Err(crate::Error::Unsupported { + message: format!( + "Batch incremental Diff does not support comparing column type {:?}", + left.data_type() + ), + }) +} + /// Whether a primary-key split must go through the sort-merge reader. /// /// Mirrors Java `PrimaryKeyTableRawFileSplitReadProvider#match`: a raw read @@ -731,4 +1438,28 @@ mod tests { "directly-constructed read of a query-auth.enabled table must fail closed" ); } + + #[test] + fn test_diff_rejects_nested_and_decimal_types() { + use crate::spec::{ArrayType, DecimalType, IntType}; + + let decimal = DataField::new( + 1, + "amount".to_string(), + DataType::Decimal(DecimalType::new(10, 2).unwrap()), + ); + let nested = DataField::new( + 2, + "tags".to_string(), + DataType::Array(ArrayType::new(DataType::Int(IntType::new()))), + ); + assert!(matches!( + ensure_diff_supported_read_type(&[decimal]), + Err(crate::Error::Unsupported { message }) if message.contains("nested or decimal") + )); + assert!(matches!( + ensure_diff_supported_read_type(&[nested]), + Err(crate::Error::Unsupported { message }) if message.contains("nested or decimal") + )); + } } diff --git a/crates/paimon/src/table/table_scan.rs b/crates/paimon/src/table/table_scan.rs index c8c0f199..4c2c0f1a 100644 --- a/crates/paimon/src/table/table_scan.rs +++ b/crates/paimon/src/table/table_scan.rs @@ -987,6 +987,20 @@ impl<'a> TableScan<'a> { } } + /// Plan before/after full-snapshot splits for batch incremental Diff. + pub(crate) async fn plan_snapshot_diff( + &self, + before: &Snapshot, + after: &Snapshot, + ) -> crate::Result<(Plan, Plan)> { + match &self.0 { + TableScanKind::Paimon(scan) => scan.plan_snapshot_diff(before, after).await, + TableScanKind::Format(_) => Err(crate::Error::Unsupported { + message: "Format tables do not support incremental Diff scan".to_string(), + }), + } + } + #[cfg(test)] fn apply_limit_pushdown(&self, splits: Vec) -> Vec { match &self.0 { @@ -1535,6 +1549,118 @@ impl<'a> PaimonTableScan<'a> { .await } + /// Plan before/after full-snapshot states for Diff incremental scan. + /// + /// Loads full manifest entries for both snapshots, rejects bucket rescale, + /// prunes files that are identical and unaffected by exclusive changes, then + /// builds splits via the shared snapshot planning path. + pub(crate) async fn plan_snapshot_diff( + &self, + before: &Snapshot, + after: &Snapshot, + ) -> crate::Result<(Plan, Plan)> { + self.ensure_query_auth_allowed()?; + let mut before_entries = self.plan_manifest_entries(before).await?; + let mut after_entries = self.plan_manifest_entries(after).await?; + Self::validate_diff_bucket_layout(&before_entries, &after_entries)?; + Self::prune_unchanged_diff_files(&mut before_entries, &mut after_entries); + let data_evolution_read_field_ids = self.projected_read_field_ids()?; + let before_plan = self + .plan_snapshot_from_entries( + before.clone(), + before_entries, + data_evolution_read_field_ids.as_ref(), + None, + ) + .await?; + let after_plan = self + .plan_snapshot_from_entries( + after.clone(), + after_entries, + data_evolution_read_field_ids.as_ref(), + None, + ) + .await?; + Ok((before_plan, after_plan)) + } + + fn validate_diff_bucket_layout( + before_entries: &[ManifestEntry], + after_entries: &[ManifestEntry], + ) -> crate::Result<()> { + use std::collections::{BTreeMap, BTreeSet}; + + fn totals(entries: &[ManifestEntry]) -> BTreeMap, BTreeSet> { + let mut result: BTreeMap, BTreeSet> = BTreeMap::new(); + for entry in entries { + result + .entry(entry.partition().to_vec()) + .or_default() + .insert(entry.total_buckets()); + } + result + } + + let before = totals(before_entries); + let after = totals(after_entries); + for (partition, before_totals) in &before { + let Some(after_totals) = after.get(partition) else { + continue; + }; + if before_totals != after_totals { + return Err(crate::Error::Unsupported { + message: + "Batch incremental Diff does not support bucket rescale between snapshots" + .to_string(), + }); + } + } + Ok(()) + } + + fn prune_unchanged_diff_files(before: &mut Vec, after: &mut Vec) { + // Common files (byte-equal manifest entries) that do not overlap any + // exclusive before/after file key-range can be dropped from both sides. + let mut common: Vec = before + .iter() + .filter(|entry| after.iter().any(|other| other == *entry)) + .cloned() + .collect(); + if common.is_empty() { + return; + } + + let exclusive_after: Vec<&ManifestEntry> = after + .iter() + .filter(|entry| !common.iter().any(|c| c == *entry)) + .collect(); + let exclusive_before: Vec<&ManifestEntry> = before + .iter() + .filter(|entry| !common.iter().any(|c| c == *entry)) + .collect(); + common.retain(|entry| { + let overlaps_exclusive = |other: &&ManifestEntry| { + Self::key_ranges_overlap_bytes( + &entry.file().min_key, + &entry.file().max_key, + &other.file().min_key, + &other.file().max_key, + ) + }; + !exclusive_after.iter().any(overlaps_exclusive) + && !exclusive_before.iter().any(overlaps_exclusive) + }); + if common.is_empty() { + return; + } + before.retain(|entry| !common.iter().any(|c| c == entry)); + after.retain(|entry| !common.iter().any(|c| c == entry)); + } + + fn key_ranges_overlap_bytes(min_a: &[u8], max_a: &[u8], min_b: &[u8], max_b: &[u8]) -> bool { + min_a <= max_b && min_b <= max_a + } + /// Read entries from a single manifest list (delta or changelog) with /// partition / bucket filter pushdown matching the full scan path. async fn plan_manifest_list_entries( diff --git a/crates/paimon/tests/audit_log_table_test.rs b/crates/paimon/tests/audit_log_table_test.rs index 8839673e..9dc2ca08 100644 --- a/crates/paimon/tests/audit_log_table_test.rs +++ b/crates/paimon/tests/audit_log_table_test.rs @@ -320,9 +320,73 @@ async fn audit_log_exposes_sequence_number_when_enabled() { assert!(rows.iter().all(|(_, seq, _, _)| *seq >= 0)); } +async fn audit_diff_rows( + table: &paimon::table::Table, + start: i64, + end: i64, +) -> Vec<(String, i32, i32)> { + let audit = AuditLogTable::new(table.clone()); + let plan = audit + .new_incremental_scan(IncrementalScanMode::Diff, start, end) + .plan() + .await + .unwrap(); + let batches: Vec = audit.to_arrow(&plan).unwrap().try_collect().await.unwrap(); + collect_audit_rows(&batches) +} + +fn assert_rows_contain(rows: &[(String, i32, i32)], expected: &[(&str, i32, i32)]) { + for (kind, id, value) in expected { + assert!( + rows.iter() + .any(|(k, i, v)| k == *kind && i == id && v == value), + "missing rowkind={kind} id={id} value={value} in {rows:?}" + ); + } +} + +fn assert_rows_exclude(rows: &[(String, i32, i32)], excluded: &[(&str, i32, i32)]) { + for (kind, id, value) in excluded { + assert!( + !rows + .iter() + .any(|(k, i, v)| k == *kind && i == id && v == value), + "unexpected rowkind={kind} id={id} value={value} in {rows:?}" + ); + } +} + +#[tokio::test] +async fn audit_log_diff_scan_emits_row_level_delete_insert_and_updates() { + let table_path = "memory:/audit_log/diff_range"; + let (file_io, table) = memory_table( + table_path, + pk_schema(&[ + ("changelog-producer", "none"), + ("merge-engine", "deduplicate"), + ("bucket", "1"), + ]), + ); + setup_dirs(&file_io, table_path).await; + persist_table_schema(&file_io, table_path, table.schema()).await; + + write_batch(&table, &make_batch(vec![1, 2], vec![10, 20])).await; + write_batch(&table, &make_batch(vec![2, 3], vec![25, 30])).await; + + let rows = audit_diff_rows(&table, 1, 2).await; + assert_eq!( + rows, + vec![ + ("+I".to_string(), 3, 30), + ("+U".to_string(), 2, 25), + ("-U".to_string(), 2, 20), + ] + ); +} + #[tokio::test] -async fn audit_log_diff_mode_is_unsupported() { - let table_path = "memory:/audit_log/diff_unsupported"; +async fn audit_log_diff_same_snapshot_range_returns_no_rows() { + let table_path = "memory:/audit_log/diff_same_snapshot"; let (file_io, table) = memory_table( table_path, pk_schema(&[ @@ -333,17 +397,215 @@ async fn audit_log_diff_mode_is_unsupported() { ); setup_dirs(&file_io, table_path).await; persist_table_schema(&file_io, table_path, table.schema()).await; + write_batch(&table, &make_batch(vec![1], vec![10])).await; - write_batch(&table, &make_batch(vec![2], vec![20])).await; - let audit = AuditLogTable::new(table.clone()); - let err = audit + let rows = audit_diff_rows(&table, 1, 1).await; + assert!( + rows.is_empty(), + "same start/end snapshot should yield empty diff" + ); +} + +#[tokio::test] +async fn audit_log_diff_insert_only_emits_plus_i() { + let table_path = "memory:/audit_log/diff_insert_only"; + let (file_io, table) = memory_table( + table_path, + pk_schema(&[ + ("changelog-producer", "none"), + ("merge-engine", "deduplicate"), + ("bucket", "1"), + ]), + ); + setup_dirs(&file_io, table_path).await; + persist_table_schema(&file_io, table_path, table.schema()).await; + + write_batch(&table, &make_batch(vec![1], vec![10])).await; + write_batch(&table, &make_batch(vec![1, 2], vec![10, 20])).await; + + let rows = audit_diff_rows(&table, 1, 2).await; + assert_eq!(rows, vec![("+I".to_string(), 2, 20)]); +} + +#[tokio::test] +async fn audit_log_diff_update_only_emits_minus_u_and_plus_u_from_before_after() { + let table_path = "memory:/audit_log/diff_update_only"; + let (file_io, table) = memory_table( + table_path, + pk_schema(&[ + ("changelog-producer", "none"), + ("merge-engine", "deduplicate"), + ("bucket", "1"), + ]), + ); + setup_dirs(&file_io, table_path).await; + persist_table_schema(&file_io, table_path, table.schema()).await; + + write_batch(&table, &make_batch(vec![1], vec![10])).await; + write_batch(&table, &make_batch(vec![1], vec![20])).await; + + let rows = audit_diff_rows(&table, 1, 2).await; + assert_eq!( + rows, + vec![("+U".to_string(), 1, 20), ("-U".to_string(), 1, 10),] + ); +} + +#[tokio::test] +async fn audit_log_diff_delete_via_input_delete_row() { + // Diff compares materialized PK state; input -D removes a key without compact. + let table_path = "memory:/audit_log/diff_delete_input"; + let (file_io, table) = memory_table( + table_path, + pk_schema(&[ + ("changelog-producer", "input"), + ("merge-engine", "deduplicate"), + ("bucket", "1"), + ]), + ); + setup_dirs(&file_io, table_path).await; + persist_table_schema(&file_io, table_path, table.schema()).await; + + write_batch( + &table, + &make_batch_with_kinds(vec![1, 2], vec![10, 20], vec![0, 0]), + ) + .await; + write_batch(&table, &make_batch_with_kinds(vec![1], vec![10], vec![3])).await; + + let rows = audit_diff_rows(&table, 1, 2).await; + assert_rows_contain(&rows, &[("-D", 1, 10)]); + assert_rows_exclude(&rows, &[("+I", 1, 10), ("-U", 1, 10), ("+U", 1, 10)]); +} + +#[tokio::test] +async fn audit_log_diff_mixed_delete_insert_update_without_compact() { + let table_path = "memory:/audit_log/diff_mixed"; + let (file_io, table) = memory_table( + table_path, + pk_schema(&[ + ("changelog-producer", "input"), + ("merge-engine", "deduplicate"), + ("bucket", "1"), + ]), + ); + setup_dirs(&file_io, table_path).await; + persist_table_schema(&file_io, table_path, table.schema()).await; + + write_batch( + &table, + &make_batch_with_kinds(vec![1, 2, 3], vec![10, 20, 30], vec![0, 0, 0]), + ) + .await; + write_batch( + &table, + &make_batch_with_kinds(vec![1, 2, 4], vec![10, 25, 40], vec![3, 2, 0]), + ) + .await; + + let rows = audit_diff_rows(&table, 1, 2).await; + assert_rows_contain( + &rows, + &[("-D", 1, 10), ("-U", 2, 20), ("+U", 2, 25), ("+I", 4, 40)], + ); + assert_rows_exclude(&rows, &[("+I", 3, 30), ("-D", 3, 30)]); +} + +#[tokio::test] +async fn audit_log_diff_processes_multiple_bucket_pairs() { + let table_path = "memory:/audit_log/diff_multi_bucket"; + let (file_io, table) = memory_table( + table_path, + pk_schema(&[ + ("changelog-producer", "none"), + ("merge-engine", "deduplicate"), + ("bucket", "4"), + ]), + ); + setup_dirs(&file_io, table_path).await; + persist_table_schema(&file_io, table_path, table.schema()).await; + + write_batch(&table, &make_batch(vec![1, 8], vec![10, 80])).await; + write_batch(&table, &make_batch(vec![1, 8], vec![11, 81])).await; + + let plan = table + .new_read_builder() .new_incremental_scan(IncrementalScanMode::Diff, 1, 2) .plan() .await - .unwrap_err(); + .unwrap(); + let diff_pairs = plan + .splits() + .iter() + .filter(|split| matches!(split, paimon::table::IncrementalSplit::DiffPair { .. })) + .count(); assert!( - matches!(err, paimon::Error::Unsupported { .. }), - "expected Unsupported for Diff audit plan, got {err:?}" + diff_pairs >= 2, + "expected multiple (partition,bucket) diff pairs, got {diff_pairs}" + ); + + let rows = audit_diff_rows(&table, 1, 2).await; + assert_rows_contain( + &rows, + &[("-U", 1, 10), ("+U", 1, 11), ("-U", 8, 80), ("+U", 8, 81)], ); } + +#[tokio::test] +async fn audit_log_diff_with_sequence_number_enabled_exposes_ordered_columns() { + use std::collections::HashMap; + + let table_path = "memory:/audit_log/diff_sequence"; + let (file_io, table) = memory_table( + table_path, + pk_schema(&[ + ("changelog-producer", "none"), + ("merge-engine", "deduplicate"), + ("bucket", "1"), + ]) + .copy_with_options(HashMap::from([( + "table-read.sequence-number.enabled".to_string(), + "true".to_string(), + )])), + ); + setup_dirs(&file_io, table_path).await; + persist_table_schema(&file_io, table_path, table.schema()).await; + + write_batch(&table, &make_batch(vec![1], vec![10])).await; + write_batch(&table, &make_batch(vec![1], vec![20])).await; + + let audit = AuditLogTable::new(table.clone()); + let field_names: Vec = audit + .fields() + .unwrap() + .into_iter() + .map(|f| f.name().to_string()) + .collect(); + assert_eq!( + field_names, + vec![ + "rowkind".to_string(), + SEQUENCE_NUMBER_FIELD_NAME.to_string(), + "id".to_string(), + "value".to_string(), + ] + ); + + let plan = audit + .new_incremental_scan(IncrementalScanMode::Diff, 1, 2) + .plan() + .await + .unwrap(); + let batches: Vec = audit.to_arrow(&plan).unwrap().try_collect().await.unwrap(); + let rows = collect_audit_rows_with_sequence(&batches); + assert_eq!(rows.len(), 2); + assert_rows_contain( + &rows + .iter() + .map(|(k, _s, i, v)| (k.clone(), *i, *v)) + .collect::>(), + &[("-U", 1, 10), ("+U", 1, 20)], + ); + assert!(rows.iter().all(|(_, seq, _, _)| *seq >= 0)); +} diff --git a/crates/paimon/tests/incremental_batch_scan_test.rs b/crates/paimon/tests/incremental_batch_scan_test.rs index bc007851..28525b33 100644 --- a/crates/paimon/tests/incremental_batch_scan_test.rs +++ b/crates/paimon/tests/incremental_batch_scan_test.rs @@ -429,10 +429,104 @@ async fn incremental_changelog_scan_applies_partition_filter_from_read_builder() assert_eq!(collect_pairs(&batches), vec![(1, 10)]); } -/// Diff mode remains unsupported in this PR. #[tokio::test] -async fn diff_mode_is_unsupported() { - let table_path = "memory:/incremental_batch/diff_unsupported"; +async fn diff_between_snapshots_returns_after_image_rows() { + let table_path = "memory:/incremental_batch/diff_after_image"; + let (file_io, table) = memory_table( + table_path, + pk_schema(&[ + ("changelog-producer", "none"), + ("merge-engine", "deduplicate"), + ("bucket", "1"), + ]), + ); + setup_dirs(&file_io, table_path).await; + persist_table_schema(&file_io, table_path, table.schema()).await; + + write_batch(&table, &make_batch(vec![1, 2], vec![10, 20])).await; + write_batch(&table, &make_batch(vec![2, 3], vec![25, 30])).await; + + let rows = read_incremental_pairs(&table, IncrementalScanMode::Diff, 1, 2).await; + assert_eq!(rows, vec![(2, 25), (3, 30)]); +} + +#[tokio::test] +async fn diff_identical_rows_are_skipped_from_after_image() { + let table_path = "memory:/incremental_batch/diff_identical"; + let (file_io, table) = memory_table( + table_path, + pk_schema(&[ + ("changelog-producer", "none"), + ("merge-engine", "deduplicate"), + ("bucket", "1"), + ]), + ); + setup_dirs(&file_io, table_path).await; + persist_table_schema(&file_io, table_path, table.schema()).await; + + write_batch(&table, &make_batch(vec![1], vec![10])).await; + write_batch(&table, &make_batch(vec![1, 2], vec![10, 20])).await; + + let rows = read_incremental_pairs(&table, IncrementalScanMode::Diff, 1, 2).await; + assert_eq!(rows, vec![(2, 20)]); +} + +#[tokio::test] +async fn diff_rejects_start_before_earliest_snapshot() { + let table_path = "memory:/incremental_batch/diff_earliest"; + let (file_io, table) = memory_table( + table_path, + pk_schema(&[ + ("changelog-producer", "none"), + ("merge-engine", "deduplicate"), + ("bucket", "1"), + ]), + ); + setup_dirs(&file_io, table_path).await; + persist_table_schema(&file_io, table_path, table.schema()).await; + write_batch(&table, &make_batch(vec![1], vec![10])).await; + + let err = plan_incremental(&table, IncrementalScanMode::Diff, 0, 1) + .await + .unwrap_err(); + assert!( + matches!(err, paimon::Error::DataInvalid { .. }), + "expected DataInvalid, got {err:?}" + ); +} + +#[tokio::test] +async fn diff_rejects_non_deduplicate_merge_engine() { + for merge_engine in ["partial-update", "aggregation", "first-row"] { + let table_path = format!("memory:/incremental_batch/diff_engine_{merge_engine}"); + let (file_io, table) = memory_table( + &table_path, + pk_schema(&[ + ("changelog-producer", "none"), + ("merge-engine", merge_engine), + ("bucket", "1"), + ]), + ); + setup_dirs(&file_io, &table_path).await; + persist_table_schema(&file_io, &table_path, table.schema()).await; + write_batch(&table, &make_batch(vec![1], vec![10])).await; + write_batch(&table, &make_batch(vec![2], vec![20])).await; + + let err = plan_incremental(&table, IncrementalScanMode::Diff, 1, 2) + .await + .unwrap_err(); + assert!( + matches!(err, paimon::Error::Unsupported { .. }), + "merge-engine={merge_engine} expected Unsupported, got {err:?}" + ); + } +} + +#[tokio::test] +async fn diff_rejects_bucket_rescale_between_snapshots() { + use std::collections::HashMap; + + let table_path = "memory:/incremental_batch/diff_bucket_rescale"; let (file_io, table) = memory_table( table_path, pk_schema(&[ @@ -444,14 +538,19 @@ async fn diff_mode_is_unsupported() { setup_dirs(&file_io, table_path).await; persist_table_schema(&file_io, table_path, table.schema()).await; write_batch(&table, &make_batch(vec![1], vec![10])).await; + + let table = table.copy_with_options(HashMap::from([("bucket".to_string(), "2".to_string())])); + persist_table_schema(&file_io, table_path, table.schema()).await; write_batch(&table, &make_batch(vec![2], vec![20])).await; - // Non-empty range so planning reaches plan_diff (empty range short-circuits). let err = plan_incremental(&table, IncrementalScanMode::Diff, 1, 2) .await .unwrap_err(); assert!( - matches!(err, paimon::Error::Unsupported { .. }), - "expected Unsupported for Diff, got {err:?}" + matches!( + err, + paimon::Error::Unsupported { ref message } if message.contains("bucket rescale") + ), + "expected Unsupported for bucket rescale, got {err:?}" ); } From c5a7071610b1ad1820ac8cd36f8bdca332950cbf Mon Sep 17 00:00:00 2001 From: pandas886 <123344357+pandas886@users.noreply.github.com> Date: Fri, 17 Jul 2026 10:20:10 +0800 Subject: [PATCH 2/6] fix: correct batch incremental diff reads --- crates/paimon/src/table/audit_log_table.rs | 1 + crates/paimon/src/table/incremental_scan.rs | 130 ++++- crates/paimon/src/table/kv_file_reader.rs | 38 +- crates/paimon/src/table/table_read.rs | 250 ++++++--- crates/paimon/src/table/table_scan.rs | 99 ++-- crates/paimon/tests/audit_log_table_test.rs | 112 ++++- .../tests/incremental_batch_scan_test.rs | 475 +++++++++++++++++- 7 files changed, 948 insertions(+), 157 deletions(-) diff --git a/crates/paimon/src/table/audit_log_table.rs b/crates/paimon/src/table/audit_log_table.rs index f048e9b1..a6b6e5fe 100644 --- a/crates/paimon/src/table/audit_log_table.rs +++ b/crates/paimon/src/table/audit_log_table.rs @@ -81,6 +81,7 @@ impl AuditLogTable { } pub fn to_arrow(&self, plan: &IncrementalPlan) -> crate::Result { + plan.validate()?; let read = self.wrapped.new_read_builder().new_read()?; read.to_audit_log_arrow(plan) } diff --git a/crates/paimon/src/table/incremental_scan.rs b/crates/paimon/src/table/incremental_scan.rs index ea4b1fbe..cdfc15fa 100644 --- a/crates/paimon/src/table/incremental_scan.rs +++ b/crates/paimon/src/table/incremental_scan.rs @@ -47,7 +47,7 @@ pub enum IncrementalScanMode { #[derive(Debug, Clone)] pub enum IncrementalSplit { Data(DataSplit), - /// Per-(partition, bucket) diff pair. Memory bounded by one bucket's data. + /// Per-(partition, bucket) diff pair. DiffPair { before: Vec, after: Vec, @@ -66,6 +66,80 @@ impl IncrementalPlan { Self { mode, splits } } + pub fn try_new( + mode: IncrementalScanMode, + splits: Vec, + ) -> crate::Result { + let plan = Self::new(mode, splits); + plan.validate()?; + Ok(plan) + } + + /// Validate the plan at every point it crosses into a reader. + /// + /// `new` is retained for source compatibility, so callers can still build + /// an invalid plan. Readers must call this method instead of assuming a + /// plan came from the scanner. + pub fn validate(&self) -> crate::Result<()> { + if self.mode == IncrementalScanMode::Auto { + return Err(crate::Error::DataInvalid { + message: "Incremental plan mode Auto must be resolved before consumption" + .to_string(), + source: None, + }); + } + if self.mode == IncrementalScanMode::Diff { + let mut before_snapshot_id = None; + let mut after_snapshot_id = None; + for split in &self.splits { + let IncrementalSplit::DiffPair { before, after } = split else { + return Err(crate::Error::DataInvalid { + message: "Diff incremental plan contains a Data split".to_string(), + source: None, + }); + }; + validate_diff_pair(before, after)?; + if let Some(snapshot_id) = before.first().map(DataSplit::snapshot_id) { + if before_snapshot_id.is_some_and(|expected| expected != snapshot_id) { + return Err(crate::Error::DataInvalid { + message: "Diff plan contains different before snapshots".to_string(), + source: None, + }); + } + before_snapshot_id = Some(snapshot_id); + } + if let Some(snapshot_id) = after.first().map(DataSplit::snapshot_id) { + if after_snapshot_id.is_some_and(|expected| expected != snapshot_id) { + return Err(crate::Error::DataInvalid { + message: "Diff plan contains different after snapshots".to_string(), + source: None, + }); + } + after_snapshot_id = Some(snapshot_id); + } + } + if let (Some(before), Some(after)) = (before_snapshot_id, after_snapshot_id) { + if before >= after { + return Err(crate::Error::DataInvalid { + message: "Diff plan before snapshot must be earlier than after snapshot" + .to_string(), + source: None, + }); + } + } + } else if self + .splits + .iter() + .any(|split| matches!(split, IncrementalSplit::DiffPair { .. })) + { + return Err(crate::Error::DataInvalid { + message: "Non-Diff incremental plan contains a DiffPair".to_string(), + source: None, + }); + } + Ok(()) + } + /// Resolved mode (`Auto` already collapsed to `Delta` / `Changelog`). pub fn mode(&self) -> IncrementalScanMode { self.mode @@ -86,6 +160,49 @@ impl IncrementalPlan { } } +pub(crate) fn validate_diff_pair(before: &[DataSplit], after: &[DataSplit]) -> crate::Result<()> { + if before + .iter() + .chain(after) + .any(|split| split.row_ranges().is_some()) + { + return Err(crate::Error::DataInvalid { + message: "Diff pair must not contain physical row ranges".to_string(), + source: None, + }); + } + let first = before.first().or(after.first()); + let Some(first) = first else { + return Ok(()); + }; + for side in [before, after] { + if let Some(first_in_side) = side.first() { + if side + .iter() + .any(|split| split.snapshot_id() != first_in_side.snapshot_id()) + { + return Err(crate::Error::DataInvalid { + message: "Diff pair side contains splits from different snapshots".to_string(), + source: None, + }); + } + } + } + for split in before.iter().chain(after) { + if split.partition() != first.partition() + || split.bucket() != first.bucket() + || split.bucket_path() != first.bucket_path() + || split.total_buckets() != first.total_buckets() + { + return Err(crate::Error::DataInvalid { + message: "Diff pair contains splits from different partition buckets".to_string(), + source: None, + }); + } + } + Ok(()) +} + /// Batch incremental scan over a snapshot id range. pub struct IncrementalScan<'a> { table: &'a Table, @@ -202,7 +319,7 @@ impl<'a> IncrementalScan<'a> { let plan = self.scan.plan_snapshot_delta(&snapshot).await?; splits.extend(plan.splits().iter().cloned().map(IncrementalSplit::Data)); } - Ok(IncrementalPlan::new(mode, splits)) + IncrementalPlan::try_new(mode, splits) } async fn plan_changelog(&self, mode: IncrementalScanMode) -> crate::Result { @@ -220,10 +337,15 @@ impl<'a> IncrementalScan<'a> { let plan = self.scan.plan_snapshot_changelog(&snapshot).await?; splits.extend(plan.splits().iter().cloned().map(IncrementalSplit::Data)); } - Ok(IncrementalPlan::new(mode, splits)) + IncrementalPlan::try_new(mode, splits) } async fn plan_diff(&self, mode: IncrementalScanMode) -> crate::Result { + if self.table.schema().primary_keys().is_empty() { + return Err(crate::Error::Unsupported { + message: "Batch incremental Diff requires a table with primary keys".to_string(), + }); + } let core_options = CoreOptions::new(self.table.schema().options()); if core_options.merge_engine()? != crate::spec::MergeEngine::Deduplicate { return Err(crate::Error::Unsupported { @@ -269,6 +391,6 @@ impl<'a> IncrementalScan<'a> { splits.push(IncrementalSplit::DiffPair { before, after }); } - Ok(IncrementalPlan::new(mode, splits)) + IncrementalPlan::try_new(mode, splits) } } diff --git a/crates/paimon/src/table/kv_file_reader.rs b/crates/paimon/src/table/kv_file_reader.rs index ec221af5..8cb79471 100644 --- a/crates/paimon/src/table/kv_file_reader.rs +++ b/crates/paimon/src/table/kv_file_reader.rs @@ -73,6 +73,8 @@ pub(crate) struct KeyValueReadConfig { pub merge_engine: MergeEngine, pub sequence_fields: Vec, pub read_batch_size: usize, + /// Merge files from all supplied splits into one globally key-sorted stream. + pub merge_splits: bool, } /// Keep only the conjuncts of `predicates` that reference primary-key columns, @@ -368,7 +370,15 @@ impl KeyValueFileReader { } } - let splits: Vec = data_splits.to_vec(); + let split_groups: Vec> = if self.config.merge_splits { + vec![data_splits.to_vec()] + } else { + data_splits + .iter() + .cloned() + .map(|split| vec![split]) + .collect() + }; let file_io = self.file_io; let merge_engine = self.config.merge_engine; let schema_manager = self.config.schema_manager; @@ -391,21 +401,23 @@ impl KeyValueFileReader { let merge_output_schema = build_target_arrow_schema(&merge_output_fields)?; Ok(try_stream! { - for split in &splits { + for split_group in &split_groups { // DV mode should not reach KeyValueFileReader. - if split - .data_deletion_files() - .is_some_and(|files| files.iter().any(Option::is_some)) - { - Err(Error::Unsupported { - message: "KeyValueFileReader does not support deletion vectors".to_string(), - })?; + for split in split_group { + if split + .data_deletion_files() + .is_some_and(|files| files.iter().any(Option::is_some)) + { + Err(Error::Unsupported { + message: "KeyValueFileReader does not support deletion vectors".to_string(), + })?; + } } - // Create one stream per data file. let mut file_streams: Vec = Vec::new(); - for file_meta in split.data_files().to_vec() { + for split in split_group { + for file_meta in split.data_files().to_vec() { let data_fields: Option> = if file_meta.schema_id != table_schema_id { let data_schema = schema_manager.schema(file_meta.schema_id).await?; Some(data_schema.fields().to_vec()) @@ -428,7 +440,7 @@ impl KeyValueFileReader { file_meta, data_fields, None, - None, + split.row_ranges().map(|ranges| ranges.to_vec()), )?; #[cfg(test)] let stream = if let Some(batch_sizes) = input_batch_sizes.clone() { @@ -443,6 +455,7 @@ impl KeyValueFileReader { stream }; file_streams.push(stream); + } } if file_streams.is_empty() { @@ -809,6 +822,7 @@ mod tests { .map(|field| field.to_string()) .collect(), read_batch_size: core_options.read_batch_size().unwrap(), + merge_splits: false, }, ) .with_input_batch_sizes(input_batch_sizes.clone()); diff --git a/crates/paimon/src/table/table_read.rs b/crates/paimon/src/table/table_read.rs index 6f78a8af..98400e3a 100644 --- a/crates/paimon/src/table/table_read.rs +++ b/crates/paimon/src/table/table_read.rs @@ -29,7 +29,10 @@ use crate::spec::{ VALUE_KIND_FIELD_ID, VALUE_KIND_FIELD_NAME, }; use crate::DataSplit; -use arrow_array::{builder::StringBuilder, Array, ArrayRef, RecordBatch, StringArray, UInt32Array}; +use arrow_array::{ + builder::StringBuilder, Array, ArrayRef, RecordBatch, RecordBatchOptions, StringArray, + UInt32Array, +}; use arrow_schema::Schema as ArrowSchema; use arrow_select::concat::concat as arrow_concat; use arrow_select::take::take; @@ -145,6 +148,7 @@ impl<'a> TableRead<'a> { &self, plan: &IncrementalPlan, ) -> crate::Result { + plan.validate()?; match &self.0 { TableReadKind::Paimon(read) => read.to_incremental_arrow(plan), TableReadKind::Format(_) => Err(crate::Error::Unsupported { @@ -163,6 +167,7 @@ impl<'a> TableRead<'a> { &self, plan: &IncrementalPlan, ) -> crate::Result { + plan.validate()?; match &self.0 { TableReadKind::Paimon(read) => read.to_audit_log_arrow(plan), TableReadKind::Format(_) => Err(crate::Error::Unsupported { @@ -260,16 +265,7 @@ impl<'a> PaimonTableRead<'a> { &self, plan: &IncrementalPlan, ) -> crate::Result { - let pairs: Vec<(Vec, Vec)> = plan - .splits() - .iter() - .filter_map(|s| match s { - IncrementalSplit::DiffPair { before, after } => { - Some((before.clone(), after.clone())) - } - _ => None, - }) - .collect(); + let pairs = diff_pairs(plan)?; let parallel = CoreOptions::new(self.table.schema().options()).diff_parallelism(); let table = self.table.clone(); let read_type = self.read_type.clone(); @@ -280,18 +276,19 @@ impl<'a> PaimonTableRead<'a> { let table = table.clone(); let read_type = read_type.clone(); let data_predicates = data_predicates.clone(); - async move { + let worker: ArrowRecordBatchStream = Box::pin(async_stream::try_stream! { let pair_read = PaimonTableRead::new(&table, read_type, data_predicates); - pair_read.to_diff_after_image_stream(&before, &after) - } + let mut pair_stream = pair_read.to_diff_after_image_stream(&before, &after)?; + while let Some(batch) = pair_stream.next().await { + yield batch?; + } + }); + worker })) - .buffer_unordered(parallel); - while let Some(stream_result) = workers.next().await { - let mut pair_stream = stream_result?; - while let Some(batch) = pair_stream.next().await { - yield batch?; - } + .flatten_unordered(parallel); + while let Some(batch) = workers.next().await { + yield batch?; } })) } @@ -307,7 +304,11 @@ impl<'a> PaimonTableRead<'a> { self.audit_raw_stream(plan, !self.table.schema().primary_keys().is_empty()) } IncrementalScanMode::Changelog => self.audit_raw_stream(plan, true), - IncrementalScanMode::Auto => unreachable!("Auto resolved during plan()"), + IncrementalScanMode::Auto => Err(crate::Error::DataInvalid { + message: "Incremental plan mode Auto must be resolved before consumption" + .to_string(), + source: None, + }), } } @@ -316,6 +317,7 @@ impl<'a> PaimonTableRead<'a> { plan: &IncrementalPlan, has_value_kind: bool, ) -> crate::Result { + plan.validate()?; let data_splits = plan.data_splits(); let user_read_type = self.read_type.clone(); let include_sequence = audit_sequence_number_enabled(self.table); @@ -401,16 +403,7 @@ impl<'a> PaimonTableRead<'a> { } fn audit_diff_stream(&self, plan: &IncrementalPlan) -> crate::Result { - let pairs: Vec<(Vec, Vec)> = plan - .splits() - .iter() - .filter_map(|s| match s { - IncrementalSplit::DiffPair { before, after } => { - Some((before.clone(), after.clone())) - } - _ => None, - }) - .collect(); + let pairs = diff_pairs(plan)?; let parallel = CoreOptions::new(self.table.schema().options()).diff_parallelism(); let table = self.table.clone(); let read_type = self.read_type.clone(); @@ -421,17 +414,19 @@ impl<'a> PaimonTableRead<'a> { let table = table.clone(); let read_type = read_type.clone(); let data_predicates = data_predicates.clone(); - async move { + let worker: ArrowRecordBatchStream = Box::pin(async_stream::try_stream! { let pair_read = PaimonTableRead::new(&table, read_type, data_predicates); - pair_read.to_audit_log_arrow_for_diff(&before, &after) - } + let mut pair_stream = + pair_read.to_audit_log_arrow_for_diff(&before, &after)?; + while let Some(batch) = pair_stream.next().await { + yield batch?; + } + }); + worker })) - .buffer_unordered(parallel); - while let Some(stream_result) = workers.next().await { - let mut pair_stream = stream_result?; - while let Some(batch) = pair_stream.next().await { - yield batch?; - } + .flatten_unordered(parallel); + while let Some(batch) = workers.next().await { + yield batch?; } })) } @@ -441,11 +436,11 @@ impl<'a> PaimonTableRead<'a> { before: &[DataSplit], after: &[DataSplit], ) -> crate::Result { - ensure_diff_supported_read_type(&self.read_type)?; let include_sequence = audit_sequence_number_enabled(self.table); let audit_schema = audit_schema_for_read_type(&self.read_type, include_sequence)?; - let mut diff_read_type = self.read_type.clone(); + let mut diff_read_type = self.table.schema().fields().to_vec(); + ensure_diff_supported_read_type(&diff_read_type)?; if include_sequence { diff_read_type.insert( 0, @@ -526,23 +521,43 @@ impl<'a> PaimonTableRead<'a> { before: &[DataSplit], after: &[DataSplit], ) -> crate::Result { - ensure_diff_supported_read_type(&self.read_type)?; - let key_indices = primary_key_indices(self.table, &self.read_type)?; - let value_indices = value_indices_for_diff(self.table, &self.read_type); + 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)?; + 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: Vec = (0..self.read_type.len()).collect(); + let output_col_indices = self + .read_type + .iter() + .map(|field| { + diff_read_type + .iter() + .position(|candidate| candidate.id() == field.id()) + .ok_or_else(|| crate::Error::DataInvalid { + message: format!("Diff read missing projected column '{}'", field.name()), + source: None, + }) + }) + .collect::>>()?; let table = self.table.clone(); - let read_type = self.read_type.clone(); let data_predicates = self.data_predicates.clone(); let before = before.to_vec(); let after = after.to_vec(); Ok(Box::pin(async_stream::try_stream! { let core_options = CoreOptions::new(table.schema().options()); - let pair_read = PaimonTableRead::new(&table, read_type, data_predicates); - let before_stream = pair_read.read_pk_sorted_for_diff(&before, &core_options)?; - let after_stream = pair_read.read_pk_sorted_for_diff(&after, &core_options)?; + let pair_read = PaimonTableRead::new(&table, diff_read_type.clone(), data_predicates); + let before_stream = pair_read.read_pk_sorted_for_diff_with_type( + &before, + &core_options, + &diff_read_type, + )?; + let after_stream = pair_read.read_pk_sorted_for_diff_with_type( + &after, + &core_options, + &diff_read_type, + )?; let mut bc = ArrowCursor::new(before_stream).await?; let mut ac = ArrowCursor::new(after_stream).await?; let mut builder = @@ -577,14 +592,6 @@ impl<'a> PaimonTableRead<'a> { })) } - fn read_pk_sorted_for_diff( - &self, - splits: &[DataSplit], - core_options: &CoreOptions, - ) -> crate::Result { - self.read_pk_sorted_for_diff_with_type(splits, core_options, &self.read_type) - } - fn read_pk_sorted_for_diff_with_type( &self, splits: &[DataSplit], @@ -621,6 +628,8 @@ impl<'a> PaimonTableRead<'a> { .iter() .map(|s| s.to_string()) .collect(), + read_batch_size: core_options.read_batch_size()?, + merge_splits: true, }, ); reader.read(splits) @@ -755,6 +764,7 @@ impl<'a> PaimonTableRead<'a> { .map(|s| s.to_string()) .collect(), read_batch_size: core_options.read_batch_size()?, + merge_splits: false, }, ); reader.read(splits) @@ -1069,6 +1079,7 @@ impl DiffAfterImageBatchBuilder { } fn flush(&mut self) -> crate::Result { + let row_count = self.len; let mut columns = Vec::with_capacity(self.col_indices.len()); for &col_idx in &self.col_indices { let taken: Vec = self @@ -1097,7 +1108,8 @@ impl DiffAfterImageBatchBuilder { self.row_indices.clear(); self.pinned_batches.clear(); self.len = 0; - RecordBatch::try_new(self.schema.clone(), columns).map_err(|e| { + let options = RecordBatchOptions::new().with_row_count(Some(row_count)); + RecordBatch::try_new_with_options(self.schema.clone(), columns, &options).map_err(|e| { crate::Error::UnexpectedError { message: format!("Failed to build diff after-image batch: {e}"), source: Some(Box::new(e)), @@ -1106,6 +1118,26 @@ impl DiffAfterImageBatchBuilder { } } +fn diff_pairs(plan: &IncrementalPlan) -> crate::Result, Vec)>> { + plan.validate()?; + if plan.mode() != IncrementalScanMode::Diff { + return Err(crate::Error::DataInvalid { + message: "Diff reader requires a Diff incremental plan".to_string(), + source: None, + }); + } + plan.splits() + .iter() + .map(|split| match split { + IncrementalSplit::DiffPair { before, after } => Ok((before.clone(), after.clone())), + IncrementalSplit::Data(_) => Err(crate::Error::DataInvalid { + message: "Diff incremental plan contains a Data split".to_string(), + source: None, + }), + }) + .collect() +} + fn diff_output_col_indices( batch: &RecordBatch, read_type: &[DataField], @@ -1155,7 +1187,7 @@ fn primary_key_indices(table: &Table, read_type: &[DataField]) -> crate::Result< .iter() .position(|field| field.name() == pk) .ok_or_else(|| crate::Error::DataInvalid { - message: format!("Primary key column '{pk}' missing from read projection"), + message: format!("Primary key column '{pk}' missing from Diff comparison schema"), source: None, })?; indices.push(idx); @@ -1168,8 +1200,9 @@ fn ensure_diff_supported_read_type(read_type: &[DataField]) -> crate::Result<()> if !is_diff_supported_type(field.data_type()) { return Err(crate::Error::Unsupported { message: format!( - "Batch incremental Diff does not support nested or decimal column '{}'", - field.name() + "Batch incremental Diff does not support column '{}' of type {:?}", + field.name(), + field.data_type() ), }); } @@ -1178,9 +1211,18 @@ fn ensure_diff_supported_read_type(read_type: &[DataField]) -> crate::Result<()> } fn is_diff_supported_type(data_type: &DataType) -> bool { - !matches!( + matches!( data_type, - DataType::Decimal(_) | DataType::Array(_) | DataType::Map(_) | DataType::Row(_) + DataType::Boolean(_) + | DataType::TinyInt(_) + | DataType::SmallInt(_) + | DataType::Int(_) + | DataType::BigInt(_) + | DataType::Float(_) + | DataType::Double(_) + | DataType::Char(_) + | DataType::VarChar(_) + | DataType::Date(_) ) } @@ -1260,6 +1302,13 @@ fn scalar_compare( Int8Array, StringArray, UInt16Array, UInt32Array, UInt64Array, UInt8Array, }; + match (left.is_null(left_row), right.is_null(right_row)) { + (true, true) => return Ok(Ordering::Equal), + (true, false) => return Ok(Ordering::Less), + (false, true) => return Ok(Ordering::Greater), + (false, false) => {} + } + macro_rules! compare { ($ty:ty, $getter:expr) => { if let (Some(a), Some(b)) = ( @@ -1293,19 +1342,23 @@ fn scalar_compare( left.as_any().downcast_ref::(), right.as_any().downcast_ref::(), ) { - return Ok(a - .value(left_row) - .partial_cmp(&b.value(right_row)) - .unwrap_or(Ordering::Equal)); + let (left, right) = (a.value(left_row), b.value(right_row)); + return Ok(if left.is_nan() && right.is_nan() { + Ordering::Equal + } else { + left.total_cmp(&right) + }); } if let (Some(a), Some(b)) = ( left.as_any().downcast_ref::(), right.as_any().downcast_ref::(), ) { - return Ok(a - .value(left_row) - .partial_cmp(&b.value(right_row)) - .unwrap_or(Ordering::Equal)); + let (left, right) = (a.value(left_row), b.value(right_row)); + return Ok(if left.is_nan() && right.is_nan() { + Ordering::Equal + } else { + left.total_cmp(&right) + }); } Err(crate::Error::Unsupported { @@ -1440,8 +1493,8 @@ mod tests { } #[test] - fn test_diff_rejects_nested_and_decimal_types() { - use crate::spec::{ArrayType, DecimalType, IntType}; + fn test_diff_rejects_types_without_comparator_support() { + use crate::spec::{ArrayType, DecimalType, IntType, TimestampType}; let decimal = DataField::new( 1, @@ -1453,13 +1506,58 @@ mod tests { "tags".to_string(), DataType::Array(ArrayType::new(DataType::Int(IntType::new()))), ); + let timestamp = DataField::new( + 3, + "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("nested or decimal") + Err(crate::Error::Unsupported { message }) if message.contains("amount") )); assert!(matches!( ensure_diff_supported_read_type(&[nested]), - Err(crate::Error::Unsupported { message }) if message.contains("nested or decimal") + 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] + fn test_diff_scalar_compare_distinguishes_null_and_nan_values() { + use arrow_array::{Float32Array, Int32Array}; + + let null = Int32Array::from(vec![None]); + let zero = Int32Array::from(vec![Some(0)]); + assert_eq!( + scalar_compare(&null, 0, &zero, 0).unwrap(), + Ordering::Less, + "NULL -> 0 must be reported as a changed value" + ); + + let nan = Float32Array::from(vec![f32::NAN]); + let one = Float32Array::from(vec![1.0]); + assert_ne!( + scalar_compare(&nan, 0, &one, 0).unwrap(), + Ordering::Equal, + "NaN must not hide a change to a finite value" + ); + + let negative_nan = Float32Array::from(vec![f32::from_bits(0xffc0_0001)]); + assert_eq!( + scalar_compare(&nan, 0, &negative_nan, 0).unwrap(), + Ordering::Equal, + "all NaN representations must compare equal like Java Float.compare" + ); + + let negative_zero = Float32Array::from(vec![-0.0]); + let positive_zero = Float32Array::from(vec![0.0]); + assert_ne!( + scalar_compare(&negative_zero, 0, &positive_zero, 0).unwrap(), + Ordering::Equal, + "signed zero must remain distinguishable like Java Float.compare" + ); + } } diff --git a/crates/paimon/src/table/table_scan.rs b/crates/paimon/src/table/table_scan.rs index 4c2c0f1a..fa58dad9 100644 --- a/crates/paimon/src/table/table_scan.rs +++ b/crates/paimon/src/table/table_scan.rs @@ -1023,6 +1023,9 @@ struct PaimonTableScan<'a> { /// When set, the scan will try to return only enough splits to satisfy the limit. limit: Option, row_ranges: Option>, + /// Diff compares complete logical states, so it must not accept physical + /// row-range pruning from an explicit range or a global-index lookup. + row_range_optimization_disabled: bool, /// When true, disables level-0 filtering so all files are visible. /// Used by non-read paths (overwrite, truncate, writer restore) that need /// the complete file set. Normal read scans leave this as `false`. @@ -1046,6 +1049,7 @@ impl<'a> PaimonTableScan<'a> { bucket_predicate, limit, row_ranges, + row_range_optimization_disabled: false, scan_all_files: false, projected_read_field_ids: None, } @@ -1074,6 +1078,12 @@ impl<'a> PaimonTableScan<'a> { self } + fn without_row_range_optimization(mut self) -> Self { + self.row_ranges = None; + self.row_range_optimization_disabled = true; + self + } + pub(super) fn with_projected_read_field_ids( mut self, projected_read_field_ids: Option>, @@ -1552,34 +1562,30 @@ impl<'a> PaimonTableScan<'a> { /// Plan before/after full-snapshot states for Diff incremental scan. /// /// Loads full manifest entries for both snapshots, rejects bucket rescale, - /// prunes files that are identical and unaffected by exclusive changes, then - /// builds splits via the shared snapshot planning path. + /// then builds splits via the shared snapshot planning path. Diff keeps the + /// complete state on both sides because serialized key bytes do not preserve + /// the logical ordering required for safe overlap pruning. pub(crate) async fn plan_snapshot_diff( &self, before: &Snapshot, after: &Snapshot, ) -> crate::Result<(Plan, Plan)> { self.ensure_query_auth_allowed()?; - let mut before_entries = self.plan_manifest_entries(before).await?; - let mut after_entries = self.plan_manifest_entries(after).await?; + let before_entries = self.plan_manifest_entries(before).await?; + let after_entries = self.plan_manifest_entries(after).await?; Self::validate_diff_bucket_layout(&before_entries, &after_entries)?; - Self::prune_unchanged_diff_files(&mut before_entries, &mut after_entries); - let data_evolution_read_field_ids = self.projected_read_field_ids()?; - let before_plan = self - .plan_snapshot_from_entries( - before.clone(), - before_entries, - data_evolution_read_field_ids.as_ref(), - None, - ) + // A limit hint cannot be pushed into either side of a Diff: truncating + // the states independently can both hide changes and invent them. + let mut full_state_scan = self.clone(); + full_state_scan.limit = None; + // Row ranges identify physical positions in individual files, whereas + // Diff compares complete logical states across both snapshots. + full_state_scan = full_state_scan.without_row_range_optimization(); + let before_plan = full_state_scan + .plan_snapshot_from_entries(before.clone(), before_entries, None, None) .await?; - let after_plan = self - .plan_snapshot_from_entries( - after.clone(), - after_entries, - data_evolution_read_field_ids.as_ref(), - None, - ) + let after_plan = full_state_scan + .plan_snapshot_from_entries(after.clone(), after_entries, None, None) .await?; Ok((before_plan, after_plan)) } @@ -1618,49 +1624,6 @@ impl<'a> PaimonTableScan<'a> { Ok(()) } - fn prune_unchanged_diff_files(before: &mut Vec, after: &mut Vec) { - // Common files (byte-equal manifest entries) that do not overlap any - // exclusive before/after file key-range can be dropped from both sides. - let mut common: Vec = before - .iter() - .filter(|entry| after.iter().any(|other| other == *entry)) - .cloned() - .collect(); - if common.is_empty() { - return; - } - - let exclusive_after: Vec<&ManifestEntry> = after - .iter() - .filter(|entry| !common.iter().any(|c| c == *entry)) - .collect(); - let exclusive_before: Vec<&ManifestEntry> = before - .iter() - .filter(|entry| !common.iter().any(|c| c == *entry)) - .collect(); - common.retain(|entry| { - let overlaps_exclusive = |other: &&ManifestEntry| { - Self::key_ranges_overlap_bytes( - &entry.file().min_key, - &entry.file().max_key, - &other.file().min_key, - &other.file().max_key, - ) - }; - !exclusive_after.iter().any(overlaps_exclusive) - && !exclusive_before.iter().any(overlaps_exclusive) - }); - if common.is_empty() { - return; - } - before.retain(|entry| !common.iter().any(|c| c == entry)); - after.retain(|entry| !common.iter().any(|c| c == entry)); - } - - fn key_ranges_overlap_bytes(min_a: &[u8], max_a: &[u8], min_b: &[u8], max_b: &[u8]) -> bool { - min_a <= max_b && min_b <= max_a - } - /// Read entries from a single manifest list (delta or changelog) with /// partition / bucket filter pushdown matching the full scan path. async fn plan_manifest_list_entries( @@ -4325,4 +4288,14 @@ mod tests { "a dynamic override must not disable query-auth" ); } + + #[test] + fn diff_full_state_disables_global_index_row_range_optimization() { + assert!(!super::should_use_global_index_row_range_optimization( + true, true, true, true, + )); + assert!(super::should_use_global_index_row_range_optimization( + false, true, true, true, + )); + } } diff --git a/crates/paimon/tests/audit_log_table_test.rs b/crates/paimon/tests/audit_log_table_test.rs index 9dc2ca08..a4d2e397 100644 --- a/crates/paimon/tests/audit_log_table_test.rs +++ b/crates/paimon/tests/audit_log_table_test.rs @@ -23,7 +23,7 @@ use paimon::spec::{ DataType, IntType, Schema, TableSchema, VarCharType, ROW_KIND_FIELD_ID, ROW_KIND_FIELD_NAME, SEQUENCE_NUMBER_FIELD_NAME, }; -use paimon::table::{AuditLogTable, IncrementalScanMode}; +use paimon::table::{AuditLogTable, IncrementalPlan, IncrementalScanMode, IncrementalSplit}; use common::incremental_helpers::{ make_batch, make_batch_with_kinds, memory_table, persist_table_schema, pk_schema, setup_dirs, @@ -552,6 +552,62 @@ async fn audit_log_diff_processes_multiple_bucket_pairs() { ); } +#[tokio::test] +async fn audit_log_diff_merges_multiple_splits_per_bucket_by_primary_key() { + let table_path = "memory:/audit_log/diff_multi_split_bucket"; + let (file_io, table) = memory_table( + table_path, + pk_schema(&[ + ("changelog-producer", "none"), + ("merge-engine", "deduplicate"), + ("bucket", "1"), + ("target-file-size", "1b"), + ("source.split.target-size", "1b"), + ("source.split.open-file-cost", "1b"), + ("num-sorted-run.compaction-trigger", "100"), + ]), + ); + setup_dirs(&file_io, table_path).await; + persist_table_schema(&file_io, table_path, table.schema()).await; + + write_batch(&table, &make_batch(vec![1], vec![10])).await; + write_batch(&table, &make_batch(vec![3], vec![30])).await; + write_batch(&table, &make_batch(vec![1], vec![11])).await; + + let audit = AuditLogTable::new(table.clone()); + let plan = audit + .new_incremental_scan(IncrementalScanMode::Diff, 2, 3) + .plan() + .await + .unwrap(); + let scrambled = plan + .splits() + .iter() + .cloned() + .map(|split| match split { + IncrementalSplit::DiffPair { mut before, after } => { + assert!(before.len() >= 2, "test requires multiple before splits"); + assert!(after.len() >= 2, "test requires multiple after splits"); + before.reverse(); + IncrementalSplit::DiffPair { before, after } + } + other => other, + }) + .collect(); + let scrambled = IncrementalPlan::try_new(IncrementalScanMode::Diff, scrambled).unwrap(); + + let batches: Vec = audit + .to_arrow(&scrambled) + .unwrap() + .try_collect() + .await + .unwrap(); + assert_eq!( + collect_audit_rows(&batches), + vec![("+U".to_string(), 1, 11), ("-U".to_string(), 1, 10),] + ); +} + #[tokio::test] async fn audit_log_diff_with_sequence_number_enabled_exposes_ordered_columns() { use std::collections::HashMap; @@ -609,3 +665,57 @@ async fn audit_log_diff_with_sequence_number_enabled_exposes_ordered_columns() { ); assert!(rows.iter().all(|(_, seq, _, _)| *seq >= 0)); } + +#[tokio::test] +async fn audit_log_rejects_invalid_incremental_plan_at_consumption() { + let table_path = "memory:/audit_log/invalid_incremental_plan"; + let (file_io, table) = memory_table( + table_path, + pk_schema(&[("merge-engine", "deduplicate"), ("bucket", "1")]), + ); + setup_dirs(&file_io, table_path).await; + persist_table_schema(&file_io, table_path, table.schema()).await; + let audit = AuditLogTable::new(table.clone()); + + let invalid_kind = IncrementalPlan::new( + IncrementalScanMode::Delta, + vec![IncrementalSplit::DiffPair { + before: Vec::new(), + after: Vec::new(), + }], + ); + let err = match audit.to_arrow(&invalid_kind) { + Ok(_) => panic!("invalid plans must fail instead of producing an empty audit stream"), + Err(err) => err, + }; + assert!( + matches!(err, paimon::Error::DataInvalid { ref message, .. } if message.contains("DiffPair")), + "invalid plans must fail instead of producing an empty audit stream: {err:?}" + ); + + let auto = IncrementalPlan::new(IncrementalScanMode::Auto, Vec::new()); + let err = match audit.to_arrow(&auto) { + Ok(_) => panic!("Auto plans must fail at consumption"), + Err(err) => err, + }; + assert!( + matches!(err, paimon::Error::DataInvalid { ref message, .. } if message.contains("Auto")), + "Auto plans must fail at consumption: {err:?}" + ); + + let err = IncrementalPlan::try_new(IncrementalScanMode::Auto, Vec::new()).unwrap_err(); + assert!( + matches!(err, paimon::Error::DataInvalid { ref message, .. } if message.contains("Auto")), + "try_new must reject unresolved Auto plans: {err:?}" + ); + + let read = table.new_read_builder().new_read().unwrap(); + let err = match read.to_incremental_arrow(&auto) { + Ok(_) => panic!("the direct incremental reader must validate plans too"), + Err(err) => err, + }; + assert!( + matches!(err, paimon::Error::DataInvalid { ref message, .. } if message.contains("Auto")), + "the direct incremental reader must validate plans too: {err:?}" + ); +} diff --git a/crates/paimon/tests/incremental_batch_scan_test.rs b/crates/paimon/tests/incremental_batch_scan_test.rs index 28525b33..13b3cde2 100644 --- a/crates/paimon/tests/incremental_batch_scan_test.rs +++ b/crates/paimon/tests/incremental_batch_scan_test.rs @@ -17,7 +17,7 @@ mod common; -use arrow_array::{Array, Int32Array, RecordBatch}; +use arrow_array::{Array, Int32Array, RecordBatch, StringArray}; use futures::TryStreamExt; use paimon::table::IncrementalScanMode; @@ -471,6 +471,330 @@ async fn diff_identical_rows_are_skipped_from_after_image() { assert_eq!(rows, vec![(2, 20)]); } +#[tokio::test] +async fn diff_projection_without_primary_key_still_compares_full_rows() { + let table_path = "memory:/incremental_batch/diff_projection_without_pk"; + let (file_io, table) = memory_table( + table_path, + pk_schema(&[ + ("changelog-producer", "none"), + ("merge-engine", "deduplicate"), + ("bucket", "1"), + ]), + ); + setup_dirs(&file_io, table_path).await; + persist_table_schema(&file_io, table_path, table.schema()).await; + + write_batch(&table, &make_batch(vec![1], vec![10])).await; + write_batch(&table, &make_batch(vec![1], vec![20])).await; + + let mut builder = table.new_read_builder(); + builder.with_projection(&["value"]).unwrap(); + let plan = builder + .new_incremental_scan(IncrementalScanMode::Diff, 1, 2) + .plan() + .await + .unwrap(); + let batches: Vec = builder + .new_read() + .unwrap() + .to_incremental_arrow(&plan) + .unwrap() + .try_collect() + .await + .unwrap(); + let values: Vec = batches + .iter() + .flat_map(|batch| { + batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap() + .values() + .iter() + .copied() + }) + .collect(); + assert_eq!(values, vec![20]); +} + +#[tokio::test] +async fn diff_change_outside_projection_is_not_missed() { + let table_path = "memory:/incremental_batch/diff_unprojected_change"; + let (file_io, table) = memory_table( + table_path, + pk_schema(&[ + ("changelog-producer", "none"), + ("merge-engine", "deduplicate"), + ("bucket", "1"), + ]), + ); + setup_dirs(&file_io, table_path).await; + persist_table_schema(&file_io, table_path, table.schema()).await; + + write_batch(&table, &make_batch(vec![1], vec![10])).await; + write_batch(&table, &make_batch(vec![1], vec![20])).await; + + 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 = builder + .new_read() + .unwrap() + .to_incremental_arrow(&plan) + .unwrap() + .try_collect() + .await + .unwrap(); + let ids: Vec = batches + .iter() + .flat_map(|batch| { + batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap() + .values() + .iter() + .copied() + }) + .collect(); + assert_eq!(ids, vec![1]); +} + +#[tokio::test] +async fn diff_null_to_zero_is_reported_as_change() { + use arrow_schema::{DataType as ArrowDataType, Field, Schema as ArrowSchema}; + use paimon::spec::{DataType, IntType, Schema, TableSchema}; + use std::sync::Arc; + + let table_path = "memory:/incremental_batch/diff_null_to_zero"; + let schema = Schema::builder() + .column("id", DataType::Int(IntType::new())) + .column("value", DataType::Int(IntType::with_nullable(true))) + .primary_key(["id"]) + .option("changelog-producer", "none") + .option("merge-engine", "deduplicate") + .option("bucket", "1") + .option("bucket-key", "id") + .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_nullable_batch = |value| { + RecordBatch::try_new( + Arc::new(ArrowSchema::new(vec![ + Field::new("id", ArrowDataType::Int32, false), + Field::new("value", ArrowDataType::Int32, true), + ])), + vec![ + Arc::new(Int32Array::from(vec![1])), + Arc::new(Int32Array::from(vec![value])), + ], + ) + .unwrap() + }; + write_batch(&table, &make_nullable_batch(None)).await; + write_batch(&table, &make_nullable_batch(Some(0))).await; + + let rows = read_incremental_pairs(&table, IncrementalScanMode::Diff, 1, 2).await; + assert_eq!(rows, vec![(1, 0)]); +} + +#[tokio::test] +async fn diff_ignores_scan_limit_when_planning_full_states() { + let table_path = "memory:/incremental_batch/diff_limit"; + let (file_io, table) = memory_table( + table_path, + pk_schema(&[ + ("changelog-producer", "none"), + ("merge-engine", "deduplicate"), + ("bucket", "4"), + ]), + ); + setup_dirs(&file_io, table_path).await; + persist_table_schema(&file_io, table_path, table.schema()).await; + + write_batch(&table, &make_batch(vec![1, 8], vec![10, 80])).await; + write_batch(&table, &make_batch(vec![1, 8], vec![11, 81])).await; + + let mut builder = table.new_read_builder(); + builder.with_limit(1); + let plan = builder + .new_incremental_scan(IncrementalScanMode::Diff, 1, 2) + .plan() + .await + .unwrap(); + let pair_count = plan + .splits() + .iter() + .filter(|split| matches!(split, paimon::table::IncrementalSplit::DiffPair { .. })) + .count(); + assert!(pair_count >= 2, "limit must not truncate Diff state pairs"); + + let batches: Vec = builder + .new_read() + .unwrap() + .to_incremental_arrow(&plan) + .unwrap() + .try_collect() + .await + .unwrap(); + assert_eq!(collect_pairs(&batches), vec![(1, 11), (8, 81)]); +} + +#[tokio::test] +async fn diff_empty_projection_preserves_changed_row_count() { + let table_path = "memory:/incremental_batch/diff_empty_projection"; + let (file_io, table) = memory_table( + table_path, + pk_schema(&[ + ("changelog-producer", "none"), + ("merge-engine", "deduplicate"), + ("bucket", "1"), + ]), + ); + setup_dirs(&file_io, table_path).await; + persist_table_schema(&file_io, table_path, table.schema()).await; + + write_batch(&table, &make_batch(vec![1], vec![10])).await; + write_batch(&table, &make_batch(vec![1], vec![20])).await; + + let mut builder = table.new_read_builder(); + builder.with_projection(&[]).unwrap(); + let plan = builder + .new_incremental_scan(IncrementalScanMode::Diff, 1, 2) + .plan() + .await + .unwrap(); + let batches: Vec = builder + .new_read() + .unwrap() + .to_incremental_arrow(&plan) + .unwrap() + .try_collect() + .await + .unwrap(); + assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::(), 1); + assert!(batches.iter().all(|batch| batch.num_columns() == 0)); +} + +#[tokio::test] +async fn diff_ignores_row_ranges_when_planning_full_states() { + use paimon::table::{AuditLogTable, RowRange}; + + let table_path = "memory:/incremental_batch/diff_row_ranges"; + let (file_io, table) = memory_table( + table_path, + pk_schema(&[ + ("changelog-producer", "none"), + ("merge-engine", "deduplicate"), + ("bucket", "1"), + ("row-tracking.enabled", "true"), + ("target-file-size", "1b"), + ("source.split.target-size", "1b"), + ("source.split.open-file-cost", "1b"), + ("num-sorted-run.compaction-trigger", "100"), + ]), + ); + setup_dirs(&file_io, table_path).await; + persist_table_schema(&file_io, table_path, table.schema()).await; + + write_batch(&table, &make_batch(vec![1], vec![10])).await; + write_batch(&table, &make_batch(vec![3], vec![30])).await; + write_batch(&table, &make_batch(vec![1], vec![11])).await; + + let mut builder = table.new_read_builder(); + builder.with_row_ranges(vec![RowRange::new(1, 2)]); + let plan = builder + .new_incremental_scan(IncrementalScanMode::Diff, 2, 3) + .plan() + .await + .unwrap(); + assert!(plan.splits().iter().all(|split| match split { + paimon::table::IncrementalSplit::DiffPair { before, after } => before + .iter() + .chain(after) + .all(|split| split.row_ranges().is_none()), + paimon::table::IncrementalSplit::Data(_) => false, + })); + let batches: Vec = AuditLogTable::new(table.clone()) + .to_arrow(&plan) + .unwrap() + .try_collect() + .await + .unwrap(); + let rowkinds: Vec<&str> = batches + .iter() + .flat_map(|batch| { + batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap() + .iter() + .flatten() + }) + .collect(); + assert_eq!( + rowkinds, + vec!["-U", "+U"], + "Diff must compare complete states rather than physical row ranges" + ); +} + +#[tokio::test] +async fn diff_reads_more_than_128_files_in_one_side() { + let table_path = "memory:/incremental_batch/diff_many_files"; + let (file_io, table) = memory_table( + table_path, + pk_schema(&[ + ("changelog-producer", "none"), + ("merge-engine", "deduplicate"), + ("bucket", "1"), + ("target-file-size", "1b"), + ("source.split.target-size", "1b"), + ("source.split.open-file-cost", "1b"), + ("num-sorted-run.compaction-trigger", "1000"), + ]), + ); + setup_dirs(&file_io, table_path).await; + persist_table_schema(&file_io, table_path, table.schema()).await; + + for id in 0..129 { + write_batch(&table, &make_batch(vec![id], vec![10])).await; + } + write_batch(&table, &make_batch(vec![0], vec![11])).await; + + let plan = plan_incremental(&table, IncrementalScanMode::Diff, 129, 130) + .await + .unwrap(); + let file_count = plan + .splits() + .iter() + .map(|split| match split { + paimon::table::IncrementalSplit::DiffPair { before, .. } => before + .iter() + .map(|split| split.data_files().len()) + .sum::(), + paimon::table::IncrementalSplit::Data(_) => 0, + }) + .sum::(); + assert!(file_count > 128, "test requires more than 128 before files"); + + assert_eq!( + read_incremental_pairs(&table, IncrementalScanMode::Diff, 129, 130).await, + vec![(0, 11)] + ); +} + #[tokio::test] async fn diff_rejects_start_before_earliest_snapshot() { let table_path = "memory:/incremental_batch/diff_earliest"; @@ -522,6 +846,155 @@ async fn diff_rejects_non_deduplicate_merge_engine() { } } +#[tokio::test] +async fn diff_rejects_table_without_primary_keys() { + use paimon::spec::{DataType, IntType, Schema, TableSchema}; + + let table_path = "memory:/incremental_batch/diff_without_primary_keys"; + let schema = Schema::builder() + .column("id", DataType::Int(IntType::new())) + .column("value", DataType::Int(IntType::new())) + .option("changelog-producer", "none") + .option("merge-engine", "deduplicate") + .option("bucket", "1") + .option("bucket-key", "id") + .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; + write_batch(&table, &make_batch(vec![1], vec![10])).await; + write_batch(&table, &make_batch(vec![2], vec![20])).await; + + let err = plan_incremental(&table, IncrementalScanMode::Diff, 1, 2) + .await + .unwrap_err(); + assert!( + matches!(err, paimon::Error::Unsupported { ref message } if message.contains("primary keys")), + "expected Unsupported for a table without primary keys, got {err:?}" + ); +} + +#[test] +fn incremental_plan_rejects_data_split_in_diff_mode() { + use paimon::spec::BinaryRow; + use paimon::table::{DataSplitBuilder, IncrementalPlan, IncrementalSplit}; + + let split = DataSplitBuilder::new() + .with_snapshot(1) + .with_partition(BinaryRow::new(0)) + .with_bucket(0) + .with_bucket_path("memory:/incremental_batch/bucket-0".to_string()) + .with_total_buckets(1) + .with_data_files(Vec::new()) + .build() + .unwrap(); + let err = IncrementalPlan::try_new( + IncrementalScanMode::Diff, + vec![IncrementalSplit::Data(split)], + ) + .unwrap_err(); + assert!( + matches!(err, paimon::Error::DataInvalid { ref message, .. } if message.contains("Data split")), + "Diff plan must reject Data splits instead of silently skipping them: {err:?}" + ); +} + +#[test] +fn incremental_plan_rejects_diff_pair_with_mismatched_bucket_metadata() { + use paimon::spec::BinaryRow; + use paimon::table::{DataSplitBuilder, IncrementalPlan, IncrementalSplit}; + + let split = |bucket| { + DataSplitBuilder::new() + .with_snapshot(1) + .with_partition(BinaryRow::new(0)) + .with_bucket(bucket) + .with_bucket_path(format!("memory:/incremental_batch/bucket-{bucket}")) + .with_total_buckets(2) + .with_data_files(Vec::new()) + .build() + .unwrap() + }; + let err = IncrementalPlan::try_new( + IncrementalScanMode::Diff, + vec![IncrementalSplit::DiffPair { + before: vec![split(0)], + after: vec![split(1)], + }], + ) + .unwrap_err(); + assert!( + matches!(err, paimon::Error::DataInvalid { ref message, .. } if message.contains("partition buckets")), + "Diff plan must reject pairs that cross partition buckets: {err:?}" + ); +} + +#[test] +fn incremental_plan_rejects_partial_or_inconsistent_diff_states() { + use paimon::spec::BinaryRow; + use paimon::table::{DataSplitBuilder, IncrementalPlan, IncrementalSplit, RowRange}; + + let split = |snapshot, with_row_ranges| { + let mut builder = DataSplitBuilder::new() + .with_snapshot(snapshot) + .with_partition(BinaryRow::new(0)) + .with_bucket(0) + .with_bucket_path("memory:/incremental_batch/bucket-0".to_string()) + .with_total_buckets(1) + .with_data_files(Vec::new()); + if with_row_ranges { + builder = builder.with_row_ranges(vec![RowRange::new(0, 1)]); + } + builder.build().unwrap() + }; + + let err = IncrementalPlan::try_new( + IncrementalScanMode::Diff, + vec![IncrementalSplit::DiffPair { + before: vec![split(1, true)], + after: vec![split(2, false)], + }], + ) + .unwrap_err(); + assert!( + matches!(err, paimon::Error::DataInvalid { ref message, .. } if message.contains("row ranges")), + "Diff plan must reject partial physical row ranges: {err:?}" + ); + + let err = IncrementalPlan::try_new( + IncrementalScanMode::Diff, + vec![IncrementalSplit::DiffPair { + before: vec![split(2, false)], + after: vec![split(1, false)], + }], + ) + .unwrap_err(); + assert!( + matches!(err, paimon::Error::DataInvalid { ref message, .. } if message.contains("earlier")), + "Diff plan must reject reversed snapshot states: {err:?}" + ); + + let err = IncrementalPlan::try_new( + IncrementalScanMode::Diff, + vec![ + IncrementalSplit::DiffPair { + before: vec![split(1, false)], + after: Vec::new(), + }, + IncrementalSplit::DiffPair { + before: vec![split(2, false)], + after: Vec::new(), + }, + ], + ) + .unwrap_err(); + assert!( + matches!(err, paimon::Error::DataInvalid { ref message, .. } if message.contains("before snapshots")), + "Diff plan must reject mixed before snapshots: {err:?}" + ); +} + #[tokio::test] async fn diff_rejects_bucket_rescale_between_snapshots() { use std::collections::HashMap; From a84befa0b87b1a5a7fcf6291a875f774ceca5bc5 Mon Sep 17 00:00:00 2001 From: pandas886 <123344357+pandas886@users.noreply.github.com> Date: Sat, 25 Jul 2026 00:12:08 +0800 Subject: [PATCH 3/6] fix: preserve full-state incremental diff scans --- crates/paimon/src/table/table_scan.rs | 32 +++++++++++++++++++-------- 1 file changed, 23 insertions(+), 9 deletions(-) diff --git a/crates/paimon/src/table/table_scan.rs b/crates/paimon/src/table/table_scan.rs index fa58dad9..4b677633 100644 --- a/crates/paimon/src/table/table_scan.rs +++ b/crates/paimon/src/table/table_scan.rs @@ -681,6 +681,18 @@ fn global_index_detail_data_ranges(entries: &[ManifestEntry]) -> Vec { ) } +fn should_use_global_index_row_range_optimization( + row_range_optimization_disabled: bool, + data_evolution_enabled: bool, + global_index_enabled: bool, + has_data_predicates: bool, +) -> bool { + !row_range_optimization_disabled + && data_evolution_enabled + && global_index_enabled + && has_data_predicates +} + fn should_skip_level_zero_for_scan( scan_all_files: bool, has_primary_keys: bool, @@ -1323,10 +1335,12 @@ impl<'a> PaimonTableScan<'a> { core_options: &CoreOptions, data_evolution_enabled: bool, ) -> crate::Result> { - if data_evolution_enabled - && core_options.global_index_enabled() - && !self.data_predicates.is_empty() - { + if should_use_global_index_row_range_optimization( + self.row_range_optimization_disabled, + data_evolution_enabled, + core_options.global_index_enabled(), + !self.data_predicates.is_empty(), + ) { Ok(Some(GlobalIndexScanSettings { search_mode: core_options.global_index_search_mode()?, thread_num: core_options.global_index_thread_num()?, @@ -1571,9 +1585,6 @@ impl<'a> PaimonTableScan<'a> { after: &Snapshot, ) -> crate::Result<(Plan, Plan)> { self.ensure_query_auth_allowed()?; - let before_entries = self.plan_manifest_entries(before).await?; - let after_entries = self.plan_manifest_entries(after).await?; - Self::validate_diff_bucket_layout(&before_entries, &after_entries)?; // A limit hint cannot be pushed into either side of a Diff: truncating // the states independently can both hide changes and invent them. let mut full_state_scan = self.clone(); @@ -1581,11 +1592,14 @@ impl<'a> PaimonTableScan<'a> { // Row ranges identify physical positions in individual files, whereas // Diff compares complete logical states across both snapshots. full_state_scan = full_state_scan.without_row_range_optimization(); + let before_entries = full_state_scan.plan_manifest_entries(before).await?; + let after_entries = full_state_scan.plan_manifest_entries(after).await?; + Self::validate_diff_bucket_layout(&before_entries, &after_entries)?; let before_plan = full_state_scan - .plan_snapshot_from_entries(before.clone(), before_entries, None, None) + .plan_snapshot_from_entries(before.clone(), before_entries, None, None, None, None) .await?; let after_plan = full_state_scan - .plan_snapshot_from_entries(after.clone(), after_entries, None, None) + .plan_snapshot_from_entries(after.clone(), after_entries, None, None, None, None) .await?; Ok((before_plan, after_plan)) } From 2647a4982dce04708b8a9eecd068fd5e5b55b960 Mon Sep 17 00:00:00 2001 From: pandas886 <123344357+pandas886@users.noreply.github.com> Date: Sat, 25 Jul 2026 00:12:08 +0800 Subject: [PATCH 4/6] docs: add batch incremental read API guide --- docs/mkdocs.yml | 1 + docs/src/incremental-reading.md | 94 +++++++++++++++++++++++++++++++++ 2 files changed, 95 insertions(+) create mode 100644 docs/src/incremental-reading.md diff --git a/docs/mkdocs.yml b/docs/mkdocs.yml index 214293d0..6c3ec41c 100644 --- a/docs/mkdocs.yml +++ b/docs/mkdocs.yml @@ -48,6 +48,7 @@ theme: nav: - Home: index.md - Getting Started: getting-started.md + - Batch Incremental Reading: incremental-reading.md - SQL Integration: sql.md - Performance: benchmark.md - C Integration: c-binding.md diff --git a/docs/src/incremental-reading.md b/docs/src/incremental-reading.md new file mode 100644 index 00000000..5843567c --- /dev/null +++ b/docs/src/incremental-reading.md @@ -0,0 +1,94 @@ + + +# Batch Incremental Reading + +The Rust API can plan and read changes between snapshot IDs. Snapshot ranges use +`(start_exclusive, end_inclusive]` semantics. For example, `(3, 5]` includes +snapshots 4 and 5. + +## Scan Modes + +| Mode | Behavior | +|------|----------| +| `IncrementalScanMode::Delta` | Reads data files added by `APPEND` snapshots in the range. | +| `IncrementalScanMode::Changelog` | Reads existing changelog files. It skips `OVERWRITE` snapshots and snapshots without changelog files. | +| `IncrementalScanMode::Auto` | Uses `Delta` when `changelog-producer=none`; otherwise uses `Changelog`. | +| `IncrementalScanMode::Diff` | Compares the complete table states at the start and end snapshots. | + +`Changelog` mode does not generate missing changelog files. Configure a +`changelog-producer` when writing the table if changelog reads are required. + +## Read Incremental Rows + +Build the incremental plan and pass it to `TableRead::to_incremental_arrow`: + +```rust +use futures::TryStreamExt; +use paimon::IncrementalScanMode; + +let read_builder = table.new_read_builder(); +let plan = read_builder + .new_incremental_scan(IncrementalScanMode::Diff, 3, 5) + .plan() + .await?; + +let reader = read_builder.new_read()?; +let batches = reader + .to_incremental_arrow(&plan)? + .try_collect::>() + .await?; +``` + +`Delta` and `Changelog` return rows from their planned files. `Diff` returns +after-image rows for inserted or updated keys and omits deleted keys. Projection +and filters configured on the read builder are applied to the output; `Diff` +still compares complete rows before applying projection. + +## Read Audit-Log Rows + +Use `TableRead::to_audit_log_arrow` when the output must include row kinds: + +```rust +let reader = read_builder.new_read()?; +let batches = reader + .to_audit_log_arrow(&plan)? + .try_collect::>() + .await?; +``` + +The first output column is `rowkind`. `Diff` emits `+I`, `-U`, `+U`, and `-D` +records by comparing the before and after images. If table option +`read.sequence-number.enabled=true` is set, `_SEQUENCE_NUMBER` follows +`rowkind`. + +## Diff Restrictions + +`Diff` currently: + +- requires a primary-key table with `merge-engine=deduplicate`; +- does not support deletion vectors or bucket rescaling between the two + snapshots; +- supports `BOOLEAN`, integer, floating-point, character, string, and `DATE` + columns; +- uses table option `diff.parallelism` to control concurrent + partition-and-bucket comparisons (default: `4`, minimum: `1`). + +The start snapshot must still exist because `Diff` reads both endpoint states. +An equal start and end snapshot produces an empty result. From f2e308a2d764c5bed0c8a6f31d70d2c19d827a3d Mon Sep 17 00:00:00 2001 From: pandas886 <123344357+pandas886@users.noreply.github.com> Date: Sat, 25 Jul 2026 17:07:34 +0800 Subject: [PATCH 5/6] Fix incremental read guards --- crates/paimon/src/table/kv_file_reader.rs | 105 +++++++++++++++++++++- crates/paimon/src/table/table_read.rs | 36 ++++++++ crates/paimon/src/table/table_scan.rs | 84 +++++++++++++++++ docs/src/incremental-reading.md | 2 +- 4 files changed, 225 insertions(+), 2 deletions(-) diff --git a/crates/paimon/src/table/kv_file_reader.rs b/crates/paimon/src/table/kv_file_reader.rs index 8cb79471..64ffce73 100644 --- a/crates/paimon/src/table/kv_file_reader.rs +++ b/crates/paimon/src/table/kv_file_reader.rs @@ -75,6 +75,8 @@ pub(crate) struct KeyValueReadConfig { pub read_batch_size: usize, /// Merge files from all supplied splits into one globally key-sorted stream. pub merge_splits: bool, + /// Optional cap on file streams opened by a single sort-merge group. + pub max_merge_file_streams: Option, } /// Keep only the conjuncts of `predicates` that reference primary-key columns, @@ -148,6 +150,20 @@ fn widen_partial_update_sequence_group_fields( Ok(user_fields) } +fn ensure_merge_fan_in_limit(stream_count: usize, limit: Option) -> crate::Result<()> { + if let Some(limit) = limit { + if stream_count <= limit { + return Ok(()); + } + return Err(Error::Unsupported { + message: format!( + "KeyValueFileReader refuses to merge {stream_count} file streams in one sort-merge group; maximum is {limit}. Compact the table before reading this highly fragmented group" + ), + }); + } + Ok(()) +} + impl KeyValueFileReader { pub(crate) fn new(file_io: FileIO, config: KeyValueReadConfig) -> Self { let pushdown_predicates = retain_primary_key_conjuncts( @@ -391,6 +407,7 @@ impl KeyValueFileReader { let primary_keys = self.config.primary_keys; let sequence_fields = self.config.sequence_fields; let read_batch_size = self.config.read_batch_size; + let max_merge_file_streams = self.config.max_merge_file_streams; #[cfg(test)] let input_batch_sizes = self.input_batch_sizes; @@ -413,6 +430,14 @@ impl KeyValueFileReader { })?; } } + let file_count = split_group + .iter() + .map(|split| split.data_files().len()) + .sum::(); + if file_count == 0 { + continue; + } + ensure_merge_fan_in_limit(file_count, max_merge_file_streams)?; // Create one stream per data file. let mut file_streams: Vec = Vec::new(); @@ -547,8 +572,10 @@ mod tests { use crate::catalog::Identifier; use crate::io::FileIOBuilder; use crate::spec::{ - DataType, Datum, IntType, PredicateBuilder, Schema, TableSchema, VarCharType, + stats::BinaryTableStats, BinaryRow, DataFileMeta, DataType, Datum, IntType, + PredicateBuilder, Schema, TableSchema, VarCharType, }; + use crate::table::source::DataSplitBuilder; use crate::table::table_commit::TableCommit; use crate::table::{Table, TableWrite}; use arrow_array::{Array, Int32Array, StringArray}; @@ -697,6 +724,31 @@ mod tests { .collect() } + fn dummy_data_file(name: String) -> DataFileMeta { + DataFileMeta { + file_name: name, + file_size: 128, + row_count: 1, + min_key: Vec::new(), + max_key: Vec::new(), + key_stats: BinaryTableStats::new(Vec::new(), Vec::new(), Vec::new()), + value_stats: BinaryTableStats::new(Vec::new(), Vec::new(), Vec::new()), + min_sequence_number: 0, + max_sequence_number: 0, + schema_id: 0, + level: 0, + extra_files: Vec::new(), + creation_time: None, + delete_row_count: Some(0), + embedded_index: None, + file_source: None, + value_stats_cols: None, + external_path: None, + first_row_id: None, + write_cols: None, + } + } + #[test] fn retain_primary_key_conjuncts_semantics() { let fields = vec![ @@ -776,6 +828,56 @@ mod tests { ); } + #[tokio::test] + async fn kv_merge_rejects_too_many_file_streams_on_read_path() { + let file_io = test_file_io(); + let table_path = "memory:/kv_merge_fan_in_limit"; + let table = pk_table(&file_io, table_path, &[]); + let core_options = table.schema().core_options(); + let split = DataSplitBuilder::new() + .with_snapshot(1) + .with_partition(BinaryRow::new(0)) + .with_bucket(0) + .with_bucket_path(format!("{table_path}/bucket-0")) + .with_total_buckets(1) + .with_data_files( + (0..257) + .map(|i| dummy_data_file(format!("file-{i}.parquet"))) + .collect(), + ) + .build() + .unwrap(); + let reader = KeyValueFileReader::new( + table.file_io().clone(), + KeyValueReadConfig { + table_name: table.identifier().full_name(), + table_options: table.schema().options().clone(), + schema_manager: table.schema_manager().clone(), + table_schema_id: table.schema().id(), + table_fields: table.schema().fields().to_vec(), + read_type: table.schema().fields().to_vec(), + predicates: Vec::new(), + primary_keys: table.schema().trimmed_primary_keys(), + merge_engine: core_options.merge_engine().unwrap(), + sequence_fields: Vec::new(), + read_batch_size: core_options.read_batch_size().unwrap(), + merge_splits: true, + max_merge_file_streams: Some(256), + }, + ); + + let err = reader + .read(&[split]) + .unwrap() + .try_collect::>() + .await + .unwrap_err(); + assert!( + matches!(err, Error::Unsupported { message } if message.contains("file streams")), + "KV merge must fail before opening an unbounded number of file streams" + ); + } + #[tokio::test] async fn kv_input_decode_honors_read_batch_size_without_changing_merge_batching() { let file_io = test_file_io(); @@ -823,6 +925,7 @@ mod tests { .collect(), read_batch_size: core_options.read_batch_size().unwrap(), merge_splits: false, + max_merge_file_streams: None, }, ) .with_input_batch_sizes(input_batch_sizes.clone()); diff --git a/crates/paimon/src/table/table_read.rs b/crates/paimon/src/table/table_read.rs index 98400e3a..0a179b49 100644 --- a/crates/paimon/src/table/table_read.rs +++ b/crates/paimon/src/table/table_read.rs @@ -148,6 +148,7 @@ impl<'a> TableRead<'a> { &self, plan: &IncrementalPlan, ) -> crate::Result { + self.ensure_query_auth_allowed()?; plan.validate()?; match &self.0 { TableReadKind::Paimon(read) => read.to_incremental_arrow(plan), @@ -167,6 +168,7 @@ impl<'a> TableRead<'a> { &self, plan: &IncrementalPlan, ) -> crate::Result { + self.ensure_query_auth_allowed()?; plan.validate()?; match &self.0 { TableReadKind::Paimon(read) => read.to_audit_log_arrow(plan), @@ -175,6 +177,10 @@ impl<'a> TableRead<'a> { }), } } + + fn ensure_query_auth_allowed(&self) -> crate::Result<()> { + CoreOptions::new(self.table().schema().options()).ensure_read_authorized() + } } #[derive(Debug, Clone)] @@ -630,6 +636,7 @@ impl<'a> PaimonTableRead<'a> { .collect(), read_batch_size: core_options.read_batch_size()?, merge_splits: true, + max_merge_file_streams: Some(256), }, ); reader.read(splits) @@ -765,6 +772,7 @@ impl<'a> PaimonTableRead<'a> { .collect(), read_batch_size: core_options.read_batch_size()?, merge_splits: false, + max_merge_file_streams: None, }, ); reader.read(splits) @@ -1492,6 +1500,34 @@ mod tests { ); } + #[test] + fn test_direct_incremental_read_fails_closed_when_query_auth_enabled() { + let table = query_auth_table(); + let read = TableRead::new(&table, table.schema.fields().to_vec(), Vec::new()); + let plan = IncrementalPlan::new(IncrementalScanMode::Delta, Vec::new()); + assert!( + matches!( + read.to_incremental_arrow(&plan), + Err(crate::Error::Unsupported { ref message }) if message.contains("query-auth.enabled") + ), + "directly-constructed incremental read of a query-auth.enabled table must fail closed" + ); + } + + #[test] + fn test_direct_audit_log_read_fails_closed_when_query_auth_enabled() { + let table = query_auth_table(); + let read = TableRead::new(&table, table.schema.fields().to_vec(), Vec::new()); + let plan = IncrementalPlan::new(IncrementalScanMode::Delta, Vec::new()); + assert!( + matches!( + read.to_audit_log_arrow(&plan), + Err(crate::Error::Unsupported { ref message }) if message.contains("query-auth.enabled") + ), + "directly-constructed audit-log read of a query-auth.enabled table must fail closed" + ); + } + #[test] fn test_diff_rejects_types_without_comparator_support() { use crate::spec::{ArrayType, DecimalType, IntType, TimestampType}; diff --git a/crates/paimon/src/table/table_scan.rs b/crates/paimon/src/table/table_scan.rs index 4b677633..c5494929 100644 --- a/crates/paimon/src/table/table_scan.rs +++ b/crates/paimon/src/table/table_scan.rs @@ -1585,6 +1585,20 @@ impl<'a> PaimonTableScan<'a> { after: &Snapshot, ) -> 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(), + }); + } + if self.row_ranges.is_some() { + return Err(crate::Error::Unsupported { + message: "Batch incremental Diff does not support _ROW_ID row-range filters" + .to_string(), + }); + } // A limit hint cannot be pushed into either side of a Diff: truncating // the states independently can both hide changes and invent them. let mut full_state_scan = self.clone(); @@ -3250,6 +3264,39 @@ mod tests { ) } + fn diff_test_table(table_path: &str, deletion_vectors_enabled: bool) -> Table { + let file_io = FileIOBuilder::new("memory").build().unwrap(); + let mut schema = PaimonSchema::builder() + .column("id", DataType::Int(IntType::new())) + .column("value", DataType::Int(IntType::new())) + .primary_key(["id"]) + .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"), + table_path.to_string(), + TableSchema::new(0, &schema.build().unwrap()), + None, + ) + } + + fn diff_snapshot(snapshot_id: i64) -> Snapshot { + Snapshot::builder() + .version(1) + .id(snapshot_id) + .schema_id(0) + .base_manifest_list(String::new()) + .delta_manifest_list(String::new()) + .commit_user("test-user".to_string()) + .commit_identifier(snapshot_id) + .commit_kind(CommitKind::APPEND) + .time_millis(snapshot_id as u64) + .build() + } + fn two_int_stats_row(id: Option, value: Option) -> Vec { let mut builder = BinaryRowBuilder::new(2); match id { @@ -4303,6 +4350,43 @@ 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_row_id_filters() { + let table = diff_test_table("memory:/diff_row_id_gate", false); + let mut builder = table.new_read_builder(); + let filter = Predicate::Leaf { + column: crate::spec::ROW_ID_FIELD_NAME.to_string(), + index: 0, + data_type: DataType::BigInt(crate::spec::BigIntType::new()), + op: PredicateOperator::GtEq, + literals: vec![Datum::Long(10)], + }; + builder.with_filter(filter); + let scan = builder.new_scan(); + 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("_ROW_ID")), + "Diff must reject _ROW_ID row-range filters instead of dropping them" + ); + } + #[test] fn diff_full_state_disables_global_index_row_range_optimization() { assert!(!super::should_use_global_index_row_range_optimization( diff --git a/docs/src/incremental-reading.md b/docs/src/incremental-reading.md index 5843567c..36e43db1 100644 --- a/docs/src/incremental-reading.md +++ b/docs/src/incremental-reading.md @@ -75,7 +75,7 @@ let batches = reader The first output column is `rowkind`. `Diff` emits `+I`, `-U`, `+U`, and `-D` records by comparing the before and after images. If table option -`read.sequence-number.enabled=true` is set, `_SEQUENCE_NUMBER` follows +`table-read.sequence-number.enabled=true` is set, `_SEQUENCE_NUMBER` follows `rowkind`. ## Diff Restrictions From 48377039754078e6f3a76d8cc568daf3b943af3c Mon Sep 17 00:00:00 2001 From: pandas886 <123344357+pandas886@users.noreply.github.com> Date: Sun, 26 Jul 2026 14:34:42 +0800 Subject: [PATCH 6/6] test: reject row ranges in incremental diff scans --- .../tests/incremental_batch_scan_test.rs | 45 +++++-------------- 1 file changed, 11 insertions(+), 34 deletions(-) diff --git a/crates/paimon/tests/incremental_batch_scan_test.rs b/crates/paimon/tests/incremental_batch_scan_test.rs index 13b3cde2..bd938394 100644 --- a/crates/paimon/tests/incremental_batch_scan_test.rs +++ b/crates/paimon/tests/incremental_batch_scan_test.rs @@ -17,7 +17,7 @@ mod common; -use arrow_array::{Array, Int32Array, RecordBatch, StringArray}; +use arrow_array::{Array, Int32Array, RecordBatch}; use futures::TryStreamExt; use paimon::table::IncrementalScanMode; @@ -687,8 +687,8 @@ async fn diff_empty_projection_preserves_changed_row_count() { } #[tokio::test] -async fn diff_ignores_row_ranges_when_planning_full_states() { - use paimon::table::{AuditLogTable, RowRange}; +async fn diff_rejects_row_ranges_instead_of_dropping_them() { + use paimon::table::RowRange; let table_path = "memory:/incremental_batch/diff_row_ranges"; let (file_io, table) = memory_table( @@ -713,40 +713,17 @@ async fn diff_ignores_row_ranges_when_planning_full_states() { let mut builder = table.new_read_builder(); builder.with_row_ranges(vec![RowRange::new(1, 2)]); - let plan = builder + let err = builder .new_incremental_scan(IncrementalScanMode::Diff, 2, 3) .plan() .await - .unwrap(); - assert!(plan.splits().iter().all(|split| match split { - paimon::table::IncrementalSplit::DiffPair { before, after } => before - .iter() - .chain(after) - .all(|split| split.row_ranges().is_none()), - paimon::table::IncrementalSplit::Data(_) => false, - })); - let batches: Vec = AuditLogTable::new(table.clone()) - .to_arrow(&plan) - .unwrap() - .try_collect() - .await - .unwrap(); - let rowkinds: Vec<&str> = batches - .iter() - .flat_map(|batch| { - batch - .column(0) - .as_any() - .downcast_ref::() - .unwrap() - .iter() - .flatten() - }) - .collect(); - assert_eq!( - rowkinds, - vec!["-U", "+U"], - "Diff must compare complete states rather than physical row ranges" + .unwrap_err(); + assert!( + matches!( + err, + paimon::Error::Unsupported { ref message } if message.contains("_ROW_ID") + ), + "Diff must reject _ROW_ID row-range filters instead of dropping them: {err:?}" ); }