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')

Reply via email to