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..a6b6e5fe 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, @@ -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 3aa169e9..cdfc15fa 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, } @@ -46,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, @@ -65,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 @@ -85,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, @@ -201,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 { @@ -219,13 +337,60 @@ 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 { - let _ = mode; - Err(crate::Error::Unsupported { - message: "Batch incremental Diff scan is not implemented yet".to_string(), - }) + 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 { + 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 }); + } + + 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..64ffce73 100644 --- a/crates/paimon/src/table/kv_file_reader.rs +++ b/crates/paimon/src/table/kv_file_reader.rs @@ -73,6 +73,10 @@ 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, + /// 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, @@ -146,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( @@ -368,7 +386,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; @@ -381,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; @@ -391,21 +418,31 @@ 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(), + })?; + } } - + 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(); - 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 +465,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 +480,7 @@ impl KeyValueFileReader { stream }; file_streams.push(stream); + } } if file_streams.is_empty() { @@ -534,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}; @@ -684,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![ @@ -763,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(); @@ -809,6 +924,8 @@ mod tests { .map(|field| field.to_string()) .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 6c57f6b6..0a179b49 100644 --- a/crates/paimon/src/table/table_read.rs +++ b/crates/paimon/src/table/table_read.rs @@ -29,9 +29,15 @@ 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, RecordBatchOptions, 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,12 +142,14 @@ 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, ) -> crate::Result { + self.ensure_query_auth_allowed()?; + plan.validate()?; match &self.0 { TableReadKind::Paimon(read) => read.to_incremental_arrow(plan), TableReadKind::Format(_) => Err(crate::Error::Unsupported { @@ -154,12 +162,14 @@ 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, ) -> crate::Result { + self.ensure_query_auth_allowed()?; + plan.validate()?; match &self.0 { TableReadKind::Paimon(read) => read.to_audit_log_arrow(plan), TableReadKind::Format(_) => Err(crate::Error::Unsupported { @@ -167,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)] @@ -233,9 +247,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,20 +267,54 @@ impl<'a> PaimonTableRead<'a> { self.new_data_file_reader()?.read(&data_splits) } + fn to_incremental_diff_arrow( + &self, + plan: &IncrementalPlan, + ) -> crate::Result { + 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(); + 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(); + let worker: ArrowRecordBatchStream = Box::pin(async_stream::try_stream! { + let pair_read = + PaimonTableRead::new(&table, read_type, data_predicates); + let mut pair_stream = pair_read.to_diff_after_image_stream(&before, &after)?; + while let Some(batch) = pair_stream.next().await { + yield batch?; + } + }); + worker + })) + .flatten_unordered(parallel); + while let Some(batch) = workers.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()) } 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, + }), } } @@ -277,6 +323,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); @@ -361,6 +408,240 @@ impl<'a> PaimonTableRead<'a> { })) } + fn audit_diff_stream(&self, plan: &IncrementalPlan) -> crate::Result { + 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(); + 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(); + let worker: ArrowRecordBatchStream = Box::pin(async_stream::try_stream! { + let pair_read = PaimonTableRead::new(&table, read_type, data_predicates); + 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 + })) + .flatten_unordered(parallel); + while let Some(batch) = workers.next().await { + yield batch?; + } + })) + } + + fn to_audit_log_arrow_for_diff( + &self, + before: &[DataSplit], + after: &[DataSplit], + ) -> crate::Result { + 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.table.schema().fields().to_vec(); + ensure_diff_supported_read_type(&diff_read_type)?; + 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 { + 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 = 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 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, 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 = + 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_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(), + read_batch_size: core_options.read_batch_size()?, + merge_splits: true, + max_merge_file_streams: Some(256), + }, + ); + 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(); @@ -490,6 +771,8 @@ impl<'a> PaimonTableRead<'a> { .map(|s| s.to_string()) .collect(), read_batch_size: core_options.read_batch_size()?, + merge_splits: false, + max_merge_file_streams: None, }, ); reader.read(splits) @@ -609,6 +892,491 @@ 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 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 + .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; + 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)), + } + }) + } +} + +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], + 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 Diff comparison schema"), + 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 column '{}' of type {:?}", + field.name(), + field.data_type() + ), + }); + } + } + Ok(()) +} + +fn is_diff_supported_type(data_type: &DataType) -> bool { + matches!( + data_type, + DataType::Boolean(_) + | DataType::TinyInt(_) + | DataType::SmallInt(_) + | DataType::Int(_) + | DataType::BigInt(_) + | DataType::Float(_) + | DataType::Double(_) + | DataType::Char(_) + | DataType::VarChar(_) + | DataType::Date(_) + ) +} + +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, + }; + + 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)) = ( + 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::(), + ) { + 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::(), + ) { + 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 { + 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 +1499,101 @@ mod tests { "directly-constructed read of a query-auth.enabled table must fail closed" ); } + + #[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}; + + 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()))), + ); + 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("amount") + )); + assert!(matches!( + ensure_diff_supported_read_type(&[nested]), + Err(crate::Error::Unsupported { message }) if message.contains("tags") + )); + assert!(matches!( + ensure_diff_supported_read_type(&[timestamp]), + Err(crate::Error::Unsupported { message }) if message.contains("created_at") + )); + } + + #[test] + 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 c8c0f199..c5494929 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, @@ -987,6 +999,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 { @@ -1009,6 +1035,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`. @@ -1032,6 +1061,7 @@ impl<'a> PaimonTableScan<'a> { bucket_predicate, limit, row_ranges, + row_range_optimization_disabled: false, scan_all_files: false, projected_read_field_ids: None, } @@ -1060,6 +1090,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>, @@ -1299,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()?, @@ -1535,6 +1573,85 @@ 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, + /// 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 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(); + 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_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, None, None) + .await?; + let after_plan = full_state_scan + .plan_snapshot_from_entries(after.clone(), after_entries, None, None, None, 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(()) + } + /// 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( @@ -3147,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 { @@ -4199,4 +4349,51 @@ mod tests { "a dynamic override must not disable query-auth" ); } + + #[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( + 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 8839673e..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, @@ -320,9 +320,45 @@ 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_mode_is_unsupported() { - let table_path = "memory:/audit_log/diff_unsupported"; +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(&[ @@ -333,17 +369,353 @@ 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, 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_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(&[ + ("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![2], vec![20])).await; + + 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(); + let diff_pairs = plan + .splits() + .iter() + .filter(|split| matches!(split, paimon::table::IncrementalSplit::DiffPair { .. })) + .count(); + assert!( + 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_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; + + 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 err = audit + 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_err(); + .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)); +} + +#[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::Unsupported { .. }), - "expected Unsupported for Diff audit plan, got {err:?}" + 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 bc007851..bd938394 100644 --- a/crates/paimon/tests/incremental_batch_scan_test.rs +++ b/crates/paimon/tests/incremental_batch_scan_test.rs @@ -429,10 +429,9 @@ 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(&[ @@ -443,15 +442,565 @@ 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, 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_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_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( + 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 err = builder + .new_incremental_scan(IncrementalScanMode::Diff, 2, 3) + .plan() + .await + .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:?}" + ); +} + +#[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"; + 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_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; + + let table_path = "memory:/incremental_batch/diff_bucket_rescale"; + 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 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:?}" ); } 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..36e43db1 --- /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 +`table-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.