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/bindings/python/python/pypaimon_rust/datafusion.pyi b/bindings/python/python/pypaimon_rust/datafusion.pyi index 40b75c1ea..3eecaf729 100644 --- a/bindings/python/python/pypaimon_rust/datafusion.pyi +++ b/bindings/python/python/pypaimon_rust/datafusion.pyi @@ -41,6 +41,7 @@ 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 is_streaming(self) -> bool: ... def serialize(self) -> bytes: ... class Plan: @@ -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..40eab8b25 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,9 +474,13 @@ 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`. + /// 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. + /// 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 8404823ee..35f019e88 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,16 @@ 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]} + "id": [1, 1], "value": [10, 20]} 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()) + 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() == ( + 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() == { @@ -981,7 +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() for s in scan.plan().splits()] == [s.serialize() 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/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/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..42e008d9e 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, } } @@ -1161,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) @@ -1284,7 +1313,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 +1387,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 +1557,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 +1587,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 +1606,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 +1628,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 +1655,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 +1878,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 +1923,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 +1985,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 +2100,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 +2153,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 +2332,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 +3524,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 +3556,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/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); } } 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..cf61c7ac1 100644 --- a/docs/src/python-binding.md +++ b/docs/src/python-binding.md @@ -158,12 +158,16 @@ 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. `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: