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 e39c937cff [Python] Keep append distribution native end to end (#10032)
e39c937cff is described below

commit e39c937cffb769097960b693376bc2d10055ebef
Author: Jingsong Lee <[email protected]>
AuthorDate: Mon Sep 21 11:07:40 2026 +0800

    [Python] Keep append distribution native end to end (#10032)
---
 paimon-python/README.md                            |  8 ++--
 paimon-python/pypaimon/read/table_scan.py          | 23 +++++----
 .../tests/native_plan_distribution_test.py         | 54 +++++++++++++++++++++-
 .../pypaimon/tests/native_plan_incremental_test.py | 13 +++++-
 paimon-python/pypaimon/tests/native_plan_test.py   |  4 +-
 5 files changed, 84 insertions(+), 18 deletions(-)

diff --git a/paimon-python/README.md b/paimon-python/README.md
index 6aaffefb64..4fa2d7246f 100644
--- a/paimon-python/README.md
+++ b/paimon-python/README.md
@@ -181,9 +181,11 @@ When using an unreleased 0.4.0 development wheel, rebuild 
it with these fixes;
 package version checks cannot distinguish local builds with identical versions.
 
 Append scans support `with_shard()` and `with_slice()` with Rust 0.4 or newer,
-which preserves the file order needed for positional selection; primary-key 
scans support
-bucket-based `with_shard()`. Data-evolution position selection requires the
-binding's `TableScan.with_row_position_slice()` and 
`with_row_position_shard()`.
+which plans positional selection directly into native-readable splits; 
primary-key
+scans support bucket-based `with_shard()`. Both append and data-evolution 
position
+selection require the binding's `TableScan.with_row_position_slice()` and
+`with_row_position_shard()`. Ordinary append positions follow the stats-pruned
+split/file order, while Data Evolution assigns positions before group pruning.
 Selection occurs before reader filtering and deletion vectors, so surviving row
 counts can differ between shards. Limits are applied after shard/slice 
selection.
 
diff --git a/paimon-python/pypaimon/read/table_scan.py 
b/paimon-python/pypaimon/read/table_scan.py
index bcdd203ffc..71d3577685 100755
--- a/paimon-python/pypaimon/read/table_scan.py
+++ b/paimon-python/pypaimon/read/table_scan.py
@@ -156,7 +156,9 @@ class TableScan:
             # 0.4.0 includes Python-written DV decoding and legacy bucket 
paths.
             if not native_version_at_least(0, 4, 0):
                 return False
-        if getattr(fs, 'data_evolution', False):
+        if (not self.table.is_primary_key_table
+                and (getattr(fs, 'idx_of_this_subtask', None) is not None
+                     or getattr(fs, 'start_pos_of_this_subtask', None) is not 
None)):
             if (getattr(fs, 'idx_of_this_subtask', None) is not None
                     and not native_method_available('TableScan', 
'with_row_position_shard')):
                 return False
@@ -279,7 +281,8 @@ class TableScan:
                 if fs.idx_of_this_subtask is not None:
                     extra_options['shard'] = (
                         fs.idx_of_this_subtask, fs.number_of_para_subtasks)
-            if has_distribution and fs.data_evolution and chunk_shuffle is 
None:
+            if (has_distribution and not self.table.is_primary_key_table
+                    and chunk_shuffle is None):
                 if fs.idx_of_this_subtask is not None:
                     extra_options['row_position_shard'] = (
                         fs.idx_of_this_subtask, fs.number_of_para_subtasks)
@@ -329,16 +332,12 @@ class TableScan:
                 if self.table.is_primary_key_table:
                     splits = [s for s in splits
                               if s.bucket % fs.number_of_para_subtasks == 
fs.idx_of_this_subtask]
-                elif not fs.data_evolution:
-                    from pypaimon.read.scan_distribution import shard_range, 
slice_append_splits
-                    if fs.idx_of_this_subtask is not None:
-                        start, end = shard_range(
-                            sum(s.row_count for s in splits),
-                            fs.idx_of_this_subtask, fs.number_of_para_subtasks)
-                    else:
-                        start, end = fs.start_pos_of_this_subtask, 
fs.end_pos_of_this_subtask
-                    splits = slice_append_splits(splits, start, end)
-                splits = fs._apply_push_down_limit(splits)
+                # A partial IndexedSplit plus a file-wide DV cardinality cannot
+                # reveal how many deleted rows lie inside the selected range.
+                # Keep every selected split and let the reader enforce LIMIT.
+                if not (getattr(fs, 'deletion_vectors_enabled', False)
+                        and not self.table.is_primary_key_table):
+                    splits = fs._apply_push_down_limit(splits)
             # Attach scores to the row ranges retained by native planning.
             from pypaimon.globalindex.indexed_split import IndexedSplit, 
scores_for_ranges
             from pypaimon.globalindex.vector_search_result import 
ScoredGlobalIndexResult
diff --git a/paimon-python/pypaimon/tests/native_plan_distribution_test.py 
b/paimon-python/pypaimon/tests/native_plan_distribution_test.py
index 00d3dfcd5b..cfb0b54dce 100644
--- a/paimon-python/pypaimon/tests/native_plan_distribution_test.py
+++ b/paimon-python/pypaimon/tests/native_plan_distribution_test.py
@@ -90,6 +90,7 @@ class _DistributionFixture:
               limit=None, predicate=None, projection=None):
         builder = table.copy({
             'scan.native-plan.enabled': str(native).lower(),
+            'read.native.enabled': str(native).lower(),
         }).new_read_builder()
         if limit is not None:
             builder.with_limit(limit)
@@ -108,7 +109,20 @@ class _DistributionFixture:
             'distributed native plan fell back to Python')) if native else 
