From f0626edee76fc043191e50251536184aec9b960d Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Tue, 15 Sep 2026 18:02:05 +0800 Subject: [PATCH 1/4] [core] Align incremental scans and version selection with Java --- .../python/pypaimon_rust/datafusion.pyi | 5 +- bindings/python/src/read.rs | 28 ++- bindings/python/tests/test_read.py | 37 +++- crates/paimon/src/table/data_file_reader.rs | 16 +- .../paimon/src/table/hybrid_search_builder.rs | 2 +- crates/paimon/src/table/incremental_scan.rs | 7 +- .../tests.rs | 11 +- crates/paimon/src/table/source.rs | 68 ++++--- crates/paimon/src/table/table_read.rs | 13 ++ crates/paimon/src/table/table_scan.rs | 117 +++++++---- crates/paimon/src/table/time_travel.rs | 89 ++++++--- .../tests/incremental_batch_scan_test.rs | 188 ++++++++++++++++-- docs/src/incremental-reading.md | 7 +- docs/src/python-binding.md | 18 +- 14 files changed, 467 insertions(+), 139 deletions(-) diff --git a/bindings/python/python/pypaimon_rust/datafusion.pyi b/bindings/python/python/pypaimon_rust/datafusion.pyi index 40b75c1ea..34a289d45 100644 --- a/bindings/python/python/pypaimon_rust/datafusion.pyi +++ b/bindings/python/python/pypaimon_rust/datafusion.pyi @@ -41,7 +41,8 @@ class Split: def __init__(self, state: bytes) -> None: ... def row_count(self) -> int: ... # Java SplitSerializer v1 binary: a DataSplit (v8), or an IndexedSplit when the split has row ranges. - def serialize(self) -> bytes: ... + def is_streaming(self) -> bool: ... + def serialize(self, *, allow_streaming: bool = False) -> bytes: ... class Plan: def snapshot_id(self) -> Optional[int]: @@ -81,7 +82,7 @@ class ReadBuilder: ... def new_scan(self) -> TableScan: ... def new_incremental_scan(self, start_snapshot_id: int, end_snapshot_id: int) -> TableScan: - """Plan APPEND deltas in (start, end] together, merging primary-key versions across snapshots. + """Plan APPEND deltas in (start, end] together, preserving physical change events. Snapshot IDs are used, not timestamps. The end snapshot must exist. Row-position slicing and sharding use the combined delta batch as their position space. diff --git a/bindings/python/src/read.rs b/bindings/python/src/read.rs index 3c3a58dd7..05802ff3e 100644 --- a/bindings/python/src/read.rs +++ b/bindings/python/src/read.rs @@ -135,12 +135,14 @@ impl PyReadBuilder { // conflict error via its Java-parity silent fallback, so the strict // gate below would otherwise misattribute the failure to a single // selector. Surface the real conflict, listing the keys the user set. + // scan.version must first be adapted by the core: Java allows it to + // overwrite a selector of the same kind after resolving tag precedence. let present: Vec<&str> = TIME_TRAVEL_SELECTORS .iter() .copied() .filter(|name| opts.contains_key(*name)) .collect(); - if present.len() > 1 { + if present.len() > 1 && !opts.contains_key("scan.version") { return Err(PyValueError::new_err(format!( "Only one time-travel selector may be set, found: {}", present.join(", ") @@ -472,10 +474,26 @@ impl PySplit { self.inner.row_count() } - /// Serialize this planned split to the Java `SplitSerializer` (v1) binary, so pypaimon (or - /// any Paimon reader) can rebuild it without re-planning. A split carrying row ranges is - /// serialized as an `IndexedSplit`. - fn serialize<'py>(&self, py: Python<'py>) -> PyResult> { + /// Whether the split must be read as physical change events. + fn is_streaming(&self) -> bool { + self.inner.is_streaming() + } + + /// Serialize to Java SplitSerializer v1, using IndexedSplit for row ranges. + /// Streaming export requires a decoder that preserves change-event semantics. + #[pyo3(signature = (*, allow_streaming=false))] + fn serialize<'py>( + &self, + py: Python<'py>, + allow_streaming: bool, + ) -> PyResult> { + // Older Python decoders silently discard the streaming byte. Require + // acknowledgement before exporting events through that shared codec. + if self.inner.is_streaming() && !allow_streaming { + return Err(PyValueError::new_err( + "Streaming splits require a stream-aware decoder; pass allow_streaming=True", + )); + } let bytes = self.inner.serialize_split_v1().map_err(to_py_err)?; Ok(PyBytes::new(py, &bytes)) } diff --git a/bindings/python/tests/test_read.py b/bindings/python/tests/test_read.py index 8404823ee..5aaffffee 100644 --- a/bindings/python/tests/test_read.py +++ b/bindings/python/tests/test_read.py @@ -694,6 +694,24 @@ def test_time_travel_by_tag_name(): assert _rows(builder.new_read().read(splits)) == 1 +def test_scan_version_adapts_before_binding_selector_validation(): + with tempfile.TemporaryDirectory() as warehouse: + ctx = _make_two_snapshot_table(warehouse) + ctx.sql("CALL sys.create_tag(table => 'tdb.t', tag => 'v1', snapshot_id => 1)") + table = PaimonCatalog({"warehouse": warehouse}).get_table("tdb.t") + for version, key in [("1", "scan.snapshot-id"), ("v1", "scan.tag-name")]: + builder = table.new_read_builder({"scan.version": version, key: "invalid"}) + plan = builder.new_scan().plan() + assert plan.snapshot_id() == 1 + assert pa.Table.from_batches(builder.new_read().read(plan.splits())).to_pydict() == { + "id": [1], "name": ["a"]} + for opts in [ + {"scan.version": "1", "scan.tag-name": "v1"}, + {"scan.version": "v1", "scan.snapshot-id": "1"}]: + with pytest.raises(ValueError, match="did not resolve"): + table.new_read_builder(opts) + + def test_time_travel_unresolved_snapshot_raises(): with tempfile.TemporaryDirectory() as warehouse: _make_two_snapshot_table(warehouse) @@ -869,7 +887,7 @@ def test_split_serialize_encodes_deletions_and_external_path(): assert b"s3://ext/data-0.parquet" in data # external path -def test_combined_incremental_plan_merges_pk_versions_and_preserves_range(): +def test_combined_incremental_plan_preserves_pk_events_and_range(): with tempfile.TemporaryDirectory() as warehouse: ctx = SQLContext() ctx.register_catalog("paimon", {"warehouse": warehouse}) @@ -885,9 +903,17 @@ def test_combined_incremental_plan_merges_pk_versions_and_preserves_range(): assert plan.snapshot_id() == 2 assert len(plan.splits()) == 1 assert pa.Table.from_batches(builder.new_read().read(plan.splits())).to_pydict() == { - "id": [1], "value": [20]} - assert [s.serialize() for s in scan.plan().splits()] == [ - s.serialize() for s in plan.splits()] + "id": [1, 1], "value": [10, 20]} + assert [s.serialize(allow_streaming=True) for s in scan.plan().splits()] == [ + s.serialize(allow_streaming=True) for s in plan.splits()] + assert all(s.is_streaming() for s in plan.splits()) + with pytest.raises(ValueError, match="stream-aware decoder"): + plan.splits()[0].serialize() + import pickle + restored = pickle.loads(pickle.dumps(plan.splits()[0])) + assert restored.is_streaming() + assert restored.serialize(allow_streaming=True) == ( + plan.splits()[0].serialize(allow_streaming=True)) selected = builder.new_incremental_scan(1, 2).plan() assert selected.snapshot_id() == 2 assert pa.Table.from_batches(builder.new_read().read(selected.splits())).to_pydict() == { @@ -981,7 +1007,8 @@ def test_incremental_row_positions_use_combined_delta_batch(): assert plan.snapshot_id() == 2 restored = [pickle.loads(pickle.dumps(split)) for split in plan.splits()] assert pa.Table.from_batches(builder.new_read().read(restored)).column("id").to_pylist() == expected - assert [s.serialize() for s in scan.plan().splits()] == [s.serialize() for s in plan.splits()] + assert [s.serialize(allow_streaming=True) for s in scan.plan().splits()] == [ + s.serialize(allow_streaming=True) for s in plan.splits()] builder.with_row_ranges([(1, 3)]).with_limit(2) plan = builder.new_incremental_scan(0, 2).with_row_position_slice(2, 5).plan() assert pa.Table.from_batches(builder.new_read().read(plan.splits())).column("id").to_pylist() == [2, 3] diff --git a/crates/paimon/src/table/data_file_reader.rs b/crates/paimon/src/table/data_file_reader.rs index 667b681e0..e5d8e7dec 100644 --- a/crates/paimon/src/table/data_file_reader.rs +++ b/crates/paimon/src/table/data_file_reader.rs @@ -2754,13 +2754,27 @@ mod tests { Vec::new(), ); let batches = reader - .read(&[split]) + .clone() + .read(std::slice::from_ref(&split)) .unwrap() .try_collect::>() .await .unwrap(); assert_eq!(collect_ids(&batches), vec![1, 3]); + // Event planners do not attach endpoint DVs, but readers still honor + // explicitly supplied DVs in a Java streaming frame. + let mut bytes = split.serialize().unwrap(); + let flag = bytes.len() - 2; + bytes[flag] = 1; + let streaming = DataSplit::deserialize(&bytes).unwrap(); + let events = reader + .read(&[streaming]) + .unwrap() + .try_collect::>() + .await + .unwrap(); + assert_eq!(collect_ids(&events), vec![1, 3]); } /// A Mosaic file and a Parquet file in the same split must both be read and concatenated. diff --git a/crates/paimon/src/table/hybrid_search_builder.rs b/crates/paimon/src/table/hybrid_search_builder.rs index 2ef9be1ae..795978b71 100644 --- a/crates/paimon/src/table/hybrid_search_builder.rs +++ b/crates/paimon/src/table/hybrid_search_builder.rs @@ -539,7 +539,7 @@ impl<'a> HybridSearchBuilder<'a> { let core = CoreOptions::new(self.table.schema().options()); // Already targeting a fixed snapshot (resolved travel copy or a selector // that resolves deterministically): every route agrees without pinning. - if self.table.has_resolved_travel_snapshot() || core.try_time_travel_selector()?.is_some() { + if self.table.has_resolved_travel_snapshot() || core.has_time_travel_selector() { return Ok(None); } // Read-latest: pin the current latest snapshot once so a concurrent commit diff --git a/crates/paimon/src/table/incremental_scan.rs b/crates/paimon/src/table/incremental_scan.rs index dc000683d..c4e640409 100644 --- a/crates/paimon/src/table/incremental_scan.rs +++ b/crates/paimon/src/table/incremental_scan.rs @@ -277,13 +277,14 @@ impl<'a> IncrementalScan<'a> { } } - /// Plan APPEND deltas as one batch, merging their manifest entries and - /// overlapping primary-key files across the whole snapshot range. + /// Plan APPEND deltas with batch split packing and streaming read semantics. + /// Each physical change is retained, including repeated keys and retracts. /// /// Unlike [`Self::plan`], this returns an ordinary [`Plan`] for a normal /// table reader. It does not preserve a separate result for each commit. /// Only Delta (or Auto resolving to Delta) is supported. The end snapshot - /// must exist, and supplies the plan's snapshot metadata and deletion vectors. + /// must exist and supplies snapshot metadata. Snapshot deletion vectors and + /// automatic global-index pruning do not apply to these historical events. pub async fn plan_combined_delta(&self) -> crate::Result { CoreOptions::new(self.table.schema().options()).ensure_read_authorized()?; let mode = self.resolve_mode(); diff --git a/crates/paimon/src/table/sorted_global_index_build_builder/tests.rs b/crates/paimon/src/table/sorted_global_index_build_builder/tests.rs index 9013fafd6..db58c4bca 100644 --- a/crates/paimon/src/table/sorted_global_index_build_builder/tests.rs +++ b/crates/paimon/src/table/sorted_global_index_build_builder/tests.rs @@ -1133,17 +1133,22 @@ async fn test_empty_global_index_ranges_skip_legacy_manifests() { assert!(plan.splits().is_empty()); let (traced_plan, trace) = read_builder.new_scan().plan_with_trace().await.unwrap(); - let delta_plan = read_builder + let delta_error = read_builder .new_scan() .plan_snapshot_delta(&snapshot) .await - .unwrap(); + .unwrap_err(); assert!(traced_plan.splits().is_empty()); assert_eq!(trace.manifest_entries_read, 0); assert_eq!(trace.final_splits, 0); assert_eq!(trace.final_files, 0); - assert!(delta_plan.splits().is_empty()); + // A snapshot index may short-circuit a batch read, but cannot hide + // an invalid historical event file from incremental planning. + assert!( + matches!(delta_error, crate::Error::DataInvalid { ref message, .. } + if message.contains("First row id")) + ); } } diff --git a/crates/paimon/src/table/source.rs b/crates/paimon/src/table/source.rs index ed8070438..4499fd011 100644 --- a/crates/paimon/src/table/source.rs +++ b/crates/paimon/src/table/source.rs @@ -500,9 +500,16 @@ pub struct DataSplit { /// physical rows are exactly its logical rows (modulo deletion files). /// Mirrors Java `DataSplit#rawConvertible`. raw_convertible: bool, + #[serde(default)] + is_streaming: bool, } impl DataSplit { + /// Whether files contain change events rather than a materialized snapshot. + pub fn is_streaming(&self) -> bool { + self.is_streaming + } + pub fn snapshot_id(&self) -> i64 { self.snapshot_id } @@ -707,7 +714,7 @@ impl DataSplit { out.extend_from_slice(&d); } write_deletion_list(&mut out, self.data_deletion_files.as_deref())?; - out.push(0); // isStreaming = false + out.push(u8::from(self.is_streaming)); out.push(u8::from(self.raw_convertible)); Ok(out) } @@ -819,14 +826,7 @@ impl DataSplit { } let data_deletion_files = read_deletion_list(cur)?; - // Rust only produces and serves batch splits (isStreaming = false); a - // streaming split carries a semantic bit that Java readers branch on, so - // reject it rather than silently dropping it. - if read_u8(cur)? != 0 { - return Err(crate::Error::Unsupported { - message: "streaming DataSplit (isStreaming = true) not supported".to_string(), - }); - } + let is_streaming = read_u8(cur)? != 0; let raw_convertible = read_u8(cur)? != 0; let mut builder = DataSplitBuilder::new() @@ -836,7 +836,8 @@ impl DataSplit { .with_bucket_path(bucket_path) .with_total_buckets(total_buckets.unwrap_or(1)) .with_data_files(data_files) - .with_raw_convertible(raw_convertible); + .with_raw_convertible(raw_convertible) + .with_streaming(is_streaming); if let Some(dels) = data_deletion_files { builder = builder.with_data_deletion_files(dels); } @@ -1165,6 +1166,7 @@ pub struct DataSplitBuilder { data_deletion_files: Option>>, row_ranges: Option>, raw_convertible: bool, + is_streaming: bool, } impl DataSplitBuilder { @@ -1182,6 +1184,7 @@ impl DataSplitBuilder { // utility splits) are raw by nature; the merge-tree and // data-evolution scan paths set this explicitly per split group. raw_convertible: true, + is_streaming: false, } } @@ -1224,6 +1227,12 @@ impl DataSplitBuilder { self } + /// Preserve the Java DataSplit event-reading contract. + pub fn with_streaming(mut self, is_streaming: bool) -> Self { + self.is_streaming = is_streaming; + self + } + /// Mark whether the split can be read raw; see [`DataSplit::raw_convertible`]. pub fn with_raw_convertible(mut self, raw_convertible: bool) -> Self { self.raw_convertible = raw_convertible; @@ -1283,6 +1292,7 @@ impl DataSplitBuilder { data_deletion_files: self.data_deletion_files.map(Into::into), row_ranges: self.row_ranges.map(Into::into), raw_convertible: self.raw_convertible, + is_streaming: self.is_streaming, }) } } @@ -1951,26 +1961,26 @@ mod tests { } } - // A streaming split (isStreaming = true) is rejected rather than silently - // dropping the bit: Rust only writes/serves batch splits (isStreaming = - // false), and Java readers branch on this flag. The isStreaming byte is the - // second-to-last byte on the wire (isStreaming, then raw_convertible). #[test] - fn deserialize_rejects_streaming_split() { - let bytes = sample_v8_split().serialize().unwrap(); - assert!(DataSplit::deserialize(&bytes).is_ok()); - - let mut patched = bytes.clone(); - let pos = patched.len() - 2; // isStreaming flag - assert_eq!( - patched[pos], 0, - "fixture must serialize isStreaming = false" - ); - patched[pos] = 1; - match DataSplit::deserialize(&patched) { - Err(crate::Error::Unsupported { .. }) => {} - other => panic!("expected Unsupported, got {other:?}"), - } + fn streaming_split_preserves_java_flag_and_legacy_json_defaults() { + let batch = sample_v8_split(); + let bytes = batch.serialize().unwrap(); + assert!(!DataSplit::deserialize(&bytes).unwrap().is_streaming()); + let mut streaming_bytes = bytes.clone(); + let pos = streaming_bytes.len() - 2; + streaming_bytes[pos] = 1; + let split = DataSplit::deserialize(&streaming_bytes).unwrap(); + assert!(split.is_streaming()); + assert_eq!(split.serialize().unwrap(), streaming_bytes); + let json = serde_json::to_vec(&split).unwrap(); + assert!(serde_json::from_slice::(&json) + .unwrap() + .is_streaming()); + let mut legacy = serde_json::to_value(batch).unwrap(); + legacy.as_object_mut().unwrap().remove("is_streaming"); + let restored: DataSplit = serde_json::from_value(legacy).unwrap(); + assert!(!restored.is_streaming()); + assert_eq!(restored.serialize().unwrap(), bytes); } #[test] diff --git a/crates/paimon/src/table/table_read.rs b/crates/paimon/src/table/table_read.rs index b1926cb56..7f5f70f19 100644 --- a/crates/paimon/src/table/table_read.rs +++ b/crates/paimon/src/table/table_read.rs @@ -757,6 +757,19 @@ impl<'a> PaimonTableRead<'a> { data_splits: &[DataSplit], core_options: &CoreOptions, ) -> crate::Result { + if data_splits.iter().any(DataSplit::is_streaming) { + let streams = data_splits + .iter() + .map(|split| { + if split.is_streaming() { + self.read_raw(std::slice::from_ref(split)) + } else { + self.read_pk(std::slice::from_ref(split), core_options) + } + }) + .collect::>>()?; + return Ok(Box::pin(futures::stream::select_all(streams))); + } let merge_engine = core_options.merge_engine()?; let dv_enabled = core_options.deletion_vectors_enabled(); if matches!( diff --git a/crates/paimon/src/table/table_scan.rs b/crates/paimon/src/table/table_scan.rs index 431c84078..9f2bd6f80 100644 --- a/crates/paimon/src/table/table_scan.rs +++ b/crates/paimon/src/table/table_scan.rs @@ -111,6 +111,22 @@ async fn read_manifest_list( crate::spec::avro::from_avro_bytes_fast::(&bytes) } +fn validate_incremental_entries(entries: &[ManifestEntry]) -> crate::Result<()> { + if entries.iter().any(|entry| *entry.kind() != FileKind::Add) { + return Err(crate::Error::DataInvalid { + message: "Incremental delta manifests must contain only ADD entries".into(), + source: None, + }); + } + Ok(()) +} + +#[derive(Debug, Clone, Copy, PartialEq)] +enum IncrementalSplitMode { + Batch, + Streaming, +} + enum ManifestListSource<'a> { Snapshot(&'a Snapshot), AppendDeltas(&'a [Snapshot]), @@ -143,6 +159,7 @@ async fn read_all_manifest_entries( row_range_index: Option<&RowRangeIndex>, trace: Option<&mut ScanTrace>, ) -> crate::Result> { + let incremental = matches!(&source, ManifestListSource::AppendDeltas(_)); let (mut manifest_files, delta) = match source { ManifestListSource::Snapshot(snapshot) => futures::try_join!( read_manifest_list(file_io, table_path, snapshot.base_manifest_list()), @@ -300,7 +317,12 @@ async fn read_all_manifest_entries( }, ) .await?; - let mut all_entries = merge_manifest_entries(all_entries); + let mut all_entries = if incremental { + validate_incremental_entries(&all_entries)?; + all_entries + } else { + merge_manifest_entries(all_entries) + }; let manifest_entries_after_merge = all_entries.len(); if let Some(index) = row_range_index { let before = all_entries.len(); @@ -1098,10 +1120,15 @@ struct PaimonTableScan<'a> { /// Used by non-read paths (overwrite, truncate, writer restore) that need /// the complete file set. Normal read scans leave this as `false`. scan_all_files: bool, + incremental_split_mode: Option, projected_read_field_ids: Option>, } impl<'a> PaimonTableScan<'a> { + fn is_streaming(&self) -> bool { + self.incremental_split_mode.is_some() + } + pub(crate) fn new( table: &'a Table, partition_filter: Option, @@ -1120,6 +1147,7 @@ impl<'a> PaimonTableScan<'a> { row_position_selection: None, row_range_optimization_disabled: false, scan_all_files: false, + incremental_split_mode: None, projected_read_field_ids: None, } } @@ -1284,7 +1312,7 @@ impl<'a> PaimonTableScan<'a> { // Non-read paths (overwrite, truncate, writer restore) set scan_all_files=true // to see all files including level-0, matching Java's CommitScanner behavior. let skip_level_zero = should_skip_level_zero_for_scan( - self.scan_all_files, + self.scan_all_files || self.is_streaming(), has_primary_keys, deletion_vectors_enabled, core_options.deletion_vectors_merge_on_read(), @@ -1358,12 +1386,14 @@ impl<'a> PaimonTableScan<'a> { core_options: &CoreOptions, data_evolution_enabled: bool, ) -> crate::Result> { - 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(), - ) { + if !self.is_streaming() + && 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.scalar_index_search_mode()?, thread_num: core_options.global_index_thread_num()?, @@ -1526,8 +1556,8 @@ impl<'a> PaimonTableScan<'a> { Ok(crate::spec::MergeEngine::FirstRow) ); if has_primary_keys - && (!deletion_vectors_enabled || deletion_vectors_merge_on_read) - && !first_row + && (self.is_streaming() + || ((!deletion_vectors_enabled || deletion_vectors_merge_on_read) && !first_row)) { retain_primary_key_conjuncts( &self.data_predicates, @@ -1556,7 +1586,9 @@ impl<'a> PaimonTableScan<'a> { pub(crate) async fn plan_snapshot_delta(&self, snapshot: &Snapshot) -> crate::Result { self.ensure_query_auth_allowed()?; let data_evolution_read_field_ids = self.projected_read_field_ids()?; - self.plan_snapshot_manifest_list( + let mut scan = self.clone(); + scan.incremental_split_mode = Some(IncrementalSplitMode::Streaming); + scan.plan_snapshot_manifest_list( snapshot, snapshot.delta_manifest_list(), data_evolution_read_field_ids.as_ref(), @@ -1573,7 +1605,9 @@ impl<'a> PaimonTableScan<'a> { ) -> crate::Result { self.ensure_query_auth_allowed()?; let data_evolution_read_field_ids = self.projected_read_field_ids()?; - self.plan_snapshot_from_lists( + let mut scan = self.clone(); + scan.incremental_split_mode = Some(IncrementalSplitMode::Batch); + scan.plan_snapshot_from_lists( end_snapshot, ManifestListSource::AppendDeltas(snapshots), data_evolution_read_field_ids.as_ref(), @@ -1593,7 +1627,9 @@ impl<'a> PaimonTableScan<'a> { return Ok(Plan::new(Vec::new()).with_snapshot_id(snapshot.id())); }; let data_evolution_read_field_ids = self.projected_read_field_ids()?; - self.plan_snapshot_manifest_list( + let mut scan = self.clone(); + scan.incremental_split_mode = Some(IncrementalSplitMode::Streaming); + scan.plan_snapshot_manifest_list( snapshot, list_name, data_evolution_read_field_ids.as_ref(), @@ -1618,7 +1654,7 @@ impl<'a> PaimonTableScan<'a> { .read_index_manifest_entries( snapshot, global_index_settings.is_some(), - core_options.deletion_vectors_enabled(), + core_options.deletion_vectors_enabled() && !self.is_streaming(), ) .await?; let manifest_row_ranges = self @@ -1841,7 +1877,7 @@ impl<'a> PaimonTableScan<'a> { )?; entries.extend(manifest_entries); } - let entries = merge_manifest_entries(entries); + validate_incremental_entries(&entries)?; let entries = if let Some(index) = row_range_index { retain_manifest_entry_row_ranges(entries, index) } else { @@ -1886,7 +1922,7 @@ impl<'a> PaimonTableScan<'a> { .read_index_manifest_entries( snapshot, global_index_settings.is_some(), - core_options.deletion_vectors_enabled(), + core_options.deletion_vectors_enabled() && !self.is_streaming(), ) .await?; let manifest_row_ranges = self @@ -1948,7 +1984,8 @@ impl<'a> PaimonTableScan<'a> { let schema_manager = self.table.schema_manager(); let core_options = CoreOptions::new(self.table.schema().options()); let data_evolution_enabled = core_options.data_evolution_enabled(); - let deletion_vectors_enabled = core_options.deletion_vectors_enabled(); + let deletion_vectors_enabled = + core_options.deletion_vectors_enabled() && !self.is_streaming(); let target_split_size = core_options.source_split_target_size(); let open_file_cost = core_options.source_split_open_file_cost(); let partition_keys = self.table.schema().partition_keys(); @@ -2062,13 +2099,14 @@ impl<'a> PaimonTableScan<'a> { // without merging (stale rows are masked by DVs / level-0 is skipped), // so they keep plain size-based packing. DV merge-on-read includes L0 // files and must preserve overlapping key ranges just like ordinary MOR. - let read_merges_overlapping_keys = (!core_options.deletion_vectors_enabled() - || core_options.deletion_vectors_merge_on_read()) - && !matches!( - core_options.merge_engine(), - Ok(crate::spec::MergeEngine::FirstRow) - ); - let pk_comparator = if read_merges_overlapping_keys { + let use_key_interval_packing = self.is_streaming() + || (!core_options.deletion_vectors_enabled() + || core_options.deletion_vectors_merge_on_read()) + && !matches!( + core_options.merge_engine(), + Ok(crate::spec::MergeEngine::FirstRow) + ); + let pk_comparator = if use_key_interval_packing { KeyComparator::from_table_schema(self.table.schema()) } else { None @@ -2114,7 +2152,16 @@ impl<'a> PaimonTableScan<'a> { .and_then(|map| map.get(&PartitionBucket::new(partition, bucket))); // Data-evolution reads merge overlapping row-id groups column-wise. - let file_groups: Vec = if data_evolution_enabled { + let file_groups: Vec = if self.incremental_split_mode + == Some(IncrementalSplitMode::Streaming) + && !data_evolution_enabled + && (!self.table.schema().primary_keys().is_empty() || core_options.bucket() >= 0) + { + vec![SplitGroup { + files: data_files, + raw_convertible: true, + }] + } else if data_evolution_enabled { let (row_id_groups, groups_pruned_by_row_ranges) = data_evolution_row_range_groups(data_files, effective_row_ranges.as_deref())?; if let Some(trace) = trace.as_deref_mut() { @@ -2284,7 +2331,8 @@ impl<'a> PaimonTableScan<'a> { .with_bucket_path(bucket_path.clone()) .with_total_buckets(total_buckets) .with_data_files(file_group) - .with_raw_convertible(raw_convertible); + .with_raw_convertible(raw_convertible) + .with_streaming(self.is_streaming()); if let Some(files) = data_deletion_files { builder = builder.with_data_deletion_files(files); } @@ -3475,7 +3523,7 @@ mod tests { } #[tokio::test] - async fn test_snapshot_delta_row_range_pruning_does_not_resurrect_deleted_witness() { + async fn test_snapshot_delta_rejects_manifest_deletes_before_row_range_pruning() { let table_path = "memory:/de_delta_row_range_netting"; let table = data_evolution_test_table(table_path, two_column_schema(0, "id", "payload")); setup_scan_trace_dirs(&table).await; @@ -3507,19 +3555,14 @@ mod tests { .unwrap(); let mut read_builder = table.new_read_builder(); read_builder.with_row_ranges(vec![RowRange::new(120, 129)]); - let plan = read_builder + let error = read_builder .new_scan() .plan_snapshot_delta(&snapshot) .await - .unwrap(); - - assert_eq!( - plan.splits() - .iter() - .flat_map(|split| split.data_files()) - .map(|file| file.file_name.as_str()) - .collect::>(), - vec!["anchor.parquet"] + .unwrap_err(); + assert!( + matches!(error, crate::Error::DataInvalid { ref message, .. } + if message.contains("only ADD")) ); } diff --git a/crates/paimon/src/table/time_travel.rs b/crates/paimon/src/table/time_travel.rs index cfb80a3e2..8187dce4c 100644 --- a/crates/paimon/src/table/time_travel.rs +++ b/crates/paimon/src/table/time_travel.rs @@ -35,6 +35,35 @@ pub(crate) async fn travel_to_snapshot( tag_manager: &TagManager, options: &HashMap, ) -> crate::Result> { + // Java adapts scan.version before checking mutually exclusive selectors. + // It overwrites the same selector kind, but preserves conflicts with others. + let mut adapted; + let options = if let Some(version) = options.get("scan.version") { + adapted = options.clone(); + adapted.remove("scan.version"); + let (key, value) = if tag_manager.tag_exists(version).await? { + ("scan.tag-name", version.clone()) + } else if let Some(watermark) = version.strip_prefix(WATERMARK_PREFIX) { + let value = watermark.parse::().map_err(|e| Error::DataInvalid { + message: format!("scan.version '{version}' has an invalid watermark value."), + source: Some(Box::new(e)), + })?; + ("scan.watermark", value.to_string()) + } else if !version.is_empty() && version.bytes().all(|b| b.is_ascii_digit()) { + ("scan.snapshot-id", version.clone()) + } else { + return Err(Error::DataInvalid { + message: format!( + "scan.version '{version}' is not a valid tag name or snapshot id." + ), + source: None, + }); + }; + adapted.insert(key.to_string(), value); + &adapted + } else { + options + }; let core_options = CoreOptions::new(options); match core_options.try_time_travel_selector()? { @@ -50,33 +79,7 @@ pub(crate) async fn travel_to_snapshot( Some(TimeTravelSelector::Watermark(w)) => { resolve_watermark(snapshot_manager, w).await.map(Some) } - Some(TimeTravelSelector::Version { - value: v, - option_name, - }) => { - // Match Java TimeTravelUtil.adaptScanVersion: tag first, then the - // `watermark-` prefix, then snapshot id. - if tag_manager.tag_exists(v).await? { - resolve_tag(tag_manager, v).await.map(Some) - } else if let Some(raw_watermark) = v.strip_prefix(WATERMARK_PREFIX) { - let watermark = raw_watermark - .parse::() - .map_err(|e| Error::DataInvalid { - message: format!("{option_name} '{v}' has an invalid watermark value."), - source: Some(Box::new(e)), - })?; - resolve_watermark(snapshot_manager, watermark) - .await - .map(Some) - } else if let Ok(id) = v.parse::() { - snapshot_manager.get_snapshot(id).await.map(Some) - } else { - Err(Error::DataInvalid { - message: format!("{option_name} '{v}' is not a valid tag name or snapshot id."), - source: None, - }) - } - } + Some(TimeTravelSelector::Version { .. }) => unreachable!("scan.version was adapted above"), Some(TimeTravelSelector::SnapshotId { value: v, option_name, @@ -621,6 +624,38 @@ mod tests { ); } + #[tokio::test] + async fn scan_version_overwrites_same_selector_before_validation() { + let (io, path) = setup_watermark_table().await; + let table = make_table(&io, &path, schema_v0()); + let sm = table.snapshot_manager(); + let tm = table.tag_manager(); + tm.create("before", &sm.get_snapshot(1).await.unwrap()) + .await + .unwrap(); + for (version, key, expected) in [ + ("1", "scan.snapshot-id", 1), + ("before", "scan.tag-name", 1), + ("watermark-150", "scan.watermark", 3), + ] { + let opts = options(&[("scan.version", version), (key, "invalid-overridden-value")]); + let snapshot = super::travel_to_snapshot(&sm, &tm, &opts) + .await + .unwrap() + .unwrap(); + assert_eq!(snapshot.id(), expected); + assert_eq!(opts[key], "invalid-overridden-value"); + let traveled = table.copy_with_time_travel(opts).await.unwrap(); + assert_eq!(traveled.travel_snapshot().map(|s| s.id()), Some(expected)); + } + for opts in [ + options(&[("scan.version", "1"), ("scan.tag-name", "before")]), + options(&[("scan.version", "before"), ("scan.snapshot-id", "1")]), + ] { + assert!(super::travel_to_snapshot(&sm, &tm, &opts).await.is_err()); + } + } + #[tokio::test] async fn test_watermark_without_matching_snapshot_fails_at_scan() { let (file_io, table_path) = setup_watermark_table().await; diff --git a/crates/paimon/tests/incremental_batch_scan_test.rs b/crates/paimon/tests/incremental_batch_scan_test.rs index 2b9768cfe..647b32d22 100644 --- a/crates/paimon/tests/incremental_batch_scan_test.rs +++ b/crates/paimon/tests/incremental_batch_scan_test.rs @@ -1017,10 +1017,9 @@ async fn diff_rejects_bucket_rescale_between_snapshots() { ); } -/// An ordinary incremental batch must merge overlapping PK files across commits; -/// the existing per-commit Delta API must still expose both versions. +/// Java packs incremental batch files together but reads every physical event. #[tokio::test] -async fn combined_delta_merges_primary_key_versions_and_preserves_delta_api() { +async fn combined_delta_preserves_primary_key_versions_and_delta_api() { let path = "memory:/incremental_batch/combined_pk"; let (file_io, table) = memory_table(path, pk_schema(&[("changelog-producer", "none")])); setup_dirs(&file_io, path).await; @@ -1038,10 +1037,12 @@ async fn combined_delta_merges_primary_key_versions_and_preserves_delta_api() { assert_eq!( plan.splits().len(), 1, - "overlapping PK files need one merge reader" + "combined delta retains Java batch split packing" ); assert_eq!(plan.splits()[0].data_files().len(), 2); assert!(plan.splits().iter().all(|split| split.snapshot_id() == 2)); + assert!(plan.splits().iter().all(|split| split.is_streaming())); + assert!(!plan.splits()[0].raw_convertible()); let batches: Vec = builder .new_read() .unwrap() @@ -1050,7 +1051,7 @@ async fn combined_delta_merges_primary_key_versions_and_preserves_delta_api() { .try_collect() .await .unwrap(); - assert_eq!(collect_pairs(&batches), vec![(1, 20)]); + assert_eq!(collect_pairs(&batches), vec![(1, 10), (1, 20)]); assert_eq!( read_incremental_pairs(&table, IncrementalScanMode::Delta, 0, 2).await, vec![(1, 10), (1, 20)] @@ -1135,7 +1136,7 @@ async fn combined_delta_skips_overwrite_but_retains_endpoint_metadata() { } #[tokio::test] -async fn combined_delta_merges_add_delete_entries_across_appends() { +async fn combined_delta_rejects_manifest_deletes_in_append_snapshots() { let path = "memory:/incremental_batch/combined_rewrite"; let (file_io, table) = memory_table(path, pk_schema(&[])); setup_dirs(&file_io, path).await; @@ -1155,12 +1156,11 @@ async fn combined_delta_merges_add_delete_entries_across_appends() { .await .unwrap(); let mut messages = second.prepare_commit().await.unwrap(); - let new_name = messages[0].new_files[0].file_name.clone(); messages[0].deleted_files.push(old_file); builder.new_commit().commit(messages).await.unwrap(); // The Rust rewrite writer labels deletions OVERWRITE. Construct an APPEND // metadata fixture over its real ADD/DELETE manifests to exercise the batch - // planner's entry cancellation independently of that writer policy. + // planner's Java ADD-only validation independently of that writer policy. let manager = table.snapshot_manager(); let snapshot = manager.get_snapshot(2).await.unwrap(); let mut metadata = serde_json::to_value(snapshot).unwrap(); @@ -1171,19 +1171,12 @@ async fn combined_delta_merges_add_delete_entries_across_appends() { .write(bytes::Bytes::from(serde_json::to_vec(&metadata).unwrap())) .await .unwrap(); - let plan = table + let result = table .new_read_builder() .new_incremental_scan(IncrementalScanMode::Delta, 0, 2) .plan_combined_delta() - .await - .unwrap(); - let files: Vec<_> = plan - .splits() - .iter() - .flat_map(|split| split.data_files()) - .map(|file| file.file_name.as_str()) - .collect(); - assert_eq!(files, vec![new_name.as_str()]); + .await; + assert!(matches!(result, Err(paimon::Error::DataInvalid { .. }))); } #[tokio::test] @@ -1376,7 +1369,7 @@ async fn combined_delta_preserves_partition_filter_and_projection_across_appends .try_collect::>() .await .unwrap(); - assert_eq!(collect_pairs(&batches), vec![(1, 99)]); + assert_eq!(collect_pairs(&batches), vec![(1, 10), (1, 99)]); } #[tokio::test] @@ -1443,3 +1436,160 @@ async fn combined_delta_rejects_missing_history_and_auto_changelog() { .unwrap(); assert_eq!(collect_pairs(&batches), vec![(3, 30)]); } + +/// Every merge engine must expose L0 events, even when ordinary DV/first-row +/// batch reads hide them. Residual value filters must precede the row limit. +#[tokio::test] +async fn delta_event_reads_ignore_batch_merge_engine_and_endpoint_indexes() { + use paimon::spec::{Datum, PredicateBuilder}; + for engine in ["deduplicate", "partial-update", "aggregation", "first-row"] { + for dv in ["false", "true"] { + if dv == "true" && engine != "deduplicate" { + continue; // Only deduplicate supports deletion vectors. + } + let path = format!("memory:/incremental_batch/events_{engine}_{dv}"); + let mut options = vec![ + ("merge-engine", engine), + ("deletion-vectors.enabled", dv), + ("source.split.target-size", "1b"), + ]; + if engine == "aggregation" { + options.push(("fields.value.aggregate-function", "sum")); + } + let (io, table) = memory_table(&path, pk_schema(&options)); + setup_dirs(&io, &path).await; + persist_table_schema(&io, &path, table.schema()).await; + write_batch(&table, &make_batch(vec![1], vec![10])).await; + write_batch(&table, &make_batch(vec![1], vec![20])).await; + // A latest-state index must not even be opened for event reads. + let manager = table.snapshot_manager(); + let mut metadata = + serde_json::to_value(manager.get_snapshot(2).await.unwrap()).unwrap(); + metadata["indexManifest"] = serde_json::json!("missing-endpoint-index"); + io.new_output(&manager.snapshot_path(2)) + .unwrap() + .write(bytes::Bytes::from(serde_json::to_vec(&metadata).unwrap())) + .await + .unwrap(); + for selected in [None, Some(10), Some(20)] { + let mut builder = table.new_read_builder(); + if let Some(value) = selected { + builder.with_filter( + PredicateBuilder::new(table.schema().fields()) + .equal("value", Datum::Int(value)) + .unwrap(), + ); + builder.with_limit(1); + } + let plan = builder + .new_incremental_scan(IncrementalScanMode::Delta, 0, 2) + .plan_combined_delta() + .await + .unwrap(); + assert!(plan + .splits() + .iter() + .all(|s| s.is_streaming() && s.data_deletion_files().is_none())); + let batches = builder + .new_read() + .unwrap() + .to_arrow(plan.splits()) + .unwrap() + .try_collect::>() + .await + .unwrap(); + let expected = selected.map_or_else(|| vec![(1, 10), (1, 20)], |v| vec![(1, v)]); + assert_eq!(collect_pairs(&batches), expected, "{engine}, dv={dv}"); + // The separate per-commit API follows the same no-DV event contract. + let per_commit = builder + .new_incremental_scan(IncrementalScanMode::Delta, 0, 2) + .plan() + .await + .unwrap(); + assert!(per_commit + .data_splits() + .iter() + .all(|s| s.is_streaming() && s.data_deletion_files().is_none())); + let events = builder + .new_read() + .unwrap() + .to_incremental_arrow(&per_commit) + .unwrap() + .try_collect::>() + .await + .unwrap(); + assert_eq!( + collect_pairs(&events), + expected, + "per-commit {engine}, dv={dv}" + ); + } + } + } +} + +#[tokio::test] +async fn delta_keeps_physical_retractions_and_changelog_row_kinds() { + use arrow_array::StringArray; + let path = "memory:/incremental_batch/retracts"; + let (io, table) = memory_table(path, pk_schema(&[("changelog-producer", "input")])); + setup_dirs(&io, path).await; + persist_table_schema(&io, path, table.schema()).await; + for (value, kind) in [(10, 0), (10, 1), (20, 2), (20, 3)] { + write_batch( + &table, + &make_batch_with_kinds(vec![1], vec![value], vec![kind]), + ) + .await; + } + let builder = table.new_read_builder(); + let combined = builder + .new_incremental_scan(IncrementalScanMode::Delta, 0, 4) + .plan_combined_delta() + .await + .unwrap(); + let batches = builder + .new_read() + .unwrap() + .to_arrow(combined.splits()) + .unwrap() + .try_collect::>() + .await + .unwrap(); + assert_eq!( + collect_pairs(&batches), + vec![(1, 10), (1, 10), (1, 20), (1, 20)] + ); + for mode in [IncrementalScanMode::Delta, IncrementalScanMode::Changelog] { + let plan = builder + .new_incremental_scan(mode, 0, 4) + .plan() + .await + .unwrap(); + assert!(plan.data_splits().iter().all(|s| s.is_streaming())); + let batches = builder + .new_read() + .unwrap() + .to_audit_log_arrow(&plan) + .unwrap() + .try_collect::>() + .await + .unwrap(); + let mut kinds: Vec<_> = batches + .iter() + .flat_map(|batch| { + batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap() + .iter() + .map(|s| s.unwrap().to_owned()) + .collect::>() + }) + .collect(); + kinds.sort(); + assert_eq!(kinds, vec!["+I", "+U", "-D", "-U"]); + } + assert!(read_current_pairs(&table).await.is_empty()); +} diff --git a/docs/src/incremental-reading.md b/docs/src/incremental-reading.md index 36e43db1a..8e91db226 100644 --- a/docs/src/incremental-reading.md +++ b/docs/src/incremental-reading.md @@ -56,7 +56,12 @@ let batches = reader .await?; ``` -`Delta` and `Changelog` return rows from their planned files. `Diff` returns +`Delta` and `Changelog` retain physical events from their planned files, including +repeated keys and retracts, without snapshot deletion vectors or automatic +global-index pruning. Their `DataSplit::is_streaming()` flag also preserves this +contract when passed to the ordinary `TableRead::to_arrow` reader. Combined delta +planning (`plan_combined_delta`) packs the entire window using Java batch split +rules while retaining these event semantics. `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. diff --git a/docs/src/python-binding.md b/docs/src/python-binding.md index 56605eeaa..89c831213 100644 --- a/docs/src/python-binding.md +++ b/docs/src/python-binding.md @@ -158,12 +158,18 @@ batches = rb.new_read().read(plan.splits()) ``` The range is `(start_snapshot_id, end_snapshot_id]`. Only APPEND snapshots -contribute delta manifests. Their entries are merged together before building -splits, so overlapping primary-key versions across those snapshots use one merge -reader. This is an ordinary batch `Plan`, not a per-commit changelog. The end -snapshot must exist and supplies snapshot metadata and deletion vectors, even -when the range produces no splits. Existing builder filters, projections, and -limits also apply. +contribute delta manifests. Files use Java batch split packing, with streaming +read semantics: repeated primary keys and physical retracts remain separate +rows. Readers do not merge these events into the window's final table state. +The end snapshot must exist and supplies snapshot metadata, including for empty +results. Snapshot deletion vectors and automatic global indexes are not applied +to historical events. Builder filters, projections, and limits still apply. + +`split.is_streaming()` identifies this read contract. Exporting its Java binary +encoding requires `split.serialize(allow_streaming=True)` and a decoder that +preserves the streaming flag. Calling `serialize()` without this acknowledgement +rejects streaming splits, protecting older Python decoders that discard the flag. +Batch splits continue to support `serialize()` without arguments. For Data Evolution tables, select half-open row positions or one balanced shard on a scan: From 20e79263b9cb9f6379f25fa49dadd18c3cc29c9b Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Tue, 15 Sep 2026 18:52:39 +0800 Subject: [PATCH 2/4] [test] Align scan planning parity tests with Java incremental semantics --- .../paimon/tests/scan_planning_parity_test.rs | 82 ++++++++++++------- 1 file changed, 51 insertions(+), 31 deletions(-) diff --git a/crates/paimon/tests/scan_planning_parity_test.rs b/crates/paimon/tests/scan_planning_parity_test.rs index 8396b7fe5..78111e39c 100644 --- a/crates/paimon/tests/scan_planning_parity_test.rs +++ b/crates/paimon/tests/scan_planning_parity_test.rs @@ -368,7 +368,7 @@ async fn position_selection_precedes_deletions_and_preserves_historical_reads() } #[tokio::test] -async fn combined_delta_uses_window_end_deletion_vectors_after_repeated_deletes() { +async fn combined_delta_preserves_events_across_repeated_endpoint_deletes() { for bucket_local in ["false", "true"] { let table = evolution_table_with_options( "memory:/planning_parity/delta_dv", @@ -382,7 +382,7 @@ async fn combined_delta_uses_window_end_deletion_vectors_after_repeated_deletes( append_ids(&table, 3, 6).await; delete_ids(&table, &[1, 4]).await; delete_ids(&table, &[5]).await; - for (end, expected) in [(2, vec![3, 4, 5]), (3, vec![3, 5]), (4, vec![3])] { + for end in [2, 3, 4] { let plan = table .new_read_builder() .new_incremental_scan(IncrementalScanMode::Delta, 1, end) @@ -390,8 +390,11 @@ async fn combined_delta_uses_window_end_deletion_vectors_after_repeated_deletes( .await .unwrap(); assert_eq!(plan.snapshot_id(), Some(end)); - assert!(plan.splits().iter().all(|s| s.snapshot_id() == end)); - assert_eq!(read_ids(&table, &plan).await, expected); + assert!(plan.splits().iter().all(|s| s.snapshot_id() == end + && s.is_streaming() + && s.data_deletion_files().is_none())); + // Endpoint deletes affect snapshot state, not historical APPEND events. + assert_eq!(read_ids(&table, &plan).await, vec![3, 4, 5]); } let mut builder = table.new_read_builder(); builder.with_limit(1); @@ -402,7 +405,10 @@ async fn combined_delta_uses_window_end_deletion_vectors_after_repeated_deletes( .plan_combined_delta() .await .unwrap(); - assert_eq!(read_column(&builder, &plan, 0).await, vec![5]); + // The limit is a planning hint; both selected positions belong to one group. + assert_eq!(read_column(&builder, &plan, 0).await, vec![4, 5]); + let current = table.new_read_builder().new_scan().plan().await.unwrap(); + assert_eq!(read_ids(&table, ¤t).await, vec![0, 2, 3]); } } @@ -473,23 +479,16 @@ async fn branch_snapshot_and_combined_delta_use_independent_snapshot_histories() .unwrap(); let branch = table.copy_with_branch("audit").await.unwrap(); append_ids(&table, 10, 11).await; - let overwrite = table.new_write_builder().with_overwrite(); - let mut writer = overwrite.new_write().unwrap(); - writer - .write_arrow_batch(&make_batch(vec![20], vec![200])) - .await - .unwrap(); - overwrite - .new_commit() - .overwrite(writer.prepare_commit().await.unwrap(), None) - .await - .unwrap(); - // Rust intentionally refuses branch writes. Use real data/manifests from - // another commit to model an independent APPEND at branch snapshot 2. + append_ids(&table, 20, 21).await; + // Rust intentionally refuses branch writes. Reuse a real APPEND delta for + // branch snapshot 2, with only the fork's snapshot-1 files as its base. + // Relabeling an OVERWRITE would leave DELETE entries in the delta manifest. let mut metadata = serde_json::to_value(table.snapshot_manager().get_snapshot(3).await.unwrap()).unwrap(); + assert_eq!(metadata["commitKind"], serde_json::json!("APPEND")); metadata["id"] = serde_json::json!(2); - metadata["commitKind"] = serde_json::json!("APPEND"); + metadata["baseManifestList"] = serde_json::json!(snapshot.delta_manifest_list()); + metadata["totalRecordCount"] = serde_json::json!(2); table .file_io() .new_output(&branch.snapshot_manager().snapshot_path(2)) @@ -509,7 +508,7 @@ async fn branch_snapshot_and_combined_delta_use_independent_snapshot_histories() } let plan = branch.new_read_builder().new_scan().plan().await.unwrap(); assert_eq!(plan.snapshot_id(), Some(2)); - assert_eq!(read_ids(&branch, &plan).await, vec![20]); + assert_eq!(read_ids(&branch, &plan).await, vec![0, 20]); assert!(branch .new_read_builder() .new_incremental_scan(IncrementalScanMode::Delta, 1, 3) @@ -519,7 +518,7 @@ async fn branch_snapshot_and_combined_delta_use_independent_snapshot_histories() } #[tokio::test] -async fn combined_delta_value_filter_does_not_resurrect_old_primary_key_versions() { +async fn combined_delta_value_filter_keeps_matching_historical_primary_key_events() { let path = "memory:/planning_parity/pk_filter"; let (io, table) = memory_table( path, @@ -540,7 +539,10 @@ async fn combined_delta_value_filter_does_not_resurrect_old_primary_key_versions .plan_combined_delta() .await .unwrap(); - assert_eq!(read_column(&builder, &plan, 0).await, vec![2]); + assert_eq!(read_column(&builder, &plan, 0).await, vec![1, 2]); + // Snapshot reads still merge versions before applying the value predicate. + let current = builder.new_scan().plan().await.unwrap(); + assert_eq!(read_column(&builder, ¤t, 0).await, vec![2]); } #[tokio::test] @@ -584,7 +586,7 @@ async fn empty_position_plans_preserve_selected_snapshot() { } #[tokio::test] -async fn float_primary_key_plans_merge_signed_zero_versions() { +async fn float_primary_key_snapshot_merges_versions_while_delta_keeps_events() { use arrow_array::{ArrayRef, Float32Array, Float64Array}; use paimon::spec::{DoubleType, FloatType}; for double in [false, true] { @@ -695,16 +697,25 @@ async fn float_primary_key_plans_merge_signed_zero_versions() { } } actual.sort(); - assert_eq!( - actual, + let expected = if combined { vec![ ("+0".into(), 1, 10), + ("-0".into(), 1, 20), ("-0".into(), 1, 200), + ("-0".into(), 2, 30), ("-0".into(), 2, 300), - ("1".into(), 1, 500) - ], - "double={double}, combined={combined}" - ); + ("1".into(), 1, 50), + ("1".into(), 1, 500), + ] + } else { + vec![ + ("+0".into(), 1, 10), + ("-0".into(), 1, 200), + ("-0".into(), 2, 300), + ("1".into(), 1, 500), + ] + }; + assert_eq!(actual, expected, "double={double}, combined={combined}"); } } } @@ -844,7 +855,7 @@ async fn python_dv_manifest_layout_resolves_legacy_canonical_and_external_files( } #[tokio::test] -async fn multiple_dv_index_files_in_one_bucket_keep_all_deletions() { +async fn multiple_dv_index_files_apply_to_snapshot_but_not_delta_events() { let table = evolution_table_with_options( "memory:/planning_parity/multiple_dvs", &[("deletion-vectors.enabled", "true")], @@ -884,6 +895,15 @@ async fn multiple_dv_index_files_in_one_bucket_keep_all_deletions() { } else { builder.new_scan().plan().await.unwrap() }; - assert_eq!(read_ids(&table, &plan).await, vec![0, 2, 3, 5]); + let expected = if combined { + assert!(plan + .splits() + .iter() + .all(|s| s.is_streaming() && s.data_deletion_files().is_none())); + vec![0, 1, 2, 3, 4, 5] + } else { + vec![0, 2, 3, 5] + }; + assert_eq!(read_ids(&table, &plan).await, expected); } } From 2d3d8dcca8dd126f7d01e7e85432b6905ca35077 Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Tue, 15 Sep 2026 19:07:16 +0800 Subject: [PATCH 3/4] [python] Remove the streaming split serialization opt-in --- .../python/python/pypaimon_rust/datafusion.pyi | 2 +- bindings/python/src/read.rs | 16 ++-------------- bindings/python/tests/test_read.py | 15 +++++++-------- docs/src/python-binding.md | 8 +++----- 4 files changed, 13 insertions(+), 28 deletions(-) diff --git a/bindings/python/python/pypaimon_rust/datafusion.pyi b/bindings/python/python/pypaimon_rust/datafusion.pyi index 34a289d45..3eecaf729 100644 --- a/bindings/python/python/pypaimon_rust/datafusion.pyi +++ b/bindings/python/python/pypaimon_rust/datafusion.pyi @@ -42,7 +42,7 @@ class Split: def row_count(self) -> int: ... # Java SplitSerializer v1 binary: a DataSplit (v8), or an IndexedSplit when the split has row ranges. def is_streaming(self) -> bool: ... - def serialize(self, *, allow_streaming: bool = False) -> bytes: ... + def serialize(self) -> bytes: ... class Plan: def snapshot_id(self) -> Optional[int]: diff --git a/bindings/python/src/read.rs b/bindings/python/src/read.rs index 05802ff3e..40eab8b25 100644 --- a/bindings/python/src/read.rs +++ b/bindings/python/src/read.rs @@ -480,20 +480,8 @@ impl PySplit { } /// Serialize to Java SplitSerializer v1, using IndexedSplit for row ranges. - /// Streaming export requires a decoder that preserves change-event semantics. - #[pyo3(signature = (*, allow_streaming=false))] - fn serialize<'py>( - &self, - py: Python<'py>, - allow_streaming: bool, - ) -> PyResult> { - // Older Python decoders silently discard the streaming byte. Require - // acknowledgement before exporting events through that shared codec. - if self.inner.is_streaming() && !allow_streaming { - return Err(PyValueError::new_err( - "Streaming splits require a stream-aware decoder; pass allow_streaming=True", - )); - } + /// Preserves the streaming flag for physical change-event reads. + fn serialize<'py>(&self, py: Python<'py>) -> PyResult> { let bytes = self.inner.serialize_split_v1().map_err(to_py_err)?; Ok(PyBytes::new(py, &bytes)) } diff --git a/bindings/python/tests/test_read.py b/bindings/python/tests/test_read.py index 5aaffffee..35f019e88 100644 --- a/bindings/python/tests/test_read.py +++ b/bindings/python/tests/test_read.py @@ -904,16 +904,15 @@ def test_combined_incremental_plan_preserves_pk_events_and_range(): assert len(plan.splits()) == 1 assert pa.Table.from_batches(builder.new_read().read(plan.splits())).to_pydict() == { "id": [1, 1], "value": [10, 20]} - assert [s.serialize(allow_streaming=True) for s in scan.plan().splits()] == [ - s.serialize(allow_streaming=True) for s in plan.splits()] + assert [s.serialize() for s in scan.plan().splits()] == [ + s.serialize() for s in plan.splits()] assert all(s.is_streaming() for s in plan.splits()) - with pytest.raises(ValueError, match="stream-aware decoder"): - plan.splits()[0].serialize() + assert plan.splits()[0].serialize()[-2] == 1 # Java isStreaming flag import pickle restored = pickle.loads(pickle.dumps(plan.splits()[0])) assert restored.is_streaming() - assert restored.serialize(allow_streaming=True) == ( - plan.splits()[0].serialize(allow_streaming=True)) + assert restored.serialize() == ( + plan.splits()[0].serialize()) selected = builder.new_incremental_scan(1, 2).plan() assert selected.snapshot_id() == 2 assert pa.Table.from_batches(builder.new_read().read(selected.splits())).to_pydict() == { @@ -1007,8 +1006,8 @@ def test_incremental_row_positions_use_combined_delta_batch(): assert plan.snapshot_id() == 2 restored = [pickle.loads(pickle.dumps(split)) for split in plan.splits()] assert pa.Table.from_batches(builder.new_read().read(restored)).column("id").to_pylist() == expected - assert [s.serialize(allow_streaming=True) for s in scan.plan().splits()] == [ - s.serialize(allow_streaming=True) for s in plan.splits()] + assert [s.serialize() for s in scan.plan().splits()] == [ + s.serialize() for s in plan.splits()] builder.with_row_ranges([(1, 3)]).with_limit(2) plan = builder.new_incremental_scan(0, 2).with_row_position_slice(2, 5).plan() assert pa.Table.from_batches(builder.new_read().read(plan.splits())).column("id").to_pylist() == [2, 3] diff --git a/docs/src/python-binding.md b/docs/src/python-binding.md index 89c831213..cf61c7ac1 100644 --- a/docs/src/python-binding.md +++ b/docs/src/python-binding.md @@ -165,11 +165,9 @@ The end snapshot must exist and supplies snapshot metadata, including for empty results. Snapshot deletion vectors and automatic global indexes are not applied to historical events. Builder filters, projections, and limits still apply. -`split.is_streaming()` identifies this read contract. Exporting its Java binary -encoding requires `split.serialize(allow_streaming=True)` and a decoder that -preserves the streaming flag. Calling `serialize()` without this acknowledgement -rejects streaming splits, protecting older Python decoders that discard the flag. -Batch splits continue to support `serialize()` without arguments. +`split.is_streaming()` identifies this read contract. `split.serialize()` exports +both batch and streaming splits to Java binary encoding, preserving the streaming +flag. For Data Evolution tables, select half-open row positions or one balanced shard on a scan: From ba637264b8607aa054a26759becd7fffdc270955 Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Tue, 15 Sep 2026 19:43:25 +0800 Subject: [PATCH 4/4] [ci] Fix time-travel assertions and scope Lumina test targets --- .github/workflows/ci.yml | 4 +- crates/integration_tests/tests/read_tables.rs | 57 ++++++++-------- .../datafusion/tests/read_tables.rs | 65 ++++++++++--------- crates/paimon/src/table/table_scan.rs | 5 +- 4 files changed, 72 insertions(+), 59 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 217fccafa..5f2336a38 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -262,7 +262,7 @@ jobs: - name: Core Lumina Native Build Test if: matrix.suite == 'lumina' run: > - cargo test --locked -p paimon + cargo test --locked -p paimon --lib table::lumina_index_build_builder::tests::test_execute_writes_lumina_index_manifest --features fulltext,vortex -- --ignored --exact @@ -270,7 +270,7 @@ jobs: - name: DataFusion Lumina Build Query E2E Test if: matrix.suite == 'lumina' run: > - cargo test --locked -p paimon-datafusion + cargo test --locked -p paimon-datafusion --test read_tables --features vortex vector_search_tests::test_lumina_build_then_vector_search_query -- --ignored --exact diff --git a/crates/integration_tests/tests/read_tables.rs b/crates/integration_tests/tests/read_tables.rs index f867f7ccb..7b76612af 100644 --- a/crates/integration_tests/tests/read_tables.rs +++ b/crates/integration_tests/tests/read_tables.rs @@ -2957,34 +2957,39 @@ async fn test_time_travel_conflicting_selectors_fail() { let catalog = create_file_system_catalog(); let table = get_table_from_catalog(&catalog, "time_travel_table").await; - let conflicted = table.copy_with_options(HashMap::from([ - ("scan.version".to_string(), "snapshot1".to_string()), - ("scan.timestamp-millis".to_string(), "1234".to_string()), - ])); - - let plan_err = conflicted - .new_read_builder() - .new_scan() - .plan() - .await - .expect_err("conflicting time-travel selectors should fail"); + // Java resolves scan.version before validating conflicts with other selectors. + for (version, selector) in [ + ("snapshot1", "scan.tag-name"), + ("1", "scan.snapshot-id"), + ("watermark-1", "scan.watermark"), + ] { + let conflicted = table.copy_with_options(HashMap::from([ + ("scan.version".to_string(), version.to_string()), + ("scan.timestamp-millis".to_string(), "1234".to_string()), + ])); - match plan_err { - Error::DataInvalid { message, .. } => { - assert!( - message.contains("Only one time-travel selector may be set"), - "unexpected conflict error: {message}" - ); - assert!( - message.contains("scan.version"), - "conflict error should mention scan.version: {message}" - ); - assert!( - message.contains("scan.timestamp-millis"), - "conflict error should mention scan.timestamp-millis: {message}" - ); + let plan_err = conflicted + .new_read_builder() + .new_scan() + .plan() + .await + .expect_err("conflicting time-travel selectors should fail"); + + match plan_err { + Error::DataInvalid { message, .. } => { + assert!( + message.contains("Only one time-travel selector may be set"), + "unexpected conflict error for version {version}: {message}" + ); + for key in [selector, "scan.timestamp-millis"] { + assert!( + message.contains(key), + "conflict error should mention {key}: {message}" + ); + } + } + other => panic!("unexpected error for version {version}: {other:?}"), } - other => panic!("unexpected error: {other:?}"), } } diff --git a/crates/integrations/datafusion/tests/read_tables.rs b/crates/integrations/datafusion/tests/read_tables.rs index e2d9189b3..f867b07d8 100644 --- a/crates/integrations/datafusion/tests/read_tables.rs +++ b/crates/integrations/datafusion/tests/read_tables.rs @@ -1008,38 +1008,45 @@ async fn time_travel_schema_evolution() { #[tokio::test] async fn test_time_travel_conflicting_selectors_fail() { - // When both scan.version and scan.timestamp-millis are set on the same - // provider, Paimon rejects the combination at scan time. - let provider = create_provider_with_options( - "time_travel_table", - HashMap::from([ - ("scan.version".to_string(), "1".to_string()), - ("scan.timestamp-millis".to_string(), "1234".to_string()), - ]), - ) - .await; + // Java resolves scan.version before validating conflicts with other selectors. + for (version, selector) in [ + ("snapshot1", "scan.tag-name"), + ("1", "scan.snapshot-id"), + ("watermark-1", "scan.watermark"), + ] { + let provider = create_provider_with_options( + "time_travel_table", + HashMap::from([ + ("scan.version".to_string(), version.to_string()), + ("scan.timestamp-millis".to_string(), "1234".to_string()), + ]), + ) + .await; - let ctx = create_context().await; - ctx.register_temp_table("paimon.default.time_travel_table", Arc::new(provider)) - .expect("Failed to register temp table"); + let ctx = create_context().await; + ctx.register_temp_table("paimon.default.time_travel_table", Arc::new(provider)) + .expect("Failed to register temp table"); - let err = ctx - .sql("SELECT id, name FROM paimon.default.time_travel_table") - .await - .expect("query should parse") - .collect() - .await - .expect_err("conflicting time-travel selectors should fail"); + let err = ctx + .sql("SELECT id, name FROM paimon.default.time_travel_table") + .await + .expect("query should parse") + .collect() + .await + .expect_err("conflicting time-travel selectors should fail"); - let message = err.to_string(); - assert!( - message.contains("Only one time-travel selector may be set"), - "unexpected conflict error: {message}" - ); - assert!( - message.contains("scan.version"), - "conflict error should mention scan.version: {message}" - ); + let message = err.to_string(); + assert!( + message.contains("Only one time-travel selector may be set"), + "unexpected conflict error for version {version}: {message}" + ); + for key in [selector, "scan.timestamp-millis"] { + assert!( + message.contains(key), + "conflict error should mention {key}: {message}" + ); + } + } } #[tokio::test] diff --git a/crates/paimon/src/table/table_scan.rs b/crates/paimon/src/table/table_scan.rs index 9f2bd6f80..42e008d9e 100644 --- a/crates/paimon/src/table/table_scan.rs +++ b/crates/paimon/src/table/table_scan.rs @@ -1189,8 +1189,9 @@ impl<'a> PaimonTableScan<'a> { /// Plan the full scan: resolve snapshot (via options or latest), then read manifests and build DataSplits. /// /// Time travel is resolved from table options: - /// - only one of `scan.version`, `scan.timestamp-millis`, `scan.watermark`, - /// `scan.snapshot-id`, `scan.tag-name` may be set + /// - `scan.version` is resolved first, overwriting the same selector kind; + /// only one of `scan.timestamp-millis`, `scan.watermark`, `scan.snapshot-id`, + /// `scan.tag-name` may remain after resolution /// - `scan.version` → tag name (if exists) → `watermark-` → snapshot /// id (if parseable) → error (ambiguous by design, like SQL `VERSION AS OF`) /// - `scan.snapshot-id` → snapshot id only (never a tag lookup)