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 13af729902 [python] Plan data evolution row ranges with pypaimon-rust 
(#9038)
13af729902 is described below

commit 13af7299024d4a3784b196012991752bd0e36f00
Author: XiaoHongbo <[email protected]>
AuthorDate: Thu Aug 6 12:03:26 2026 +0800

    [python] Plan data evolution row ranges with pypaimon-rust (#9038)
---
 .github/workflows/paimon-python-checks.yml         |  2 +-
 paimon-python/pypaimon/read/native_plan.py         |  7 +-
 paimon-python/pypaimon/read/table_scan.py          | 37 ++++++--
 .../pypaimon/tests/native_plan_integration_test.py | 51 +++++++++++
 paimon-python/pypaimon/tests/native_plan_test.py   | 98 ++++++++++++++++++++--
 5 files changed, 176 insertions(+), 19 deletions(-)

diff --git a/.github/workflows/paimon-python-checks.yml 
b/.github/workflows/paimon-python-checks.yml
index bad89021ac..836a1aa665 100755
--- a/.github/workflows/paimon-python-checks.yml
+++ b/.github/workflows/paimon-python-checks.yml
@@ -34,7 +34,7 @@ env:
   JDK_VERSION: 8
   MAVEN_OPTS: -Dmaven.wagon.httpconnectionManager.ttlSeconds=30 
-Dmaven.wagon.http.retryHandler.requestSentEnabled=true
   LUMINA_DATA_VERSION: 0.1.0
-  PYPAIMON_RUST_REV: 4df2bcc2e8d245aafba98c8014e710098ef7ac6b
+  PYPAIMON_RUST_REV: b8d590521eb4cfde916ce3589cda88e75f31046a
 
 
 concurrency:
diff --git a/paimon-python/pypaimon/read/native_plan.py 
b/paimon-python/pypaimon/read/native_plan.py
index 54c87e3594..55b91d93b8 100644
--- a/paimon-python/pypaimon/read/native_plan.py
+++ b/paimon-python/pypaimon/read/native_plan.py
@@ -22,7 +22,7 @@ Predicates and limits are pushed into Rust planning. The 
normal pypaimon reader
 still applies them while reading, so pushdown remains an optimization.
 """
 
-from typing import List, Optional
+from typing import List, Optional, Tuple
 
 from pypaimon.common.options.config import CatalogOptions
 from pypaimon.common.options.core_options import CoreOptions
@@ -171,7 +171,8 @@ def native_plan(
         table,
         predicate: Optional[Predicate] = None,
         limit: Optional[int] = None,
-        projection: Optional[List[str]] = None) -> List[Split]:
+        projection: Optional[List[str]] = None,
+        row_ranges: Optional[List[Tuple[int, int]]] = None) -> List[Split]:
     """Plan with pypaimon_rust and return the decoded pypaimon splits.
 
     Native conversion or planning failures are handled by TableScan, which
@@ -190,6 +191,8 @@ def native_plan(
         builder = builder.with_filter(_predicate_to_native(predicate))
     if limit is not None:
         builder = builder.with_limit(limit)
+    if row_ranges is not None:
+        builder = builder.with_row_ranges(row_ranges)
     rust_splits = builder.new_scan().plan().splits()
     pfields = _partition_fields(table)
     # Trimmed primary keys decode per-file min/max keys (PK merge-on-read).
diff --git a/paimon-python/pypaimon/read/table_scan.py 
b/paimon-python/pypaimon/read/table_scan.py
index af1a939f2f..fc50ece0f8 100755
--- a/paimon-python/pypaimon/read/table_scan.py
+++ b/paimon-python/pypaimon/read/table_scan.py
@@ -17,7 +17,7 @@
 
 import json as _json
 import logging
-from typing import Optional, Tuple
+from typing import List, Optional, Tuple
 
 from pypaimon.catalog.catalog_exception import TableNoPermissionException
 from pypaimon.common.identifier import UNKNOWN_DATABASE
@@ -95,9 +95,9 @@ class TableScan:
 
     def _native_plan_supported_impl(self) -> bool:
         """Fall back to the Python scanner for scans native can't carry:
-        shard/slice, chunk-shuffle, global-index/row-ranges, first-row
-        merge-engine (Rust drops L0), deletion vectors, postpone bucket
-        (drops synthetic buckets),
+        shard/slice, chunk-shuffle, explicit row ranges, scored or primary-key
+        global-index results, first-row merge-engine (Rust drops L0), deletion
+        vectors, postpone bucket,
         a primary-key table whose trimmed PK is empty (PK equals the partition
         key; native may mark splits raw-convertible and skip merge), dynamic
         bucket / cross-partition PK tables (unconfirmed Rust parity), a stale
@@ -113,8 +113,8 @@ class TableScan:
         if (getattr(fs, 'idx_of_this_subtask', None) is not None
                 or getattr(fs, 'start_pos_of_this_subtask', None) is not None
                 or getattr(fs, 'chunk_shuffle', None) is not None
-                or getattr(fs, '_global_index_result', None) is not None
                 or getattr(fs, '_row_ranges', None) is not None
+                or not self._native_global_index_result_supported()
                 or getattr(fs, 'deletion_vectors_enabled', False)
                 or getattr(fs, 'only_read_real_buckets', False)):
             return False
@@ -173,16 +173,36 @@ class TableScan:
             return False
         return not options.contains(CoreOptions.INCREMENTAL_BETWEEN_TIMESTAMP)
 
+    def _native_global_index_result_supported(self) -> bool:
+        result = self.file_scanner._global_index_result
+        if result is None:
+            return True
+        if (self.table.is_primary_key_table
+                or not self.file_scanner.data_evolution):
+            return False
+        from pypaimon.globalindex.global_index_result import GlobalIndexResult
+        from pypaimon.globalindex.vector_search_result import 
ScoredGlobalIndexResult
+        return (isinstance(result, GlobalIndexResult)
+                and not isinstance(result, ScoredGlobalIndexResult))
+
+    def _native_global_index_row_ranges(self) -> Optional[List[Tuple[int, 
int]]]:
+        result = self.file_scanner._global_index_result
+        if result is None:
+            return None
+        return [(range_.from_, range_.to)
+                for range_ in result.results().to_range_list()]
+
     def _try_native_plan(self) -> Optional[Plan]:
         """Plan via pypaimon_rust, then drop partitions the predicate rejects.
 
         Predicate and limit are pushed into Rust planning and are still 
enforced
-        by the reader. Return None when Rust finds no splits so the caller can 
use
-        the matching Python fallback (with scan stats when requested).
+        by the reader. Empty unrestricted scans fall back to preserve snapshot
+        metadata; explicit empty row ranges are a terminal empty result.
         """
         from pypaimon.read.native_plan import native_plan
 
         try:
+            row_ranges = self._native_global_index_row_ranges()
             native_predicate = self.predicate
             if self.partition_predicate is not None:
                 native_predicate = PredicateBuilder.and_predicates([
@@ -198,9 +218,10 @@ class TableScan:
                 projection=(
                     [field.name for field in self._read_type]
                     if self._read_type is not None else None),
+                row_ranges=row_ranges,
             )
             if not splits:
-                return None
+                return Plan([]) if row_ranges is not None else None
             snapshot_id = splits[0].snapshot_id
             partition_predicate = self.file_scanner.partition_key_predicate
             if partition_predicate is not None:
diff --git a/paimon-python/pypaimon/tests/native_plan_integration_test.py 
b/paimon-python/pypaimon/tests/native_plan_integration_test.py
index e308f3a611..5eb062afad 100644
--- a/paimon-python/pypaimon/tests/native_plan_integration_test.py
+++ b/paimon-python/pypaimon/tests/native_plan_integration_test.py
@@ -17,10 +17,13 @@
 
 import tempfile
 import unittest
+from unittest.mock import patch
 
 import pyarrow as pa
 
 from pypaimon import CatalogFactory, Schema
+from pypaimon.globalindex.global_index_result import GlobalIndexResult
+from pypaimon.utils.range import Range
 
 
 def _has_native_planner():
@@ -31,6 +34,14 @@ def _has_native_planner():
     return hasattr(PaimonCatalog, 'get_table') and hasattr(Split, 'serialize')
 
 
+def _has_native_row_ranges():
+    try:
+        from pypaimon_rust.datafusion import ReadBuilder
+    except ImportError:
+        return False
+    return hasattr(ReadBuilder, 'with_row_ranges')
+
+
 @unittest.skipUnless(_has_native_planner(),
                      "pypaimon_rust with split-planning API not installed")
 class NativePlanIntegrationTest(unittest.TestCase):
@@ -179,6 +190,46 @@ class NativePlanIntegrationTest(unittest.TestCase):
             for data_file in split.files
         ))
 
+    @unittest.skipUnless(_has_native_row_ranges(),
+                         "pypaimon_rust row-range API not installed")
+    def test_data_evolution_global_index_row_ranges(self):
+        self.cat.create_table('default.de_range_t', Schema.from_pyarrow_schema(
+            self.schema, options={
+                'row-tracking.enabled': 'true',
+                'data-evolution.enabled': 'true',
+            }), False)
+        self._write('de_range_t', [
+            {'k': 1, 'v': 'a'},
+            {'k': 2, 'v': 'b'},
+            {'k': 3, 'v': 'c'},
+        ])
+        table = self.cat.get_table('default.de_range_t').copy(
+            {'scan.native-plan.enabled': 'true'})
+        builder = table.new_read_builder()
+        scan = builder.new_scan().with_global_index_result(
+            GlobalIndexResult.from_range(Range(1, 1)))
+
+        self.assertTrue(scan._native_plan_supported())
+        with patch.object(
+                scan.file_scanner, 'scan', 
side_effect=AssertionError("fallback")):
+            plan = scan.plan()
+        rows = builder.new_read().to_arrow(plan.splits()).to_pylist()
+
+        self.assertEqual(rows, [{'k': 2, 'v': 'b'}])
+        self.assertEqual(
+            [(range_.from_, range_.to)
+             for range_ in plan.splits()[0].row_ranges()],
+            [(1, 1)],
+        )
+
+        empty_scan = builder.new_scan().with_global_index_result(
+            GlobalIndexResult.create_empty())
+        with patch.object(
+                empty_scan.file_scanner, 'scan',
+                side_effect=AssertionError("fallback")):
+            empty_plan = empty_scan.plan()
+        self.assertEqual(empty_plan.splits(), [])
+
     def test_filter_is_pushed_to_native_plan(self):
         options = {
             'source.split.target-size': '1b',
diff --git a/paimon-python/pypaimon/tests/native_plan_test.py 
b/paimon-python/pypaimon/tests/native_plan_test.py
index 3abecc22f8..3f74a3fb86 100644
--- a/paimon-python/pypaimon/tests/native_plan_test.py
+++ b/paimon-python/pypaimon/tests/native_plan_test.py
@@ -28,6 +28,8 @@ from pypaimon.common.options.core_options import CoreOptions
 from pypaimon.common.options.options import Options
 from pypaimon.common.predicate import Predicate
 from pypaimon.common.predicate_builder import PredicateBuilder
+from pypaimon.globalindex.global_index_result import GlobalIndexResult
+from pypaimon.globalindex.vector_search_result import ScoredGlobalIndexResult
 from pypaimon.read.native_plan import (
     _catalog_options,
     _predicate_to_native,
@@ -38,6 +40,7 @@ from pypaimon.read.native_plan import (
 from pypaimon.read.scan_stats import ScanStats
 from pypaimon.read.table_scan import TableScan
 from pypaimon.table.bucket_mode import BucketMode
+from pypaimon.utils.range import Range
 
 
 def _scan(native_enabled, file_scanner):
@@ -121,7 +124,8 @@ class NativePlanTest(unittest.TestCase):
             plan = scan.plan()
 
         np.assert_called_once_with(
-            scan.table, predicate=None, limit=None, projection=None)
+            scan.table, predicate=None, limit=None, projection=None,
+            row_ranges=None)
         fs.scan.assert_not_called()
         self.assertEqual(plan.splits(), [keep])
 
@@ -172,11 +176,87 @@ class NativePlanTest(unittest.TestCase):
             predicate=scan.predicate,
             limit=5,
             projection=['k', 'dt'],
+            row_ranges=None,
         )
 
+    def test_plan_forwards_global_index_row_ranges(self):
+        fs = Mock(partition_key_predicate=None)
+        scan = _scan(native_enabled=True, file_scanner=fs)
+        fs.data_evolution = True
+        fs._global_index_result = GlobalIndexResult.from_ranges([
+            Range(1, 2), Range(5, 5)])
+        split = Mock(partition=Mock(values=[]), snapshot_id=3)
+
+        with patch('pypaimon.read.native_plan.native_plan', 
return_value=[split]) as np:
+            plan = scan.plan()
+
+        np.assert_called_once_with(
+            scan.table,
+            predicate=None,
+            limit=None,
+            projection=None,
+            row_ranges=[(1, 2), (5, 5)],
+        )
+        fs.scan.assert_not_called()
+        self.assertEqual(plan.splits(), [split])
+
+    def test_empty_global_index_result_does_not_fall_back(self):
+        fs = Mock(partition_key_predicate=None)
+        scan = _scan(native_enabled=True, file_scanner=fs)
+        fs.data_evolution = True
+        fs._global_index_result = GlobalIndexResult.create_empty()
+
+        with patch('pypaimon.read.native_plan.native_plan', return_value=[]) 
as np:
+            plan = scan.plan()
+
+        np.assert_called_once_with(
+            scan.table,
+            predicate=None,
+            limit=None,
+            projection=None,
+            row_ranges=[],
+        )
+        fs.scan.assert_not_called()
+        self.assertEqual(plan.splits(), [])
+
+    def test_scored_global_index_result_falls_back(self):
+        fs = Mock(partition_key_predicate=None)
+        sentinel = object()
+        fs.scan.return_value = sentinel
+        scan = _scan(native_enabled=True, file_scanner=fs)
+        fs.data_evolution = True
+        bitmap = GlobalIndexResult.from_range(Range(1, 1)).results()
+        fs._global_index_result = ScoredGlobalIndexResult.create(
+            bitmap, lambda _: 1.0)
+
+        with patch('pypaimon.read.native_plan.native_plan') as np:
+            self.assertIs(scan.plan(), sentinel)
+
+        np.assert_not_called()
+        fs.scan.assert_called_once_with()
+
+    def test_global_index_row_ranges_require_data_evolution_append_table(self):
+        result = GlobalIndexResult.from_range(Range(1, 1))
+
+        for data_evolution, primary_key in ((False, False), (True, True)):
+            with self.subTest(
+                    data_evolution=data_evolution, primary_key=primary_key):
+                fs = Mock(partition_key_predicate=None)
+                fs.scan.return_value = fallback = object()
+                scan = _scan(native_enabled=True, file_scanner=fs)
+                fs.data_evolution = data_evolution
+                fs._global_index_result = result
+                scan.table.is_primary_key_table = primary_key
+
+                with patch('pypaimon.read.native_plan.native_plan') as np:
+                    self.assertIs(scan.plan(), fallback)
+
+                np.assert_not_called()
+                fs.scan.assert_called_once_with()
+
     def test_plan_falls_back_when_scan_is_not_plain(self):
-        # Native planning does not carry shard/slice, global-index, row ranges,
-        # or incremental scans -> must fall back to the file scanner.
+        # Native planning does not carry shard/slice, explicit row ranges,
+        # arbitrary global-index results, or incremental scans.
         def check(setup):
             fs = Mock(partition_key_predicate=None)
             sentinel = object()
@@ -310,7 +390,8 @@ class NativePlanTest(unittest.TestCase):
         self.assertIs(plan, fallback_plan)
         self.assertIs(stats, fallback_stats)
         np.assert_called_once_with(
-            scan.table, predicate=None, limit=None, projection=None)
+            scan.table, predicate=None, limit=None, projection=None,
+            row_ranges=None)
         fs.scan_with_stats.assert_called_once_with()
         fs.scan.assert_not_called()
 
@@ -472,8 +553,9 @@ class NativePlanTest(unittest.TestCase):
         split = Mock()
         split.serialize.return_value = b'bytes'
         rt = Mock()
-        
rt.new_read_builder.return_value.new_scan.return_value.plan.return_value \
-            .splits.return_value = [split]
+        builder = rt.new_read_builder.return_value
+        builder.with_row_ranges.return_value = builder
+        builder.new_scan.return_value.plan.return_value.splits.return_value = 
[split]
         catalog = Mock()
         catalog.get_table.return_value = rt
 
@@ -488,13 +570,14 @@ class NativePlanTest(unittest.TestCase):
                 patch('pypaimon.read.native_plan._catalog_options', 
return_value={}), \
                 patch('pypaimon.read.native_plan.deserialize_split_v1',
                       return_value='decoded') as des:
-            result = native_plan(table)
+            result = native_plan(table, row_ranges=[(1, 2)])
 
         self.assertEqual(result, ['decoded'])
         rt.new_read_builder.assert_called_once_with({
             CoreOptions.SOURCE_SPLIT_TARGET_SIZE.key(): '1024',
             CoreOptions.SOURCE_SPLIT_OPEN_FILE_COST.key(): '128',
         })
+        builder.with_row_ranges.assert_called_once_with([(1, 2)])
         des.assert_called_once_with(b'bytes', [], kfields)
 
     def test_native_plan_requires_split_api(self):
@@ -519,6 +602,5 @@ class NativePlanTest(unittest.TestCase):
                     with self.assertRaisesRegex(RuntimeError, '0.3.0'):
                         native_plan(Mock())
 
-
 if __name__ == '__main__':
     unittest.main()

Reply via email to