ExitStack()
         with guard:
             plan = scan.plan()
-        return plan, builder.new_read().to_arrow(plan.splits(), 
parallelism=1).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(
+                    'distributed native read fell back to Python'))
+        else:
+            read_guard = ExitStack()
+        with read_guard:
+            rows = builder.new_read().to_arrow(
+                plan.splits(), parallelism=1).to_pylist()
+        return plan, rows
 
     def _assert_parity(self, table, expected, snapshot_id, ordered=True, 
**options):
         plans = []
@@ -182,6 +196,44 @@ class NativePlanDistributionTest(_DistributionFixture, 
unittest.TestCase):
             with self.subTest(start=start, end=end):
                 self._assert_parity(table, rows[start:end], 3, slice_=(start, 
end))
 
+    def test_row_tracked_append_uses_global_ranges_for_native_read(self):
+        table = self._create('append_row_tracking', {
+            'row-tracking.enabled': 'true',
+            'source.split.target-size': '1b',
+        })
+        rows = [{'k': k, 'v': str(k)} for k in range(9)]
+        for start in range(0, 9, 3):
+            self._write(table, rows[start:start + 3])
+
+        native = self._assert_parity(table, rows[2:7], 3, slice_=(2, 7))
+        selected_ranges = [
+            (range_.from_, range_.to)
+            for split in native.splits()
+            for range_ in getattr(split, 'row_ranges', lambda: [])()
+        ]
+        # The middle file is selected in full and stays a plain DataSplit;
+        # only partial boundary files need explicit global ranges.
+        self.assertEqual(selected_ranges, [(2, 2), (6, 6)])
+        for index, expected in enumerate((rows[:3], rows[3:5], rows[5:7], 
rows[7:])):
+            self._assert_parity(table, expected, 3, shard=(index, 4))
+
+    def test_append_positions_follow_partition_pruned_file_order(self):
+        schema = self.schema.append(pa.field('p', pa.string()))
+        table = self._create(
+            'append_partition_filter', {'source.split.target-size': '1b'},
+            schema, ['p'])
+        low = [{'k': k, 'v': str(k), 'p': 'low'} for k in range(3)]
+        high = [{'k': k, 'v': str(k), 'p': 'high'} for k in range(100, 103)]
+        self._write(table, low, schema)
+        self._write(table, high, schema)
+        predicate = (table.new_read_builder().new_predicate_builder()
+                     .equal('p', 'high'))
+
+        self._assert_parity(
+            table, high[1:], 2, slice_=(1, 3), predicate=predicate)
+        self._assert_parity(
+            table, high[1:2], 2, slice_=(1, 3), predicate=predicate, limit=1)
+
     def test_append_limit_is_applied_after_shard_or_slice(self):
         table = self._create('append_limit', {'source.split.target-size': 
'1b'})
         rows = [{'k': k, 'v': str(k)} for k in range(12)]
diff --git a/paimon-python/pypaimon/tests/native_plan_incremental_test.py 
b/paimon-python/pypaimon/tests/native_plan_incremental_test.py
index d633338f9e..f77a4fc6bd 100644
--- a/paimon-python/pypaimon/tests/native_plan_incremental_test.py
+++ b/paimon-python/pypaimon/tests/native_plan_incremental_test.py
@@ -88,6 +88,7 @@ def _read(table, native, window, predicate=None, limit=None, 
shard=None,
           slice_=None, row_ranges=None, with_stats=False):
     builder = table.copy({
         'scan.native-plan.enabled': str(native).lower(),
+        'read.native.enabled': str(native).lower(),
         'scan.mode': 'incremental',
         'incremental-between-timestamp': '%s,%s' % window,
     }).new_read_builder()
@@ -109,7 +110,17 @@ def _read(table, native, window, predicate=None, 
limit=None, shard=None,
                     scan.file_scanner, method, side_effect=AssertionError(
                         'incremental native plan fell back to Python')))
         plan = scan.scan_with_stats()[0] if with_stats else scan.plan()
-    result = builder.new_read().to_arrow(plan.splits(), 
parallelism=1).to_pydict()
+    with ExitStack() as stack:
+        if native:
+            assert all(
+                getattr(split, '_native_split', None) is not None
+                for split in plan.splits())
+            stack.enter_context(patch(
+                'pypaimon.read.table_read.TableRead._create_split_read',
+                side_effect=AssertionError(
+                    'incremental native read fell back to Python')))
+        result = builder.new_read().to_arrow(
+            plan.splits(), parallelism=1).to_pydict()
     return plan, [dict(zip(result, row)) for row in zip(*result.values())]
 
 
diff --git a/paimon-python/pypaimon/tests/native_plan_test.py 
b/paimon-python/pypaimon/tests/native_plan_test.py
index 8ea40f6b2b..717e69dec6 100644
--- a/paimon-python/pypaimon/tests/native_plan_test.py
+++ b/paimon-python/pypaimon/tests/native_plan_test.py
@@ -903,7 +903,9 @@ class NativePlanTest(unittest.TestCase):
                     setattr(fs, selection, 0)
                     fs.scan.return_value = fallback = object()
                     with 
patch('pypaimon.read.native_plan.native_version_at_least',
-                               return_value=available):
+                               return_value=available), patch(
+                            
'pypaimon.read.native_plan.native_method_available',
+                            return_value=True):
                         self.assertEqual(scan._native_plan_supported(), 
available)
                         if not available:
                             with 
patch('pypaimon.read.native_plan.native_plan') as native:

Reply via email to