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