This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new 0665d4757b [Python] Use native planning and reads for streaming
changelog (#10036)
0665d4757b is described below
commit 0665d4757b414cf2af51763837ba9a2738a15143
Author: Jingsong Lee <[email protected]>
AuthorDate: Mon Sep 21 13:41:28 2026 +0800
[Python] Use native planning and reads for streaming changelog (#10036)
---
paimon-python/README.md | 34 ++++---
paimon-python/pypaimon/read/native_plan.py | 14 ++-
.../pypaimon/read/streaming_table_scan.py | 18 +++-
.../tests/native_plan_capabilities_test.py | 18 +++-
.../tests/native_plan_dynamic_bucket_test.py | 43 ++++++++-
.../pypaimon/tests/native_plan_expanded_test.py | 16 +++-
.../pypaimon/tests/native_plan_incremental_test.py | 102 +++++++++++++++++++++
.../tests/native_plan_materialized_pk_test.py | 37 +++++++-
.../tests/native_plan_resolved_schema_test.py | 16 +++-
paimon-python/pypaimon/tests/native_plan_test.py | 22 +++++
.../pypaimon/tests/streaming_table_scan_test.py | 47 ++++++++++
11 files changed, 336 insertions(+), 31 deletions(-)
diff --git a/paimon-python/README.md b/paimon-python/README.md
index 7db4c25386..cb1836fc57 100644
--- a/paimon-python/README.md
+++ b/paimon-python/README.md
@@ -120,8 +120,9 @@ object-store requests.
# Native scan planning
-PyPaimon can plan splits with the optional `pypaimon-rust` package while
retaining
-the Python reader:
+PyPaimon can plan splits with the optional `pypaimon-rust` package. Planning
and
+reading are independently selectable, so native plans can still use the Python
+reader:
```python
native_table = table.copy({"scan.native-plan.enabled": "true"})
@@ -147,10 +148,11 @@ plan = builder.new_scan().plan()
rows = builder.new_read().to_arrow(plan.splits())
```
-Native reads return PyArrow batches through the Arrow C Data interface. They
-currently require untouched splits produced by the native planner and top-level
-projection. Query authorization, nested projection, and row-kind output retain
-the Python reader. For both materialized
+Native reads return PyArrow batches through the Arrow C Data interface. Native
+planner handles are used directly when available; splits refined by Python
+index or shuffle logic are serialized through the stable split contract and
+reconstructed by Rust. Nested projection and row-kind output are supported.
+Query authorization retains the Python reader. For both materialized
`to_arrow()` and streaming `to_arrow_batch_reader()` reads, the effective split
parallelism (the method argument, `read.parallelism`, or the automatic default)
runs independent Rust readers. Splits stay in input order but contiguous groups
@@ -205,11 +207,17 @@ counts can differ between shards. Limits are applied
after shard/slice selection
Timestamp incremental scans require `ReadBuilder.new_incremental_scan()` and
stream-aware splits exposing `Split.is_streaming()`. Python resolves
`(start_timestamp, end_timestamp]` to snapshot IDs; Rust packs the selected
APPEND
-deltas into one plan. Like Java, readers retain physical change events,
including
-repeated primary keys and retracts across commits. They do not merge the window
-into a final table state or apply endpoint deletion vectors or global indexes.
-Other commit kinds are excluded; the ending snapshot still supplies plan
metadata.
-Rebuild development wheels from Rust main to obtain this contract.
+deltas into one plan. Continuous streaming uses the same native path for
initial
+and delta frames. When `changelog-producer` is enabled, follow-up frames
request
+Rust's explicit `changelog` mode and read the physical changelog manifests.
+OVERWRITE changelog frames retain per-snapshot Python planning because Java
+streaming reads them while Java and Rust range-based incremental scans skip
+OVERWRITE; the resulting splits can still use native reads.
+Like Java, readers retain physical change events, including repeated primary
+keys and retracts across commits. They do not merge the window into a final
+table state or apply endpoint deletion vectors or global indexes. Other commit
+kinds are excluded; the ending snapshot still supplies plan metadata. Rebuild
+development wheels from Rust main to obtain this contract.
`scan.version` supports tags, snapshot IDs and `watermark-<value>`, resolving
tags
first and using the historical schema. Ordinary postpone-bucket batch scans can
@@ -243,7 +251,9 @@ index reader, preserving merge-required splits and the
selected snapshot.
Query authorization, first-row plans mixing L0 with merge-required
materialized files,
and precomputed primary-key global-index results still use the Python planner.
-Continuous streaming and write planning also retain their Python entrypoints.
+Continuous streaming retains Python polling and resume handling, while its
+initial, delta and physical-changelog frames can use native planning and reads.
+Write planning retains its Python entrypoint.
Native planning remains optional and is disabled by default.
# Coalesced BLOB reads
diff --git a/paimon-python/pypaimon/read/native_plan.py
b/paimon-python/pypaimon/read/native_plan.py
index 5520f9c165..4388716154 100644
--- a/paimon-python/pypaimon/read/native_plan.py
+++ b/paimon-python/pypaimon/read/native_plan.py
@@ -370,6 +370,7 @@ def native_plan(
projection: Optional[List[str]] = None,
row_ranges: Optional[List[Tuple[int, int]]] = None,
incremental_range: Optional[Tuple[int, int]] = None,
+ incremental_mode: str = 'delta',
row_position_slice: Optional[Tuple[int, int]] = None,
row_position_shard: Optional[Tuple[int, int]] = None,
chunk_shuffle: Optional[Tuple[int, int]] = None,
@@ -379,6 +380,8 @@ def native_plan(
Native conversion or planning failures are handled by TableScan, which
falls back to the Python planner.
"""
+ if incremental_range is None and incremental_mode != 'delta':
+ raise ValueError('incremental_mode requires incremental_range')
if not native_runtime_available():
raise RuntimeError(
"scan.native-plan.enabled needs pypaimon-rust>=0.3.0 (split
planning API)")
@@ -386,8 +389,15 @@ def native_plan(
_native_read_builder(table), predicate, limit, projection)
if row_ranges is not None:
builder = builder.with_row_ranges(row_ranges)
- scan = (builder.new_scan() if incremental_range is None
- else builder.new_incremental_scan(*incremental_range))
+ if incremental_range is None:
+ scan = builder.new_scan()
+ elif incremental_mode == 'delta':
+ # Keep the two-argument call compatible with runtimes predating the
+ # explicit mode API. Non-delta modes require the new binding.
+ scan = builder.new_incremental_scan(*incremental_range)
+ else:
+ scan = builder.new_incremental_scan(
+ *incremental_range, incremental_mode)
if row_position_slice is not None:
scan = scan.with_row_position_slice(*row_position_slice)
if row_position_shard is not None:
diff --git a/paimon-python/pypaimon/read/streaming_table_scan.py
b/paimon-python/pypaimon/read/streaming_table_scan.py
index 7cd19ee019..2481293a8c 100644
--- a/paimon-python/pypaimon/read/streaming_table_scan.py
+++ b/paimon-python/pypaimon/read/streaming_table_scan.py
@@ -366,8 +366,9 @@ class AsyncStreamingTableScan:
return self._create_plan_from_manifests(manifest_files, snapshot.id)
def _try_native_plan(self, expected_snapshot_id: int,
- incremental_range=None) -> Optional[Plan]:
- """Plan an initial or delta streaming frame with pypaimon-rust.
+ incremental_range=None,
+ incremental_mode='delta') -> Optional[Plan]:
+ """Plan an initial, delta or changelog streaming frame with
pypaimon-rust.
An arbitrary Python bucket predicate cannot be represented by the Rust
planner, so sharded stream consumers retain the Python plan. Any native
@@ -394,6 +395,7 @@ class AsyncStreamingTableScan:
[field.name for field in self._read_type]
if self._read_type is not None else None),
incremental_range=incremental_range,
+ incremental_mode=incremental_mode,
)
if plan.snapshot_id != expected_snapshot_id:
logging.warning(
@@ -410,6 +412,18 @@ class AsyncStreamingTableScan:
def _create_changelog_plan(self, snapshot: Snapshot) -> Plan:
"""Read from changelog_manifest_list
(changelog-producer=input/full-compaction/lookup)."""
+ # Java's range-based incremental CHANGELOG scan skips OVERWRITE,
+ # whereas ChangelogFollowUpScanner reads any follow-up snapshot with a
+ # changelog manifest. Rust's incremental API implements the former, so
+ # keep OVERWRITE on the per-snapshot Python planner.
+ if snapshot.commit_kind != 'OVERWRITE':
+ plan = self._try_native_plan(
+ snapshot.id,
+ incremental_range=(snapshot.id - 1, snapshot.id),
+ incremental_mode='changelog',
+ )
+ if plan is not None:
+ return plan
manifest_files = self._manifest_list_manager.read_changelog(snapshot)
return self._create_plan_from_manifests(manifest_files, snapshot.id)
diff --git a/paimon-python/pypaimon/tests/native_plan_capabilities_test.py
b/paimon-python/pypaimon/tests/native_plan_capabilities_test.py
index 867b3261cd..02a3264670 100644
--- a/paimon-python/pypaimon/tests/native_plan_capabilities_test.py
+++ b/paimon-python/pypaimon/tests/native_plan_capabilities_test.py
@@ -20,6 +20,7 @@
import json
import tempfile
import unittest
+from contextlib import ExitStack
from dataclasses import replace
from unittest.mock import patch
@@ -116,7 +117,9 @@ class NativePlanCapabilitiesTest(unittest.TestCase):
plans = []
for native in (False, True):
read_table = table.copy({
- 'scan.native-plan.enabled': str(native).lower()})
+ 'scan.native-plan.enabled': str(native).lower(),
+ 'read.native.enabled': str(native).lower(),
+ })
builder = read_table.new_read_builder()
if predicate is not None:
builder.with_filter(predicate)
@@ -133,7 +136,18 @@ class NativePlanCapabilitiesTest(unittest.TestCase):
plan = scan.plan()
else:
plan = scan.plan()
- rows = builder.new_read().to_arrow(plan.splits()).to_pylist()
+ if native:
+ self.assertTrue(all(
+ getattr(split, '_native_split', None) is not None
+ for split in plan.splits()))
+ read_guard = patch(
+ 'pypaimon.read.table_read.TableRead._create_split_read',
+ side_effect=AssertionError(
+ 'native capability read fell back to Python'))
+ else:
+ read_guard = ExitStack()
+ with read_guard:
+ rows = builder.new_read().to_arrow(plan.splits()).to_pylist()
self.assertEqual(sorted(rows, key=lambda row: row['k']),
expected_rows)
self.assertEqual(plan.snapshot_id, snapshot_id)
plans.append(plan)
diff --git a/paimon-python/pypaimon/tests/native_plan_dynamic_bucket_test.py
b/paimon-python/pypaimon/tests/native_plan_dynamic_bucket_test.py
index 27ca00ef6f..f4d58b6422 100644
--- a/paimon-python/pypaimon/tests/native_plan_dynamic_bucket_test.py
+++ b/paimon-python/pypaimon/tests/native_plan_dynamic_bucket_test.py
@@ -18,6 +18,7 @@
"""Real bucket growth and cross-partition updates through native planning."""
import json
+from contextlib import ExitStack
from unittest.mock import patch
import pyarrow as pa
@@ -47,7 +48,10 @@ def
test_native_pk_bucket_growth_and_partition_migration(tmp_path, cross_partiti
}), False)
def read(table, native, predicate=None, shard=None, limit=None):
- builder = table.copy({'scan.native-plan.enabled':
str(native).lower()}).new_read_builder()
+ builder = table.copy({
+ 'scan.native-plan.enabled': str(native).lower(),
+ 'read.native.enabled': str(native).lower(),
+ }).new_read_builder()
if predicate is not None:
builder.with_filter(predicate)
if limit is not None:
@@ -60,7 +64,16 @@ def
test_native_pk_bucket_growth_and_partition_migration(tmp_path, cross_partiti
plan = scan.plan()
else:
plan = scan.plan()
- rows = builder.new_read().to_arrow(plan.splits()).to_pylist()
+ if native:
+ assert all(getattr(split, '_native_split', None) is not None
+ for split in plan.splits())
+ read_guard = patch(
+ 'pypaimon.read.table_read.TableRead._create_split_read',
+ side_effect=AssertionError('dynamic-bucket native read fell
back'))
+ else:
+ read_guard = ExitStack()
+ with read_guard:
+ rows = builder.new_read().to_arrow(plan.splits()).to_pylist()
return plan, sorted(rows, key=lambda row: (row['id'], row['p']))
initial = [{'id': 1, 'p': 'a', 'v': 'old'}, {'id': 2, 'p': 'a', 'v':
'two'},
@@ -116,6 +129,7 @@ def
test_native_pk_bucket_growth_and_partition_migration(tmp_path, cross_partiti
for partition in (None, 'a', 'b'):
builder = table.copy({
'scan.native-plan.enabled': str(native).lower(),
+ 'read.native.enabled': str(native).lower(),
'incremental-between-timestamp': '0,200',
}).new_read_builder()
if partition is not None:
@@ -127,8 +141,27 @@ def
test_native_pk_bucket_growth_and_partition_migration(tmp_path, cross_partiti
else:
plan = scan.plan()
assert all(split.is_streaming for split in plan.splits())
- actual = [(row.get_field(0), row.get_field(1), row.get_field(2),
- row.get_row_kind().value)
- for row in builder.new_read().to_iterator(plan.splits())]
+ if native:
+ from pypaimon.table.row.row_kind import RowKind
+ assert all(getattr(split, '_native_split', None) is not None
+ for split in plan.splits())
+ read = builder.new_read()
+ read.include_row_kind = True
+ with patch(
+
'pypaimon.read.table_read.TableRead._create_split_read',
+ side_effect=AssertionError(
+ 'dynamic incremental native read fell back')):
+ rows = read.to_arrow(plan.splits()).to_pylist()
+ actual = [
+ (row['id'], row['p'], row['v'],
+ RowKind.from_string(row['_row_kind']).value)
+ for row in rows
+ ]
+ else:
+ actual = [
+ (row.get_field(0), row.get_field(1), row.get_field(2),
+ row.get_row_kind().value)
+ for row in builder.new_read().to_iterator(plan.splits())
+ ]
assert sorted(actual) == sorted(
event for event in events if partition is None or event[1] ==
partition)
diff --git a/paimon-python/pypaimon/tests/native_plan_expanded_test.py
b/paimon-python/pypaimon/tests/native_plan_expanded_test.py
index 06422c1c6d..244ab1cfa1 100644
--- a/paimon-python/pypaimon/tests/native_plan_expanded_test.py
+++ b/paimon-python/pypaimon/tests/native_plan_expanded_test.py
@@ -18,6 +18,7 @@
"""Native planning coverage for Java time travel, postpone and scored DE
reads."""
import json
+from contextlib import ExitStack
from unittest.mock import patch
import pyarrow as pa
@@ -63,7 +64,10 @@ def write(table, rows, schema):
def read(table, native, predicate=None, shard=None, limit=None, result=None,
ranges=None):
- builder = table.copy({'scan.native-plan.enabled':
str(native).lower()}).new_read_builder()
+ builder = table.copy({
+ 'scan.native-plan.enabled': str(native).lower(),
+ 'read.native.enabled': str(native).lower(),
+ }).new_read_builder()
if predicate is not None:
builder.with_filter(predicate)
if limit is not None:
@@ -80,7 +84,15 @@ def read(table, native, predicate=None, shard=None,
limit=None, result=None, ran
plan = scan.plan()
else:
plan = scan.plan()
- return plan, builder.new_read().to_arrow(plan.splits()).to_pylist()
+ if native:
+ read_guard = patch(
+ 'pypaimon.read.table_read.TableRead._create_split_read',
+ side_effect=AssertionError('expanded native read fell back'))
+ else:
+ read_guard = ExitStack()
+ with read_guard:
+ rows = builder.new_read().to_arrow(plan.splits()).to_pylist()
+ return plan, rows
@pytest.mark.parametrize('version,expected', [('1', 1), ('base', 1),
('watermark-150', 2)])
diff --git a/paimon-python/pypaimon/tests/native_plan_incremental_test.py
b/paimon-python/pypaimon/tests/native_plan_incremental_test.py
index f77a4fc6bd..791041b753 100644
--- a/paimon-python/pypaimon/tests/native_plan_incremental_test.py
+++ b/paimon-python/pypaimon/tests/native_plan_incremental_test.py
@@ -430,6 +430,108 @@ def
test_stream_read_builder_combines_native_nested_projection_and_row_kind(
]
[email protected]_plan
+def test_streaming_changelog_frames_use_native_plan_and_read(catalog):
+ import asyncio
+
+ table = _table(catalog, 'native_changelog', True, {
+ 'bucket': '1',
+ 'changelog-producer': 'input',
+ 'source.split.target-size': '1 b',
+ 'source.split.open-file-cost': '1 b',
+ })
+ _write(table, 100, [{'k': 1, 'v': 'a'}, {'k': 2, 'v': 'b'}])
+ _write(table, 200, [{'k': 3, 'v': 'c'}])
+
+ native_table = table.copy({
+ 'scan.native-plan.enabled': 'true',
+ 'read.native.enabled': 'true',
+ })
+ predicate = (native_table.new_read_builder().new_predicate_builder()
+ .greater_than('k', 1))
+ builder = (native_table.new_stream_read_builder()
+ .with_filter(predicate)
+ .with_projection(['v'])
+ .with_include_row_kind())
+ scan = builder.new_streaming_scan()
+ scan.next_snapshot_id = 1
+
+ async def first_two_frames():
+ plans = []
+ async for plan in scan.stream():
+ plans.append(plan)
+ if len(plans) == 2:
+ return plans
+
+ with patch.object(
+ scan, '_create_plan_from_manifests',
+ side_effect=AssertionError(
+ 'streaming changelog native plan fell back to Python')):
+ plans = asyncio.run(first_two_frames())
+
+ assert [plan.snapshot_id for plan in plans] == [1, 2]
+ assert all(
+ getattr(split, '_native_split', None) is not None
+ for plan in plans for split in plan.splits())
+ assert all(
+ file.file_name.startswith('changelog-')
+ for plan in plans for split in plan.splits() for file in split.files)
+
+ with patch(
+ 'pypaimon.read.table_read.TableRead._create_split_read',
+ side_effect=AssertionError(
+ 'streaming changelog native read fell back to Python')):
+ rows = [
+ row
+ for plan in plans
+ for row in builder.new_read().to_arrow(
+ plan.splits(), parallelism=2).to_pylist()
+ ]
+ assert rows == [
+ {'_row_kind': '+I', 'v': 'b'},
+ {'_row_kind': '+I', 'v': 'c'},
+ ]
+
+
[email protected]_plan
+def test_streaming_overwrite_changelog_uses_java_follow_up_semantics(catalog):
+ table = _table(catalog, 'native_overwrite_changelog', True, {
+ 'bucket': '1',
+ 'changelog-producer': 'input',
+ })
+ _write(table, 100, [{'k': 1, 'v': 'before'}])
+ _write(table, 200, [{'k': 2, 'v': 'after'}], overwrite=True)
+ snapshot = table.snapshot_manager().get_latest_snapshot()
+ assert snapshot.commit_kind == 'OVERWRITE'
+ assert snapshot.changelog_manifest_list is not None
+
+ native_table = table.copy({
+ 'scan.native-plan.enabled': 'true',
+ 'read.native.enabled': 'true',
+ })
+ builder = (native_table.new_stream_read_builder()
+ .with_projection(['v'])
+ .with_include_row_kind())
+ scan = builder.new_streaming_scan()
+
+ with patch.object(
+ scan, '_try_native_plan',
+ side_effect=AssertionError(
+ 'OVERWRITE must not use range-based native changelog
planning')):
+ plan = scan._create_changelog_plan(snapshot)
+
+ assert plan.snapshot_id == snapshot.id
+ assert all(
+ file.file_name.startswith('changelog-')
+ for split in plan.splits() for file in split.files)
+ with patch(
+ 'pypaimon.read.table_read.TableRead._create_split_read',
+ side_effect=AssertionError(
+ 'OVERWRITE changelog native read fell back to Python')):
+ rows = builder.new_read().to_arrow(plan.splits()).to_pylist()
+ assert rows == [{'_row_kind': '+I', 'v': 'after'}]
+
+
def test_streaming_reader_honors_explicit_split_deletion_vector(catalog,
native, tmp_path):
from pypaimon.deletionvectors.bitmap_deletion_vector import
BitmapDeletionVector
from pypaimon.table.source.deletion_file import DeletionFile
diff --git a/paimon-python/pypaimon/tests/native_plan_materialized_pk_test.py
b/paimon-python/pypaimon/tests/native_plan_materialized_pk_test.py
index 6e9471088f..c418e71173 100644
--- a/paimon-python/pypaimon/tests/native_plan_materialized_pk_test.py
+++ b/paimon-python/pypaimon/tests/native_plan_materialized_pk_test.py
@@ -15,6 +15,7 @@
# specific language governing permissions and limitations
# under the License.
+from contextlib import ExitStack
from dataclasses import replace
from pathlib import Path
from unittest.mock import patch
@@ -88,7 +89,10 @@ def
test_clustered_materialized_dv_files_use_native_raw_splits(tmp_path, engine,
pb = table.new_read_builder().new_predicate_builder()
for predicate in (None, pb.equal('id', 3), pb.equal('value', 'v8')):
for native in (False, True):
- builder = table.copy({'scan.native-plan.enabled':
str(native).lower()}).new_read_builder()
+ builder = table.copy({
+ 'scan.native-plan.enabled': str(native).lower(),
+ 'read.native.enabled': str(native).lower(),
+ }).new_read_builder()
if predicate is not None:
builder.with_filter(predicate)
scan = builder.new_scan()
@@ -98,7 +102,17 @@ def
test_clustered_materialized_dv_files_use_native_raw_splits(tmp_path, engine,
else:
plan = scan.plan()
assert all(split.raw_convertible for split in plan.splits())
- rows = builder.new_read().to_arrow(plan.splits()).to_pylist()
+ if native:
+ assert all(getattr(split, '_native_split', None) is not None
+ for split in plan.splits())
+ read_guard = patch(
+ 'pypaimon.read.table_read.TableRead._create_split_read',
+ side_effect=AssertionError(
+ 'materialized PK native read fell back'))
+ else:
+ read_guard = ExitStack()
+ with read_guard:
+ rows = builder.new_read().to_arrow(plan.splits()).to_pylist()
wanted = expected if predicate is None else ([] if predicate.field
== 'id' else [expected[1]])
assert sorted(rows, key=lambda row: row['id']) == wanted
assert plan.snapshot_id == 3
@@ -172,7 +186,10 @@ def
test_first_row_level_zero_merges_before_filtering(tmp_path, target_size, com
for predicate, expected in [(None, all_ids), (pb.equal('value', 'later'),
[]),
(pb.equal('value', 'first'), [1]),
(pb.equal('id', 2), [])]:
for native in (False, True):
- builder = table.copy({'scan.native-plan.enabled':
str(native).lower()}).new_read_builder()
+ builder = table.copy({
+ 'scan.native-plan.enabled': str(native).lower(),
+ 'read.native.enabled': str(native).lower(),
+ }).new_read_builder()
builder.with_projection(['id'])
if predicate is not None:
builder.with_filter(predicate)
@@ -183,4 +200,16 @@ def
test_first_row_level_zero_merges_before_filtering(tmp_path, target_size, com
else:
plan = scan.plan()
assert plan.snapshot_id == len(batches) + 1
- assert
sorted(builder.new_read().to_arrow(plan.splits()).column('id').to_pylist()) ==
expected
+ if native:
+ assert all(getattr(split, '_native_split', None) is not None
+ for split in plan.splits())
+ read_guard = patch(
+ 'pypaimon.read.table_read.TableRead._create_split_read',
+ side_effect=AssertionError(
+ 'first-row native read fell back'))
+ else:
+ read_guard = ExitStack()
+ with read_guard:
+ actual = builder.new_read().to_arrow(
+ plan.splits()).column('id').to_pylist()
+ assert sorted(actual) == expected
diff --git a/paimon-python/pypaimon/tests/native_plan_resolved_schema_test.py
b/paimon-python/pypaimon/tests/native_plan_resolved_schema_test.py
index 2d837694c8..cbae972b09 100644
--- a/paimon-python/pypaimon/tests/native_plan_resolved_schema_test.py
+++ b/paimon-python/pypaimon/tests/native_plan_resolved_schema_test.py
@@ -73,7 +73,10 @@ def _write(table, rows):
def _read(table, native, predicate=None, projection=None):
- table = table.copy_without_time_travel({'scan.native-plan.enabled':
str(native).lower()})
+ table = table.copy_without_time_travel({
+ 'scan.native-plan.enabled': str(native).lower(),
+ 'read.native.enabled': str(native).lower(),
+ })
builder = table.new_read_builder()
if predicate is not None:
builder.with_filter(predicate)
@@ -97,7 +100,16 @@ def _read(table, native, predicate=None, projection=None):
plan = scan.plan()
else:
plan = scan.plan()
- rows = builder.new_read().to_arrow(plan.splits()).to_pylist()
+ if native:
+ assert all(getattr(split, '_native_split', None) is not None
+ for split in plan.splits())
+ read_guard = patch(
+ 'pypaimon.read.table_read.TableRead._create_split_read',
+ side_effect=AssertionError('resolved-schema native read fell
back'))
+ else:
+ read_guard = ExitStack()
+ with read_guard:
+ rows = builder.new_read().to_arrow(plan.splits()).to_pylist()
return plan.snapshot_id, sorted(rows, key=lambda row: row['id'])
diff --git a/paimon-python/pypaimon/tests/native_plan_test.py
b/paimon-python/pypaimon/tests/native_plan_test.py
index 717e69dec6..0d9f2f581d 100644
--- a/paimon-python/pypaimon/tests/native_plan_test.py
+++ b/paimon-python/pypaimon/tests/native_plan_test.py
@@ -993,6 +993,28 @@ class NativePlanTest(unittest.TestCase):
self.assertEqual(native.call_args[1]['incremental_range'], (2, 4))
fs.scan.assert_not_called()
+ def test_native_plan_forwards_explicit_incremental_mode(self):
+ table = _scan(True, Mock()).table
+ table.table_schema = Mock(fields=[], partition_keys=[])
+ rust_plan = SimpleNamespace(splits=lambda: [], snapshot_id=lambda: 4)
+ rust_scan = Mock()
+ rust_scan.plan.return_value = rust_plan
+ builder = Mock()
+ builder.new_incremental_scan.return_value = rust_scan
+ with patch('pypaimon.read.native_plan.native_runtime_available',
+ return_value=True), patch(
+ 'pypaimon.read.native_plan._native_read_builder',
+ return_value=builder):
+ plan = native_plan(
+ table, incremental_range=(2, 4),
+ incremental_mode='changelog')
+ self.assertEqual(plan.snapshot_id, 4)
+ builder.new_incremental_scan.assert_called_once_with(
+ 2, 4, 'changelog')
+
+ with self.assertRaisesRegex(ValueError, 'incremental_range'):
+ native_plan(table, incremental_mode='changelog')
+
def test_incremental_window_outside_snapshots_is_terminal_empty(self):
fs = Mock(partition_key_predicate=None)
scan = _scan(True, fs)
diff --git a/paimon-python/pypaimon/tests/streaming_table_scan_test.py
b/paimon-python/pypaimon/tests/streaming_table_scan_test.py
index 630f1324af..04ecb48ad3 100644
--- a/paimon-python/pypaimon/tests/streaming_table_scan_test.py
+++ b/paimon-python/pypaimon/tests/streaming_table_scan_test.py
@@ -114,6 +114,53 @@ class AsyncStreamingTableScanTest(unittest.TestCase):
self.assertEqual(
native_plan.call_args.kwargs['projection'], ['payload', 'id'])
+ @patch('pypaimon.read.streaming_table_scan.ManifestListManager')
+ @patch('pypaimon.read.streaming_table_scan.ManifestFileManager')
+ @patch('pypaimon.read.native_plan.native_plan')
+ def test_changelog_scan_uses_explicit_native_mode(
+ self, native_plan, _manifest_files, _manifest_lists):
+ table, _ = _create_mock_table()
+ table.options.native_plan_enabled.return_value = True
+ table.options.changelog_producer.return_value = ChangelogProducer.INPUT
+ native_plan.return_value = Plan([], snapshot_id=5)
+ scan = AsyncStreamingTableScan(table, predicate=Mock())
+ scan._read_type = [Mock(name='payload')]
+ scan._read_type[0].name = 'payload'
+
+ plan = scan._create_changelog_plan(_create_mock_snapshot(5))
+
+ self.assertIs(plan, native_plan.return_value)
+ self.assertEqual(
+ native_plan.call_args.kwargs['incremental_range'], (4, 5))
+ self.assertEqual(
+ native_plan.call_args.kwargs['incremental_mode'], 'changelog')
+ self.assertEqual(
+ native_plan.call_args.kwargs['projection'], ['payload'])
+
+ @patch('pypaimon.read.streaming_table_scan.ManifestListManager')
+ @patch('pypaimon.read.streaming_table_scan.ManifestFileManager')
+ @patch('pypaimon.read.native_plan.native_plan')
+ def test_overwrite_changelog_keeps_per_snapshot_python_plan(
+ self, native_plan, _manifest_files, manifest_lists):
+ table, _ = _create_mock_table()
+ table.options.native_plan_enabled.return_value = True
+ table.options.changelog_producer.return_value = ChangelogProducer.INPUT
+ snapshot = _create_mock_snapshot(5, 'OVERWRITE')
+ manifests = [Mock()]
+ manifest_lists.return_value.read_changelog.return_value = manifests
+ python_plan = Plan([], snapshot_id=5)
+ scan = AsyncStreamingTableScan(table)
+
+ with patch.object(
+ scan, '_create_plan_from_manifests',
+ return_value=python_plan) as create_plan:
+ plan = scan._create_changelog_plan(snapshot)
+
+ self.assertIs(plan, python_plan)
+ native_plan.assert_not_called()
+
manifest_lists.return_value.read_changelog.assert_called_once_with(snapshot)
+ create_plan.assert_called_once_with(manifests, 5)
+
@patch('pypaimon.read.streaming_table_scan.ManifestListManager')
@patch('pypaimon.read.streaming_table_scan.ManifestFileManager')
@patch('pypaimon.read.native_plan.native_plan')