This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git


The following commit(s) were added to refs/heads/main by this push:
     new d6bc0024 Expose combined changelog scans to Python (#894)
d6bc0024 is described below

commit d6bc002403fc7f421a1ac50cabaa759a53bc107d
Author: Jingsong Lee <[email protected]>
AuthorDate: Mon Sep 21 12:34:47 2026 +0800

    Expose combined changelog scans to Python (#894)
---
 .../python/python/pypaimon_rust/datafusion.pyi     | 20 ++++-
 bindings/python/src/read.rs                        | 65 +++++++++++----
 bindings/python/tests/test_read.py                 | 75 +++++++++++++++++
 crates/paimon/src/table/incremental_scan.rs        | 36 ++++++++
 crates/paimon/tests/incremental_batch_scan_test.rs | 95 ++++++++++++++++++++++
 docs/src/python-binding.md                         | 24 ++++--
 6 files changed, 289 insertions(+), 26 deletions(-)

diff --git a/bindings/python/python/pypaimon_rust/datafusion.pyi 
b/bindings/python/python/pypaimon_rust/datafusion.pyi
index a91ad1da..dbbdaba6 100644
--- a/bindings/python/python/pypaimon_rust/datafusion.pyi
+++ b/bindings/python/python/pypaimon_rust/datafusion.pyi
@@ -92,6 +92,9 @@ class ReadBuilder:
         """Set a scan-planning row-limit hint, not an exact cap: a matching 
split is
         returned whole. Apply application-level limiting for an exact bound."""
         ...
+    def with_include_row_kind(self, include: bool) -> "ReadBuilder":
+        """Include a leading ``rowkind`` string column in native read 
results."""
+        ...
     def with_blob_parallelism(self, blob_parallelism: int) -> "ReadBuilder":
         """Set the maximum number of concurrent BLOB range reads. Must be 
positive."""
         ...
@@ -100,11 +103,20 @@ class ReadBuilder:
         """Set Data Evolution row ranges. Empty selects no rows; format tables 
are unsupported."""
         ...
     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, preserving physical 
change events.
+    def new_incremental_scan(
+        self,
+        start_snapshot_id: int,
+        end_snapshot_id: int,
+        mode: str = "delta",
+    ) -> TableScan:
+        """Plan incremental files in (start, end] as one native plan.
 
-        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.
+        ``mode`` accepts ``delta``, ``changelog`` or ``auto``. Delta reads 
APPEND
+        manifests; changelog reads physical changelog manifests; auto follows 
the
+        table's changelog-producer option. Diff is not representable as one
+        split list and is rejected. 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.
         """
         ...
     def new_read(self) -> "TableRead": ...
diff --git a/bindings/python/src/read.rs b/bindings/python/src/read.rs
index 56516bb6..120fdbf2 100644
--- a/bindings/python/src/read.rs
+++ b/bindings/python/src/read.rs
@@ -409,7 +409,7 @@ impl PyReadBuilder {
             filter: self.filter.clone(),
             row_ranges: self.row_ranges.clone(),
             case_sensitive: self.case_sensitive,
-            incremental_range: None,
+            incremental_scan: None,
             row_position_slice: None,
             row_position_shard: None,
             chunk_shuffle: None,
@@ -417,12 +417,23 @@ impl PyReadBuilder {
         }
     }
 
-    /// Plan APPEND deltas in (start_snapshot_id, end_snapshot_id] as one 
batch.
-    /// Primary-key versions are grouped across all selected snapshots.
-    fn new_incremental_scan(&self, start_snapshot_id: i64, end_snapshot_id: 
i64) -> PyTableScan {
+    /// Plan physical changes in (start_snapshot_id, end_snapshot_id] as one
+    /// ordinary split plan. Mode is `delta` by default; `changelog` reads
+    /// changelog manifest files and `auto` follows the table's producer.
+    #[pyo3(signature = (start_snapshot_id, end_snapshot_id, mode = "delta"))]
+    fn new_incremental_scan(
+        &self,
+        start_snapshot_id: i64,
+        end_snapshot_id: i64,
+        mode: &str,
+    ) -> PyResult<PyTableScan> {
         let mut scan = self.new_scan();
-        scan.incremental_range = Some((start_snapshot_id, end_snapshot_id));
-        scan
+        scan.incremental_scan = Some(PyIncrementalScan {
+            start_snapshot_id,
+            end_snapshot_id,
+            mode: parse_incremental_scan_mode(mode)?,
+        });
+        Ok(scan)
     }
 
     fn new_read(&self) -> PyTableRead {
@@ -448,7 +459,7 @@ pub struct PyTableScan {
     filter: Option<Predicate>,
     row_ranges: Option<Vec<RowRange>>,
     case_sensitive: bool,
-    incremental_range: Option<(i64, i64)>,
+    incremental_scan: Option<PyIncrementalScan>,
     row_position_slice: Option<(u64, u64)>,
     row_position_shard: Option<(u64, u64)>,
     chunk_shuffle: Option<PyChunkShuffle>,
@@ -461,6 +472,27 @@ struct PyChunkShuffle {
     chunk_size: u64,
 }
 
+#[derive(Clone, Copy)]
+struct PyIncrementalScan {
+    start_snapshot_id: i64,
+    end_snapshot_id: i64,
+    mode: IncrementalScanMode,
+}
+
+fn parse_incremental_scan_mode(mode: &str) -> PyResult<IncrementalScanMode> {
+    match mode.to_ascii_lowercase().as_str() {
+        "delta" => Ok(IncrementalScanMode::Delta),
+        "changelog" => Ok(IncrementalScanMode::Changelog),
+        "auto" => Ok(IncrementalScanMode::Auto),
+        "diff" => Err(PyValueError::new_err(
+            "incremental mode 'diff' requires before/after split pairs and is 
not supported by TableScan.plan()",
+        )),
+        _ => Err(PyValueError::new_err(format!(
+            "unsupported incremental scan mode '{mode}'; expected delta, 
changelog or auto"
+        ))),
+    }
+}
+
 impl PyTableScan {
     fn core_scan(&self) -> PyResult<paimon::table::TableScan<'_>> {
         let mut scan = self.read_builder()?.new_scan();
@@ -489,10 +521,9 @@ impl PyTableScan {
         &self,
         start: i64,
         end: i64,
+        mode: IncrementalScanMode,
     ) -> PyResult<paimon::table::IncrementalScan<'_>> {
-        let mut scan =
-            self.read_builder()?
-                .new_incremental_scan(IncrementalScanMode::Delta, start, end);
+        let mut scan = self.read_builder()?.new_incremental_scan(mode, start, 
end);
         if let Some((start, end)) = self.row_position_slice {
             scan = scan
                 .with_row_position_slice(start, end)
@@ -591,11 +622,15 @@ impl PyTableScan {
     fn plan(&self, py: Python<'_>) -> PyResult<PyPlan> {
         py.detach(|| {
             runtime().block_on(async {
-                let plan = match self.incremental_range {
-                    Some((start, end)) => {
-                        self.core_incremental_scan(start, end)?
-                            .plan_combined_delta()
-                            .await
+                let plan = match self.incremental_scan {
+                    Some(incremental) => {
+                        self.core_incremental_scan(
+                            incremental.start_snapshot_id,
+                            incremental.end_snapshot_id,
+                            incremental.mode,
+                        )?
+                        .plan_combined()
+                        .await
                     }
                     None => self.core_scan()?.plan().await,
                 };
diff --git a/bindings/python/tests/test_read.py 
b/bindings/python/tests/test_read.py
index 02eb9dc2..a6240546 100644
--- a/bindings/python/tests/test_read.py
+++ b/bindings/python/tests/test_read.py
@@ -15,6 +15,7 @@
 # specific language governing permissions and limitations
 # under the License.
 
+import json
 import pickle
 import tempfile
 
@@ -1152,6 +1153,80 @@ def 
test_combined_incremental_retains_predicate_and_projection():
         assert result.to_pydict() == {"id": [3]}
 
 
+def _native_split_file_names(splits):
+    names = []
+    for split in splits:
+        _, (state,) = split.__reduce__()
+        payload = json.loads(bytes(state))
+        names.extend(file_["_FILE_NAME"] for file_ in payload["data_files"])
+    return names
+
+
+def _make_input_changelog_table(warehouse):
+    ctx = SQLContext()
+    ctx.register_catalog("paimon", {"warehouse": warehouse})
+    ctx.sql("CREATE SCHEMA paimon.cldb")
+    ctx.sql("""CREATE TABLE paimon.cldb.t (id INT, value STRING, PRIMARY KEY 
(id))
+        WITH ('bucket' = '1', 'changelog-producer' = 'input')""")
+    ctx.sql("INSERT INTO paimon.cldb.t VALUES (1, 'a'), (2, 'b')")
+    ctx.sql("INSERT INTO paimon.cldb.t VALUES (3, 'c')")
+    return PaimonCatalog({"warehouse": warehouse}).get_table("cldb.t")
+
+
+def test_incremental_changelog_and_auto_plan_physical_changelog_files():
+    with tempfile.TemporaryDirectory() as warehouse:
+        table = _make_input_changelog_table(warehouse)
+        builder = table.new_read_builder().with_include_row_kind(True)
+
+        delta = builder.new_incremental_scan(0, 2).plan()
+        assert all(name.startswith("data-")
+                   for name in _native_split_file_names(delta.splits()))
+
+        for mode in ("changelog", "CHANGELOG", "auto"):
+            plan = builder.new_incremental_scan(0, 2, mode).plan()
+            assert plan.snapshot_id() == 2
+            assert plan.splits()
+            assert all(split.is_streaming() for split in plan.splits())
+            assert all(name.startswith("changelog-")
+                       for name in _native_split_file_names(plan.splits()))
+            actual = pa.Table.from_batches(
+                builder.new_read().read(plan.splits())).to_pydict()
+            assert actual == {
+                "rowkind": ["+I", "+I", "+I"],
+                "id": [1, 2, 3],
+                "value": ["a", "b", "c"],
+            }
+
+        second = builder.new_incremental_scan(1, 2, "changelog").plan()
+        assert pa.Table.from_batches(
+            builder.new_read().read(second.splits())).to_pydict() == {
+                "rowkind": ["+I"], "id": [3], "value": ["c"]}
+        empty = builder.new_incremental_scan(2, 2, "changelog").plan()
+        assert empty.snapshot_id() == 2
+        assert empty.splits() == []
+
+
+def test_incremental_changelog_keeps_filter_projection_and_validates_mode():
+    with tempfile.TemporaryDirectory() as warehouse:
+        table = _make_input_changelog_table(warehouse)
+        builder = (table.new_read_builder()
+                   .with_projection(["id"])
+                   .with_filter({
+                       "method": "greaterThan",
+                       "field": "id",
+                       "literals": [1],
+                   })
+                   .with_include_row_kind(True))
+        plan = builder.new_incremental_scan(0, 2, "changelog").plan()
+        assert pa.Table.from_batches(
+            builder.new_read().read(plan.splits())).to_pydict() == {
+                "rowkind": ["+I", "+I"], "id": [2, 3]}
+
+        for mode, message in (("diff", "before/after"), ("unknown", 
"expected")):
+            with pytest.raises(ValueError, match=message):
+                table.new_read_builder().new_incremental_scan(0, 2, mode)
+
+
 def _make_de_position_table(warehouse):
     ctx = SQLContext()
     ctx.register_catalog("paimon", {"warehouse": warehouse})
diff --git a/crates/paimon/src/table/incremental_scan.rs 
b/crates/paimon/src/table/incremental_scan.rs
index 19696462..2f4a225a 100644
--- a/crates/paimon/src/table/incremental_scan.rs
+++ b/crates/paimon/src/table/incremental_scan.rs
@@ -296,6 +296,42 @@ impl<'a> IncrementalScan<'a> {
         }
     }
 
+    /// Plan a Delta or Changelog range as one ordinary [`Plan`].
+    ///
+    /// This is the bridge used by readers which already consume
+    /// [`DataSplit`]s. Delta keeps its cross-snapshot packing semantics;
+    /// Changelog preserves physical changelog files and row kinds in snapshot
+    /// order. `Auto` resolves from `changelog-producer`. Diff cannot be
+    /// represented by an ordinary split list because each unit contains a
+    /// before/after pair.
+    pub async fn plan_combined(&self) -> crate::Result<Plan> {
+        match self.resolve_mode() {
+            IncrementalScanMode::Delta => self.plan_combined_delta().await,
+            IncrementalScanMode::Changelog => {
+                let incremental = self.plan().await?;
+                let mut splits = 
Vec::with_capacity(incremental.splits().len());
+                for split in incremental.splits() {
+                    match split {
+                        IncrementalSplit::Data(split) => 
splits.push(split.clone()),
+                        IncrementalSplit::DiffPair { .. } => {
+                            return Err(crate::Error::UnexpectedError {
+                                message: "DiffPair appeared in a Changelog 
incremental plan"
+                                    .to_string(),
+                                source: None,
+                            });
+                        }
+                    }
+                }
+                Ok(Plan::new(splits).with_snapshot_id(self.end_inclusive))
+            }
+            IncrementalScanMode::Diff => Err(crate::Error::Unsupported {
+                message: "Combined incremental planning does not support Diff 
mode; Diff requires before/after split pairs"
+                    .to_string(),
+            }),
+            IncrementalScanMode::Auto => unreachable!("Auto must resolve 
before planning"),
+        }
+    }
+
     /// Plan APPEND deltas with batch split packing and streaming read 
semantics.
     /// Each physical change is retained, including repeated keys and retracts.
     ///
diff --git a/crates/paimon/tests/incremental_batch_scan_test.rs 
b/crates/paimon/tests/incremental_batch_scan_test.rs
index 647b32d2..4e11ada3 100644
--- a/crates/paimon/tests/incremental_batch_scan_test.rs
+++ b/crates/paimon/tests/incremental_batch_scan_test.rs
@@ -367,6 +367,101 @@ async fn auto_uses_changelog_when_producer_is_input() {
     assert_eq!(auto, vec![(1, 10), (1, 20)]);
 }
 
+#[tokio::test]
+async fn combined_changelog_plan_keeps_physical_files_kinds_and_snapshot() {
+    use arrow_array::StringArray;
+
+    let table_path = "memory:/incremental_batch/combined_changelog";
+    let (file_io, table) = memory_table(
+        table_path,
+        pk_schema(&[
+            ("changelog-producer", "input"),
+            ("merge-engine", "deduplicate"),
+            ("bucket", "1"),
+        ]),
+    );
+    setup_dirs(&file_io, table_path).await;
+    persist_table_schema(&file_io, table_path, table.schema()).await;
+
+    write_batch(
+        &table,
+        &make_batch_with_kinds(vec![1, 1], vec![10, 20], vec![0, 2]),
+    )
+    .await;
+
+    for mode in [IncrementalScanMode::Changelog, IncrementalScanMode::Auto] {
+        let plan = table
+            .new_read_builder()
+            .new_incremental_scan(mode, 0, 1)
+            .plan_combined()
+            .await
+            .unwrap();
+        assert_eq!(plan.snapshot_id(), Some(1));
+        assert!(!plan.splits().is_empty());
+        assert!(plan.splits().iter().all(|split| split.is_streaming()));
+        assert!(plan.splits().iter().all(|split| split
+            .data_files()
+            .iter()
+            .all(|file| file.file_name.starts_with("changelog-"))));
+
+        let read = table.new_read_builder().new_read().unwrap();
+        let kind_batches: Vec<RecordBatch> = read
+            .to_arrow_with_row_kind(plan.splits())
+            .unwrap()
+            .try_collect()
+            .await
+            .unwrap();
+        let mut kinds: Vec<_> = kind_batches
+            .iter()
+            .flat_map(|batch| {
+                batch
+                    .column(0)
+                    .as_any()
+                    .downcast_ref::<StringArray>()
+                    .unwrap()
+                    .iter()
+                    .map(|kind| kind.unwrap().to_owned())
+                    .collect::<Vec<_>>()
+            })
+            .collect();
+        kinds.sort();
+        assert_eq!(kinds, vec!["+I", "+U"]);
+
+        let batches: Vec<RecordBatch> = read
+            .to_arrow(plan.splits())
+            .unwrap()
+            .try_collect()
+            .await
+            .unwrap();
+        assert_eq!(collect_pairs(&batches), vec![(1, 10), (1, 20)]);
+    }
+}
+
+#[tokio::test]
+async fn combined_incremental_plan_rejects_diff_pairs() {
+    let table_path = "memory:/incremental_batch/combined_diff_rejected";
+    let (file_io, table) = memory_table(
+        table_path,
+        pk_schema(&[("merge-engine", "deduplicate"), ("bucket", "1")]),
+    );
+    setup_dirs(&file_io, table_path).await;
+    persist_table_schema(&file_io, table_path, table.schema()).await;
+    write_batch(&table, &make_batch(vec![1], vec![10])).await;
+    write_batch(&table, &make_batch(vec![1], vec![20])).await;
+
+    let error = table
+        .new_read_builder()
+        .new_incremental_scan(IncrementalScanMode::Diff, 1, 2)
+        .plan_combined()
+        .await
+        .unwrap_err();
+    assert!(matches!(
+        error,
+        paimon::Error::Unsupported { ref message }
+            if message.contains("before/after split pairs")
+    ));
+}
+
 /// Partition filter from ReadBuilder is pushed into the changelog plan path.
 #[tokio::test]
 async fn 
incremental_changelog_scan_applies_partition_filter_from_read_builder() {
diff --git a/docs/src/python-binding.md b/docs/src/python-binding.md
index 0c6d287e..318cfb7f 100644
--- a/docs/src/python-binding.md
+++ b/docs/src/python-binding.md
@@ -157,13 +157,23 @@ plan = rb.new_incremental_scan(2, 5).plan()
 batches = rb.new_read().read(plan.splits())
 ```
 
-The range is `(start_snapshot_id, end_snapshot_id]`. Only APPEND snapshots
-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.
+The range is `(start_snapshot_id, end_snapshot_id]`. The default `delta` mode
+uses APPEND delta manifests. Use `changelog` to read physical changelog
+manifests, or `auto` to follow the table's `incremental-between` option:
+
+```python
+plan = rb.new_incremental_scan(2, 5, "changelog").plan()
+batches = rb.with_include_row_kind(True).new_read().read(plan.splits())
+```
+
+Explicit `diff` mode is rejected because a diff contains before/after split
+pairs and cannot be represented by the ordinary `Plan.splits()` contract.
+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

Reply via email to